F1 Kafka & Event-Driven Architecture
Loading learning experience...
Lecture transcript
Read the narration for F1: Kafka & Event-Driven Architecture
From Redis speed to Kafka durability
Dr. Wei: Last lecture we used Redis for sub-millisecond reads and writes and rich data structures. That solves the speed problem. But here is the stake: if your system needs an audit trail, analytics, or cross-service fan-out, a cache cannot be the source of truth. Today we are going to discover Kafka and event-driven architecture, and it will feel like a simple idea: treat events as an append-only log. One important nuance: Kafka keeps a durable event log with configurable retention and, in some topics, compaction, so it is often a source of truth for events. But depending on your domain and how long you must retain data, the authoritative system of record may still be a database owned upstream. We will go from why async helps, to how partitions and offsets work, to what replay and delivery guarantees really mean.
Sam: So Redis is for fast current state, and Kafka is for the history of what happened?
Dr. Wei: Exactly. Redis accelerates reads of state; Kafka preserves a durable stream of events that led to that state, which lets you rebuild, debug, and fan out safely. Just remember that the event log has retention or compaction policies, and a separate database may still be the long-term authoritative record depending on requirements.
Dr. Wei: By the end, you should be able to look at any system design where someone draws a box called message queue and ask the real questions: ordering, replay, retention, and failure semantics.
Why event-driven architecture exists
Dr. Wei: This slide is the motivation: synchronous calls create tight coupling. If B is slow, A becomes slow. If B is down, A might go down too. That is how you get cascading failures.
Sam: Is the main risk here that one slow dependency can stall the whole user request, even if the rest of the system is healthy?
Dr. Wei: When you see a chain like A to B to C, you should immediately worry about retries, timeouts, and backpressure. Those worries are not optional; they are the architecture now.
Sam: So the broker is basically a buffer between services?
Dr. Wei: Yes, and more than a buffer: it is a durable record of events. Buffering absorbs spikes, decoupling prevents cascades, and replay gives you a time machine for debugging and recomputation.
Kafka as a distributed commit log
Dr. Wei: Here is the core mental model: Kafka is a distributed commit log. Think of it like a file you only append to, and you never edit old lines.
Dr. Wei: Append-only and immutable makes durability and replication simpler. It also makes reprocessing possible, because the past is still there until retention removes it.
Sam: When you say ordering, is that global ordering across the whole topic?
Topics and partitions: parallelism with ordering
Dr. Wei: Conceptually, think of a topic as a stream of events. That stream is split into partitions, and each partition is its own ordered log with offsets.
Dr. Wei: A topic is just the name you publish to, like user-actions. Under the hood, it is partitioned so multiple machines can handle it.
Sam: So partitions are the reason Kafka can scale reads and writes?
Dr. Wei: Exactly. But notice the constraint: if a consumer group wants each event processed once, each partition can be owned by only one consumer in that group at a time. So your maximum parallelism inside a group is the number of partitions.
Example: fan-out with one event, many consumers
Dr. Wei: Here is the classic event-driven move in system design: publish one event, and let many downstream services react independently.
Dr. Wei: The event is small: post created. It does not try to do everything. It just states that something happened, and the log becomes the shared truth of that fact.
Sam: Why is this better than the Post service just calling Feed, Search, and Notifications directly?
Dr. Wei: Because direct calls turn one write into three dependencies and three failure modes. With events, each consumer can scale on its own, retry on its own, and even go down temporarily without breaking the Post write path.
Producers: picking a partition controls ordering
Dr. Wei: Now we make the partition choice precise. Conceptually, the partition index is roughly the hash of the key, modulo the number of partitions. The exact behavior depends on the partitioner Kafka is using.
Dr. Wei: If you key by something like user id, all events for that user land in the same partition, so the consumer sees them in the order they were appended. In many setups that means a murmur2-style hash of the key is used to pick the partition.
Sam: So ordering is something I buy by giving up parallelism for that key?
Dr. Wei: Exactly. If you do not set a key, Kafka can spread events across partitions for throughput, but then any per-user ordering is gone because events can interleave across partitions. And depending on the version and partitioner, unkeyed records may use a sticky approach that batches to a partition for a while, which is not the same as a simple hash rule.
Consumer groups, offsets, and lag
Dr. Wei: Offsets are the second core idea. Each partition is a log, and an offset is your position in that log. Lag is simply the gap between where the partition ends and where you have committed.
Dr. Wei: Quick prediction: suppose the end offset is 12,500 and your committed offset is 12,200. What is L, and operationally what does that suggest you should do next?
Sam: So the lag would be 300 messages. And if that keeps growing, I either scale the consumer group up until I hit the partition count, or speed up processing, right?
Sam: And is lag measured per partition, then summed up for the whole group when you look at dashboards?
Sam: Why does the group rule say one consumer per partition? Why not let two consumers share a partition to go faster?
Dr. Wei: Because a partition is an ordered sequence, and sharing it would require coordination to avoid double-processing and reordering. Kafka chooses a simpler contract: one consumer owns the partition, and you scale by adding partitions, not by splitting a partition.
Replay: fixing bugs by reprocessing history
Dr. Wei: Replay is why logs are powerful: you can treat the past as data. If an analytics job was wrong last Tuesday, you do not need to guess; you reprocess from Tuesday.
Dr. Wei: The bug creates wrong downstream state, but the raw events are still correct. So you correct the processing logic and reset the committed offset back to the point in time where correctness diverged.
Sam: Is resetting an offset dangerous? Could I accidentally double-count everything?
Dr. Wei: Yes, which is why replay forces you to think about idempotency and how you store outputs. And you also have a hard limit: if retention is seven days and your bug was two weeks ago, the raw events are already gone.
Delivery guarantees are really offset semantics
Dr. Wei: Most confusion about Kafka comes from guarantees. The shortcut is: guarantees are about when you commit offsets relative to processing, and what you do about duplicates.
Dr. Wei: Here is a quick crash scenario to reason with, but we will split it into two cases with clear assumptions. Case one is offsets-only semantics: assume processing has no external side effects, and we only care whether the consumer will re-read message M after a crash, based on when offsets are committed. Case two is end-to-end semantics: we do write to an external database, and to avoid duplicates we must assume either idempotent database upserts or an outbox or transactional sink that can coordinate the write with the Kafka offset commit. With those assumptions stated, we read message M, we process it, and then we crash before the offset commit. Predict what happens on restart for each case under at-most-once, at-least-once, and exactly-once: do we reprocess M, do we lose it, or do we risk a duplicate effect?
Sam: So if the commit did not happen, we will see M again after restart, which means at-least-once repeats it, and without idempotency we can duplicate the write?
Dr. Wei: Exactly. In the offsets-only case, if the commit was after processing and it never happened, then on restart at-least-once will re-read and reprocess M. If you commit before processing, as in at-most-once, a crash during processing can mean M is effectively lost because Kafka thinks you already moved past it. Now in the end-to-end case with a database write, the same restart behavior can translate into duplicate effects unless the sink is idempotent. Kafka exactly-once is typically scoped to consume-transform-produce within Kafka, where transactions make the produced output and the offset commit succeed or fail together. But for an external database, you still need idempotent writes or an outbox or transactional sink pattern; Kafka transactions alone do not automatically make the database write exactly-once.
Durability: replication, leaders, and acknowledgments
Dr. Wei: Kafka is not just an API; it is a distributed storage system. Durability comes from replication: each partition is copied across multiple brokers.
Dr. Wei: With replication, one broker is the leader for a partition and the others are followers that copy the log. A key idea is the in-sync replica set, meaning replicas that are caught up enough to take over safely. With acks set to one, the leader confirms the write as soon as it has it. With acks set to all, the leader waits until all in-sync replicas have it, trading lower risk of loss for higher latency.
Sam: What happens if the leader dies right after it acked a write?
Dr. Wei: That is exactly why ack choice matters. If the leader acked but the write was not yet on an in-sync replica, the write can be lost on failover. If you require acks set to all, a new leader chosen from the in-sync set will have the data, but your tail latency is higher.
Scaling pitfalls you should mention in interviews
Dr. Wei: If you want to sound operationally mature, these are the three pitfalls to call out without being prompted: lag, hot partitions, and message size.
Dr. Wei: If lag grows, you can add consumers, but only until you hit the partition count. Past that, you need more partitions or faster processing, not more replicas of the same consumer group.
Sam: And hot partitions happen when one key dominates, like a celebrity user id?
Dr. Wei: Exactly. One key forces all traffic to one partition, so you lose parallelism. You fix it by changing the partitioning scheme, like adding a sub-key. And for large payloads, store the blob elsewhere and send only a pointer through Kafka.
Exit ticket: partitions, parallelism, and lag
Dr. Wei: Let us test the two most interview relevant calculations: parallelism and lag. You are given partitions, consumers, and offsets.
Dr. Wei: With 6 partitions and 3 consumers in a group, the maximum parallelism is 3, because each partition can be owned by only one consumer at a time. That usually means each consumer gets about two partitions. For lag, subtract the committed offset from the end offset: twelve hundred minus nine hundred equals three hundred messages behind. For the reflection: choose at least once when you can tolerate duplicates by making processing idempotent, and choose exactly once when duplicates are unacceptable and you can afford the added complexity.
Thank you for watching!
Thanks for watching. Subscribe and share if you found this useful—see you next time!