Part 7 · 2 chapters · ~18 min
Logs, Queues and Streams
The log as the backbone: Kafka's partitions and per-key order, replication with in-sync replicas, consumer-owned offsets and groups, retention, replay and lag, and log compaction; then stream processing: stateless and stateful operators, windows, event time and watermarks, joins and exactly-once state.
14
Kafka's model and log compaction
the log as the backbone
- Partitions: the key picks the partition, and order holds only within a partition.
- Replication: a leader plus in-sync replicas, with acks=all and min.insync=2.
- Offsets belong to consumers. Commit after processing, and make processing idempotent.
- Consumer groups: within a group, each partition goes to one consumer. Every group gets everything.
- Retention and replay. Consumer lag is the health metric.
- Compaction keeps the latest record per key, which turns a topic into a changelog of current state.
code
// producer: keyed for per-account order, durable, idempotent
const producer = kafka.producer({ idempotent: true, maxInFlightRequests: 5 });
await producer.send({
topic: 'payments.events',
acks: -1, // all in-sync replicas
messages: [{ key: accountId, value: JSON.stringify(evt), headers: { 'event-id': evt.id } }],
});
// consumer: commit only after the effect (and dedup by event-id: part 6)
await consumer.run({
autoCommit: false,
eachMessage: async ({ topic, partition, message }) => {
await applyIdempotently(message); // inbox table, same transaction
await consumer.commitOffsets([{ topic, partition, offset: (BigInt(message.offset) + 1n).toString() }]);
},
});THE LOG: KAFKA'S MODEL
an append-only, partitioned, replicated log that many consumers read at their own pace
swipe the figure sideways, or tap expand for full screen
1/6
partitions
Topic and partitions: a topic (payments.events) is split into partitions; a record's key (the account id) decides its partition by hash, so all events for one account are in one partition, in order. Order is guaranteed within a partition, never across partitions.
15
Stream processing
computing over unbounded data
- Stateless transforms scale by partition.
- Stateful aggregation keeps a local store backed by a changelog.
- Windows: tumbling, hopping and session.
- Event time gives correct results, though you can never be sure a window is complete.
- Watermarks fire windows, and allowed lateness emits corrections.
- Joins run against tables or within windows. Checkpoints make state updates exactly-once.
code
-- Flink SQL: failed transfers per bank per 5 minutes, by event time, tolerating 30 s of lateness
CREATE TABLE transfers (
bank_code STRING, status STRING, created_at TIMESTAMP(3),
WATERMARK FOR created_at AS created_at - INTERVAL '30' SECOND
) WITH ('connector' = 'kafka', 'topic' = 'payments.events', 'format' = 'json', …);
SELECT bank_code, window_start, COUNT(*) AS failures
FROM TABLE(TUMBLE(TABLE transfers, DESCRIPTOR(created_at), INTERVAL '5' MINUTES))
WHERE status = 'FAILED'
GROUP BY bank_code, window_start, window_end;
-- feeds the "bank degraded" banner the Trust course part 7 shows usersSTREAM PROCESSING
computing over unbounded data: stateless transforms, windows, joins, event time and watermarks
swipe the figure sideways, or tap expand for full screen
1/6
stateless
Stateless operations: map, filter and enrich each event independently (parse, drop test traffic, add the merchant's category from a cached lookup). Scale out by partition; nothing to remember between events; restart anywhere.