Apache Kafka¶
Version and source policy
Examples target the pinned lab baseline. Check Versions & Primary Sources before applying configuration to another Kafka release.
11:58 PM. Billing's consumer for service-events is nine minutes behind and climbing. The API gateway is still writing tens of thousands of records per second — auth, billing, fraud, search, and the warehouse all read the same stream, each at its own pace. If that stream were a shared database table, the fix would be "wait" or "drop rows." Neither is acceptable tonight.
Before you read on: how do you let five independent consumers read the same firehose at five different speeds, with one of them nine minutes behind, without slowing the producer down and without losing what a faster consumer already deleted?
A request/response API stalls the producer on the slowest consumer. A database table used as a queue lets the second reader see only what the first has not yet deleted. A traditional broker that acknowledges and drops loses the evidence the moment one consumer says "done." Kafka exists because producers and consumers must evolve independently without losing the stream.
This is SaaSCo at Stage 3
SaaSCo: The Evolving Company hits this exact wall at 4 TB/day: many producers, one object store, and PUT-rate contention that a nightly Spark job never had to deal with.
What this module covers¶
| Topic | What you will be able to reason about |
|---|---|
| The log abstraction | Segments, indexes, retention versus compaction, why Kafka is not a database |
| Partitions and consumers | Partitioner, consumer groups, eager versus cooperative rebalance, lag |
| Replication and durability | ISR, acks, min.insync.replicas, unclean leader election |
| Exactly-once semantics | Idempotent producers, transactions, what EOS actually guarantees |
| Production gotchas | Hot partitions, poison pills, schema evolution, rebalance storms |
| Labs | Produce, consume, watch lag, then kill a broker and inject a poison message |
This module sits on partitioning, data movement, and batch versus stream. Downstream you will use the same topics from Flink.
The central intuition¶
Kafka is not a queue. It is a distributed, replicated, append-only log.
Producers append. Consumers read from an offset. The log stays around for a retention window (or is compacted to the latest value per key). Ten consumer groups can read the same bytes at ten different speeds. Replay is seeking, not begging the producer to send again.
Partition 0 of `service-events`
offset 0 1 2 3 4 5 6 7
e0 → e1 → e2 → e3 → e4 → e5 → e6 → e7 → (tail)
alert-processor committed offset 3
warehouse-loader committed offset 6
fraud-scorer committed offset 7 (caught up)
A queue that deletes on ack cannot do this. A database table that is not an append log fights the access pattern: sequential write, sequential read, many independent cursors.
Five systems, one event shape¶
Reuse this record throughout the module. Serialise it however you like in labs (JSON is fine); production will argue about Avro and Protobuf later.
{
"timestamp": "2024-01-15T10:03:45.123Z",
"customer_id": "cust_1842",
"user_id": "u_99102",
"service": "auth",
"endpoint": "/login",
"region": "eu-west-1",
"latency_ms": 87,
"status_code": 401,
"bytes": 512
}
| System | How Kafka shows up |
|---|---|
| SaaS analytics | Product events from hundreds of tenants. Partition by customer_id so one tenant's burst does not reorder another tenant. Warehouse and real-time dashboards are different consumer groups on the same topic. |
| Observability | Logs and traces at 500k–2M records/s. Retention is hours to days. Replay exists so a broken ClickHouse sink can catch up without asking apps to re-emit. |
| E-commerce | Orders, payments, inventory. Per-order_id (or user_id) order matters. Duplicate checkout events are a money bug — see exactly-once. |
| IoT | Millions of devices, bursty reconnects. Keys are device_id. Compaction on a "latest device state" topic is a different access pattern from the raw telemetry topic. |
| Fraud | Same login events as observability, plus a low-latency scorer. A slow model must not stall ingest. Lag on the fraud group is an SLA; lag on the archive group is a disk problem. |
The security / observability platform is the running numeric example: ~500k log lines/s, ~2M metric points/s, ~50k security events/s. Multiple processors (alert, enrich, store, train) must not couple back to the fleet of producers.
Where it sits in an architecture¶
flowchart LR
subgraph producers [Producers]
API[API / auth / billing]
Agents[Sidecars / agents]
end
subgraph kafka [Kafka cluster]
T1["topic service-events"]
T2["topic login-events"]
T3["compacted user-profile"]
end
subgraph consumers [Independent consumer groups]
Alert[Alerting]
Flink[Flink jobs]
WH[Warehouse loader]
Lake[Lake / Iceberg]
end
API --> T1
API --> T2
Agents --> T1
T1 --> Alert
T1 --> Flink
T1 --> WH
T2 --> Flink
T3 --> Flink
Flink --> Lake Producers need three things from this layer: high write throughput, durability they can name (acks, ISR), and a partition key that preserves the ordering they actually care about. Consumers need independent offsets, enough partitions to parallelise, and a retention window longer than their worst outage.
Kafka does not compute "failed logins per user in five minutes". That is Flink. Kafka stores and fans out the raw stream.
What Kafka actually guarantees (and does not)¶
| Guarantee | Scope |
|---|---|
| Order | Per partition, not per topic |
| Durability | As strong as acks × ISR × min.insync.replicas × disk |
| Delivery to a consumer group | At-most-once, at-least-once, or exactly-once for a read-process-write to Kafka — see EOS |
| Fan-out | Independent groups, independent offsets |
| Replay | Until retention or compaction removes the record |
It does not give you:
- Random access by
user_idacross partitions (scan or materialise elsewhere) - Cross-partition transactions that include Postgres, Stripe, or a webhook
- Automatic schema compatibility because you "use a registry" (the registry only helps if producers and consumers agree the compatibility mode)
Throughput intuition before you tune¶
A single partition is a single log on a single leader broker, consumed by at most one member of a group. Millions of records per second therefore means many partitions, many leaders, and many consumers — plus the disk sequential-write path that Kafka was built around.
| Scale | What usually breaks first |
|---|---|
| 10× (tens of thousands/s) | Bad keys: one customer_id owns a partition. Consumer lag on that partition only. |
| 100× | Disk and page cache. Replication traffic ≈ ingest × (RF − 1). Under-replicated partitions after a broker bounce. |
| 1000× | Partition count, controller metadata, rebalance time, request-handler threads, and the operational cost of "just add partitions" (it reshuffles keys). |
Details live in partitions and replication. The partition simulator is worth five minutes before the labs.
The cluster you are actually operating¶
A Kafka cluster is brokers plus a controller (KRaft in 3.x; ZooKeeper + controller in older deployments). Clients (bootstrap.servers) ask any broker for metadata: which broker leads service-events partition 7, where the ISR is, what the topic configs are.
You will live in a handful of objects:
| Object | Role |
|---|---|
| Topic | Named log, split into partitions |
| Partition | Ordered log; unit of parallelism and replication |
| Consumer group | Independent cursor over partitions |
| ACL / quota | Who may produce/consume, at what rate |
Producers batch (linger.ms, batch.size) and compress. Followers fetch those batches. Consumers fetch them again. Disk sequential write is the happy path; random reads of old segments (a warehouse backfill) is the unhappy path that evicts page cache from the live tail.
If you remember one operational sentence: Kafka is a disk and metadata system that happens to speak a consumer protocol. Incidents that look like "consumer lag" are often disk, ISR, or rebalance.
Contracts to write down before the first topic¶
Do not create service-events until you can fill this:
key: customer_id # order scope
partitions: 24 # headroom for consumers
RF / min.ISR: 3 / 2
acks: all
cleanup: delete, 48h # or compact for changelogs
consumers: alert, warehouse, flink-fraud # separate group ids
poison policy: DLQ (logs) / halt (payments)
schema: Avro + BACKWARD # not "JSON, we'll be careful"
E-commerce orders will differ (tighter EOS, maybe fewer partitions for order). IoT device-state will differ (compact). Observability logs.raw will differ (short retention, many partitions, no keys). Same product, five contracts.
Related: foundations — scale, observability architecture, fraud architecture.
How to study this module¶
- Read the log until "offset" and "segment" are muscle memory.
- Read partitions with the simulator open. Predict lag before you generate it in the labs.
- Read replication and be able to say what
acks=alldoes when ISR shrinks to one. - Read exactly-once sceptically. Name the pipeline it actually covers.
- Skim gotchas, then do the labs including now break it.
- Continue into Flink with the same
service-eventstopic.
You are done with Kafka as a platform when you can look at a lag graph, an ISR shrink, and a rebalance, and say which one is the incident.
Exit check — you pass this module if you can
- Predict how keys distribute across partitions, and name which key would create a hot partition before running the simulator.
- Identify hot-key risk from a traffic shape (one tenant, one device, one entity id) without needing the incident to happen first.
- Explain what a consumer-group rebalance actually does, and why it briefly pauses processing for the whole group, not just the joining/leaving consumer.
- Reason about replication and ISR: what
acks=allandmin.insync.replicasguarantee, and what happens to writes when ISR shrinks below that minimum. - State precisely where Kafka's exactly-once guarantee starts and stops (producer→topic, not topic→external sink) and what makes an end-to-end pipeline exactly-once anyway.
- Diagnose rising consumer lag from metrics alone: distinguish a hot partition, a slow sink, a rebalance storm, and a genuine capacity shortfall before touching any config.
- Decide when Kafka is unnecessary — a single consumer, low volume, or no replay requirement is often better served by a simpler queue or direct call.