Skip to main content
Mati Data

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. 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. 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. 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. 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. 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. 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. 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

  8. 8

    Concurrency and resource contention

    Schedule shared resources before buying bigger ones

    Many performance incidents are contention, not capacity: jobs that each run fine collide when they start at the same moment. Spreading the work across time is often cheaper and safer than a larger instance.

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.