12.6 Technology deep dives
6.1 Apache Kafka
Problem it solves. Be simultaneously a durable log (rereadable, replayable, retained) and a low-latency notification system — the hybrid §2 opens with.
Why wasn't RabbitMQ enough? Consumption is destructive; a new consumer can't read the past; per-message acking limits throughput; and load balancing + redelivery inevitably reorders (§1.4).
Why wasn't a database enough? Polling cost grows as the hit rate falls (§1); triggers are an afterthought.
How it works internally. Topics → partitions, each an append-only segmented log on disk. Producers pick a partition by key hash (which is how you get per-key ordering). Each partition is replicated; one replica is leader, and the ISR (in-sync replicas) set determines what counts as committed — acks=all + min.insync.replicas=2 is the durability contract that actually matters. Consumers in a group are assigned whole partitions; progress is a committed offset. Log compaction keeps the latest value per key forever (§3.2). Reads use zero-copy sendfile, which is why disk-backed throughput is so high. Transactions + idempotent producer give the internal atomic commit of §5.2.
Deployment. 3+ brokers across AZs; RF=3, min.insync.replicas=2, acks=all; unclean.leader.election.enable=false (Ch 10 §3.5 — otherwise you trade correctness for availability silently); KRaft (Raft) instead of ZooKeeper in modern versions; tiered storage to object store for long retention.
Monitoring.
- Consumer lag per partition — the metric. §2.4: the buffer is large enough that a human can fix a slow consumer before it starts missing messages, but only if you're watching
- Under-replicated partitions and ISR shrink/expand rate — a shrinking ISR is a durability degradation happening quietly
- Offline partitions (no leader available)
- Request latency p99 by request type; fetch-purgatory size
- Rebalance frequency per consumer group — frequent rebalances mean processing time exceeds
max.poll.interval.ms - Partition skew — bytes/messages per partition; a hot key means one consumer does all the work
Scaling. Partition count bounds consumer parallelism (§2.2) and increasing it changes key→partition mapping, breaking per-key ordering for existing keys. Choose generously up front.
Backup. Mirroring to a second cluster (MirrorMaker 2 / Cluster Linking) and/or tiered storage. A compacted topic used for event sourcing is a system of record and needs real backups (Ch 3 §3.4).
What actually breaks in production.
- Retention expiry before a stalled consumer catches up — silent, permanent data loss with no error, just an offset jump. §2.4's ring buffer, unmonitored.
- Unclean leader election enabled → an out-of-sync replica becomes leader → acknowledged writes vanish.
- Hot partition from a low-cardinality key — 90% of traffic on one partition, one consumer saturated, the rest idle.
- Rebalance storms when a consumer's processing time exceeds the poll interval, so the group never stabilizes and throughput goes to zero.
- Assuming global ordering. There is none across partitions (§2.1). Every "events arrived out of order" incident traces back to this.
- Increasing partitions on a keyed topic, silently breaking ordering for keys that move.
- Consumer commits offsets before processing ("at-most-once" by accident) → messages silently dropped on crash.
6.2 Debezium / CDC pipelines
Problem it solves. Make one database the leader and every derived system a follower (§3.2), eliminating dual-write races.
Why wasn't dual writing enough? §3.1's race condition and partial-failure problem — and crucially, you don't even notice it happened.
Why not query the database periodically? You miss intermediate states, you miss deletes, and the polling cost/freshness trade-off is bad.
How it works internally. A connector reads the database's replication log (MySQL binlog, Postgres logical decoding slot / pgoutput, Oracle LogMiner) and emits row-level change events with before/after images plus source metadata, into Kafka. Initial snapshot is bound to a specific log position (§3.2), and Debezium uses DBLog watermarking for incremental, non-blocking snapshots.
Deployment. Kafka Connect cluster; one connector per source database; schema registry for event schemas (Ch 5); an outbox table if you don't want the internal schema to become a public API (§3.3).
Monitoring.
- Replication slot lag / retained WAL on the source — the Postgres CDC failure: an inactive slot pins WAL forever and fills the primary's disk (Ch 6 §6.1). This is a source database outage caused by your pipeline.
- Connector status and snapshot progress
- Event lag (source commit time → Kafka append time)
- Schema-change events — every DDL is a potential downstream break (§3.3)
- Tombstone/delete event rate — an unexpected spike means someone ran a mass delete
Scaling. Bounded by the source's log generation and the connector's single-threaded decode. Filter tables aggressively.
What actually breaks.
- The replication slot filling the source disk when the connector is down or slow. This takes down the production database, not just the pipeline.
- A dropped or renamed column breaking downstream production consumers — §3.3's "schema became a public API," which "can cause a customer-facing outage."
- DDL that logical decoding can't represent, silently skipping data.
- Snapshot + stream boundary bugs — duplicates or gaps at the handover if the snapshot isn't tied to an exact log position.
- CDC on a quorum database (§3.2): per-node raw log segments that must be merged, with no single source of truth.
- Confusing CDC with event sourcing — expecting business intent from what is actually a row diff.
6.3 Apache Flink
Problem it solves. Stateful, event-time-correct stream processing with exactly-once state semantics and large windows/joins.
Why wasn't microbatching enough? §5.1: microbatching imposes a tumbling window by processing time equal to the batch size, and forces a latency/overhead trade-off. Flink's barrier-triggered checkpoints decouple fault tolerance from window size.
Why wasn't a simple consumer loop enough? Because windows, joins, and aggregations need state that must survive failure and event-time semantics that a naive loop doesn't provide.
How it works internally. A dataflow graph of operators; Chandy–Lamport-style asynchronous barrier snapshotting — barriers flow through the stream, each operator snapshots its state when barriers from all inputs align, and the aligned snapshot is written to durable storage. Watermarks carry event-time progress and trigger windows; allowed lateness handles stragglers (§4.2). State backends: heap or RocksDB (for state larger than memory), with incremental checkpoints. Two-phase commit sinks extend exactly-once past the framework boundary for sinks that support it (§5.2).
Deployment. JobManager + TaskManagers on K8s/YARN; checkpoint storage on a DFS/object store; savepoints for planned upgrades (a savepoint is how you change code without losing state).
Monitoring.
- Checkpoint duration, size, and failure rate — growing duration is the leading indicator of every Flink problem
- Backpressure per operator (Flink exposes this directly — it tells you which operator is the bottleneck)
- Watermark lag per source — how far behind event time you are, and therefore when windows will fire
- State size per operator; RocksDB compaction and memory
- Records dropped as late (§4.2's "track the number of dropped events as a metric and alert")
- Restart count and time-to-recover
Scaling. Parallelism per operator; keyed state is partitioned by key, so a hot key means a hot subtask. Rescaling requires a savepoint.
What actually breaks.
- Checkpoints taking longer than the checkpoint interval, so they overlap and eventually the job can't make progress. Usually caused by backpressure or by state that grew unbounded.
- Unbounded state — a keyed state entry per user with no TTL, a stream–stream join with too wide a window, or session windows on a key space that never stops growing.
- Watermark stalls. One idle source partition holds the watermark back and no windows ever fire — the job looks healthy and emits nothing.
- Late data silently dropped because allowed lateness is zero and nobody instrumented the drop counter.
- Nondeterministic joins across streams (§4.3) making reprocessing produce different results.
- Losing state on redeploy because someone restarted without a savepoint.
6.4 Kafka Streams / ksqlDB (and IVM engines: Materialize, RisingWave)
Problem it solves. Maintain materialized views — the "window that stretches back to the beginning of time" (§4.1c) — continuously and incrementally.
Why wasn't REFRESH MATERIALIZED VIEW enough? §4.1: poor efficiency (all data reprocessed) and poor freshness (stale until the next scheduled run).
Why weren't triggers enough? They work only when "the data is easily partitioned and the computation is naturally incremental" — and "many SQL queries can't be easily or efficiently converted to incremental computation."
How it works internally. Kafka Streams: a library, not a cluster; KTable = a compacted changelog interpreted as a table; KStream = an event stream; joins follow §4.3's three types directly. Local state in RocksDB, backed by a compacted Kafka changelog topic so it can be rebuilt anywhere (§5.4). IVM engines (Materialize, RisingWave, Feldera) go further: they compile SQL into incremental dataflow operators (differential dataflow lineage), buffering recent events in memory and periodically updating on-disk views, with reads combining both for a real-time answer.
Monitoring. View staleness / end-to-end lag; state-store size and restore time (how long to rebuild from the changelog — that's your RTO); changelog topic size; rebalance-induced state migration; for IVM engines, memory per materialized view and arrangement size.
What actually breaks.
- State restore time after a rebalance — a multi-GB RocksDB store rebuilt from a changelog topic can take many minutes, during which that partition is unavailable. Standby replicas exist for exactly this.
- Unbounded KTable growth where keys are never deleted.
- A materialized view that's expensive to maintain incrementally — some joins produce huge intermediate arrangements.
- Read-your-writes violations (§3.4) — the user writes, then reads the view before the update lands.
- Assuming the view is transactionally consistent with the source — it isn't; CDC is asynchronous.