How I reason about data-intensive systems
The principles I design by, each named after the concept it comes from and shown in my public code.
The principles
Each one comes from Designing Data-Intensive Applications or Database Internals, with a real example where my public work shows it.
- 1
Replication and change data capture
The log is the source of truth
A database's write-ahead log or binlog already records every change in commit order. Reading changes from the log, instead of querying tables for what looks new, captures deletes, keeps order, and lets any downstream table be rebuilt by replaying it.
In practice: In wb-cdc-analytics Debezium reads the PostgreSQL write-ahead log through pgoutput instead of polling tables, so deletes arrive as events with the operation and the LSN beside each row, and ClickHouse keeps an immutable event log that the marts can be rebuilt from.
wb-cdc-analytics: PostgreSQL to ClickHouse CDC in one command
- 2
Idempotence and exactly-once semantics
Every hop is at-least-once, so make duplicates harmless
Networks retry and consumers crash after writing but before acknowledging. Rather than pretend duplicates cannot happen, design every write to be idempotent: a stable natural key, and a store that collapses a repeated event into the row it repeats.
In practice: In wb-cdc-analytics the event log is keyed on the Kafka topic, partition and offset, so a redelivered message lands on the row it duplicates, and the fact table is rebuilt with delete-and-insert on the natural key.
wb-cdc-analytics: PostgreSQL to ClickHouse CDC in one command
- 3
Ordering, logical clocks and log positions
Order by the source's commit order, not by arrival
Events arrive out of order under retries and parallelism. The newest version of a row is the one with the highest log position at the source, so that position, not the arrival time, decides which version wins.
In practice: In wb-cdc-analytics the ClickHouse landing tables are ReplacingMergeTree tables versioned by the source LSN, so the newest version of a row follows PostgreSQL's commit order even when its events arrive out of order.
wb-cdc-analytics: PostgreSQL to ClickHouse CDC in one command
- 4
Transactions and atomic commit
Commit the data and the watermark together
A pipeline that records progress separately from the data it wrote can claim work that rolled back, or redo work that succeeded. Writing the data and the new watermark in one transaction makes a retried run either fully applied or not at all.
In practice: In wb-cdc-analytics the rows, the audit row and the watermark commit in one PostgreSQL transaction, so a recorded position can never claim work that rolled back, and a second run of the same data writes nothing.
wb-cdc-analytics: PostgreSQL to ClickHouse CDC in one command
- 5
Storage engines: B-trees, log-structured and merge trees
Choose the storage engine for the access pattern
B-tree row stores serve transactions and point lookups; log-structured and merge-tree engines trade write cost for scan speed and merge in the background. Each engine's merge and compaction behaviour decides when a read is correct, not only how fast it is.
In practice: ClickHouse ReplacingMergeTree deduplicates only when parts merge, so a plain SELECT can legitimately return duplicates. In wb-cdc-analytics correct reads go through one shared macro that applies FINAL and drops tombstones.
wb-cdc-analytics: PostgreSQL to ClickHouse CDC in one command
- 6
Sharding and partitioning
Partition when the data asks for it
Partitioning pays off when it prunes scans or lets retention drop whole partitions. Below that size it only creates many small parts, slower merges and, in merge-tree engines, a path to refused inserts. The trigger for partitioning belongs in writing, next to the design.
In practice: The reference CDC pipeline deliberately does not partition its tables at its current size, and documents the volume and retention conditions under which monthly partitions become worth it.
wb-cdc-analytics: PostgreSQL to ClickHouse CDC in one command
- 7
Reliability and faults
Design for the failure that raises no error
The expensive failures are the quiet ones: a component that keeps running while producing wrong data. Where a design can fail loudly or fail silently, choose loudly: reject bad rows with a reason, gate models on contracts, and alert on drift.
In practice: In wb-cdc-analytics the Kafka engine reads raw JSON and types it in materialized views, because typed Kafka-engine columns turn one parse failure into a consumer that retries forever without an error. CDC lag is measured against a heartbeat, so a stopped connector cannot look like a quiet table.
wb-cdc-analytics: PostgreSQL to ClickHouse CDC in one command
References I design from
- Designing Data-Intensive Applications, 2nd edition
Martin Kleppmann and Chris Riccomini
The vocabulary I use in design reviews: replication, sharding, transactions, consistency, and batch and stream processing, with the tradeoffs between them.
- Database Internals
Alex Petrov
How storage engines and distributed databases work underneath: B-trees and log-structured storage, then replication, consistency and consensus.