Back to System Design Index

Free Course15 lessonsLive diagramsMukul

Event Streaming in Production: 15 Lessons From the Kafka Bench

Nine years of Golang on AWS taught me one discipline: append everything, assume nothing arrives once, plus price the backlog. Each lesson starts from a real production failure, states its assumptions up front, plus ends in math you can re-run. The bill hides in the boundary.

TL;DR: Fifteen forensic streaming lessons with live diagrams, from partitions to exactly-once math plus rebalance drills. Built from production postmortems with calculators for the math.

By Mukul Kumar Mishra · Backend plus SRE course · Updated September 19, 2026

What you will be able to do

How interviews test this course: clarify scope first, estimate out loud second, break your own design third. Lesson 14 rehearses the bench.

Lesson 01 · Foundations · Append-only truth

The Log Is the Truth

Assumption first: every downstream system disagrees about what happened. The log settles it. One gaming platform routes 18 trillion messages a day through one append-only backbone, then replays any of them on demand. Databases forget. Logs remember. Design the memory first.

producers append-only log consumer A consumer B replay: any truth kept
  • Append, never overwrite. Immutable history makes every downstream rebuild a replay, not a rescue.
  • One backbone, many readers. Matchmaking plus analytics read the same events at their own pace.
  • Deep dive: 18 Trillion Messages a Day, one platform for urgent plus patient traffic.

Draw this shape for every pipeline you inherit this year. Interviews reward candidates who start from the log. Production rewards engineers who can replay the Tuesday nobody understood.

Interview room

Q1. Design event ingestion for a game with 18 trillion daily messages. Seen at: gaming plus high-scale loops.

Q2. A downstream team corrupts its state. How does the log save them? Seen at: backend plus data loops.

Truth stored. Now split it for speed. Lesson 02: Partitions →

Lesson 02 · Scale · Count partitions like money

Partitions Buy Concurrency

Assumption: one partition serves one consumer at a time. A payments team aimed at 10,000 transactions per second, then met a partition ceiling plus a 3,000-connection cap nobody had read. Partitions are concurrency you buy with coordination. Count them before traffic does.

1 topic part 1-2 part 3-4 part 5-6 cons 1-2 cons 3-4 cons 5-6 6x, capped
  • Partitions set the concurrency ceiling. More consumers than partitions means paid idleness.
  • Every partition taxes the broker. Assignment math plus metadata grow with the count, not with the data.
  • Deep dives: Razorpay Kafka plus Roblox tiers.

Size partitions from peak throughput divided by per-partition headroom, with a stated growth factor. Lesson 4 shows what happens when the count meets a rolling deploy. The ceiling is load-bearing. Read it twice.

Interview room

Q1. Size partitions for 60,000 events per second with 3x headroom. Seen at: Uber-style plus fintech loops.

Q2. Consumers idle while lag grows. What do you check first? Seen at: backend plus data loops.

Concurrency bought. Now share it across a group. Lesson 03: Groups →

Lesson 03 · Coordination · One group, many mouths

Groups Share the Load

Assumption: 120 consumers split 40,000 partitions evenly until one restarts. A consumer group is a contract: the coordinator assigns, members consume, generations mark agreement. The contract holds beautifully at rest. Every lesson after this one is about motion.

coordinator member 1-40 member 41-80 member 81-120 40k parts gen 12
  • One group id means shared fate. Every member rebalances together, which is the feature plus the hazard.
  • Generations are the agreement record. A group that never reaches a stable generation never consumes.
  • Deep dive: Kafka Rebalance Storm, the group that never settled.

Sketch your groups with member counts plus partition counts plus generation dashboards before the next deploy. The map from Lesson 1 now has owners. Owners restart. Restarts rejoin. Rejoins are Lesson 4.

Interview room

Q1. Explain consumer groups plus generations to a smart intern. Seen at: every infrastructure loop.

Q2. Two teams share one group id by accident. What breaks? Seen at: backend loops.

Load shared. Now watch a deploy pause it all. Lesson 04: Rebalance →

Lesson 04 · Failure · The world stops, eagerly

Rebalances Stop the World

Assumption: rolling deploys are zero-downtime. Roll one restart at a time across 120 consumers and each restart rejoins the group. Each rejoin triggers a full eager rebalance. Each rebalance takes long enough at 40,000 partitions that the next restart lands mid-rebalance. Modeled pause: 11 minutes. The safest rollout becomes a continuous stall.

restart 1 rebalance... restart 2! never settles triggers outpace completions, consumption stalls

When asked how deploys cause streaming outages, answer with generations. The group never reaches a stable one because triggers outpace completions. Fix the race, not the brokers.

Interview room

Q1. A rolling deploy pauses all consumers for 11 minutes. Explain the cascade. Seen at: backend plus SRE loops.

Q2. Brokers look healthy while consumption stalls. Where do you look? Seen at: SRE loops.

Storm modeled. Now remove the trigger. Lesson 05: Static ids →

Lesson 05 · Protocol · Boring fixes compose

Static Ids End Storms

Assumption: a restart is a stranger. Static membership rejects that premise. Each consumer restarts under a stable instance id plus keeps its assignment, so the rolling deploy becomes 120 uneventful restarts. Cooperative assignment softens the remaining blow: only partitions that must move, move. Nothing here is novel. That is the point.

stable id restart keeps parts no storm stable ids remove the trigger, cooperation softens the rest
  • Stable instance ids turn restarts into returns. The group stops meeting strangers every deploy.
  • Cooperative protocol turns pauses partial. Unaffected partitions keep consuming through transitions.
  • Deep dive: Kafka Consumer Rebalance Explained, the three settings that end the loop.

Audit your protocol defaults on the same schedule as capacity. Defaults are load-bearing architecture. The outage was three reasonable defaults composing at scale. The fix is five boring settings reviewed quarterly.

Interview room

Q1. Compare eager versus cooperative rebalancing at 40,000 partitions. Seen at: backend plus infra loops.

Q2. What breaks when two members share one static id? Seen at: Kafka-flavored loops.

Storms ended. Now keep events in order. Lesson 06: Key order →

Lesson 06 · Ordering · One key, one queue

Order Lives in the Key

Assumption: global order is purchasable. It is not. Ordering holds within a partition, so the key decides the queue. Hash by order id and one order never overtakes itself. Hash randomly and refunds precede charges. A chat platform learned the hot-key corollary: one popular channel concentrates 5,000 reads per second onto a single partition.

order events hash by key part 7: ordered hot key: full split plan
  • Key by what must stay ordered. Payments by order id, messages by channel, sessions by user.
  • Plan the hot-key split before the fire. Celebrity keys break naive hashing on event night.
  • Deep dive: Discord: 177 Nodes Too Many, hot partitions with receipts.

List your top ten keys by volume before launch, with a split plan for each. Capacity follows the hottest partition, never the average. The average is a rumor the hot key never heard.

Interview room

Q1. Keep per-user order for a chat app at 5,000 reads per second per channel. Seen at: Discord-style plus Meta-style loops.

Q2. One key takes 80 percent of traffic. Now what? Seen at: X-style plus high-scale loops.

Order kept. Now read the mail nobody wants. Lesson 07: Dead letters →

Lesson 07 · Failure · The queue for failures

Dead Letters Deserve Readers

Assumption: poison events are rare enough to ignore. They are not. Every message your consumer cannot process retries, fails, plus retries again, consuming the workers healthy events need. A dead-letter topic quarantines the poison after a bounded number of attempts. Every entry is a bug report from production. Teams that review them weekly find failures before customers do.

poison in retry x3 max dead letter weekly review bounded retries, then quarantine, then readers
  • Bound attempts before routing to dead letters. Unbounded poison retries are a self-inflicted stall.
  • Alert on dead-letter growth rate, not just depth. A rising slope is an incident writing its first draft.
  • Deep dive: 18 Trillion Messages a Day, patient lanes for patient problems.

Staff the dead-letter review like on-call, because it is one. Quarantine without readers is a morgue with extra steps. Readers turn poison into patches.

Interview room

Q1. Design poison-message handling for a payments pipeline. Seen at: fintech plus backend loops.

Q2. Dead letters grow 10x overnight. What happened? Seen at: SRE loops.

Poison quarantined. Now sign contracts with producers. Lesson 08: Schemas →

Lesson 08 · Evolution · Contracts, not hopes

Schemas Are Contracts

Assumption: producers plus consumers upgrade together. They never do. A field renamed on Friday breaks Saturday consumers that nobody owns anymore. A schema registry makes evolution a reviewed contract: backward plus forward compatibility checked at publish time, breaking changes rejected before they reach the log.

producer registry gate compat: log break: reject consumers safe
  • Gate publishes on compatibility. Add optional fields freely. Renames plus type changes ride a migration, not a Friday.
  • Version every event. Consumers read the version first, the payload second.
  • Deep dive: The Workspace That Ate the Warehouse, change capture across 480 shards.

Register your three highest-volume event types this week with compatibility rules written down. Ungated schemas are handshake deals with teams you have never met. The log remembers what the meeting forgot.

Interview room

Q1. Rename a field consumed by forty services without downtime. Seen at: platform plus backend loops.

Q2. A producer ships a breaking change at midnight. Who stops it? Seen at: SRE plus data loops.

Contracts signed. Now price keeping history. Lesson 09: Retention →

Lesson 09 · Storage · Hot plus warm plus cold

Retention Is Money

Assumption: retention is free because disks are cheap. Modeled reality: hot NVMe per gigabyte-month dwarfs object storage by an order of magnitude, plus every retained byte is served, replicated, plus rebalanced. One ingestion platform treats urgent matchmaking like patient analytics, then pays urgent prices for patient bytes. Tier by age. Serve hot, park cold.

hot: NVMe warm: 7 days cold: object save 10x age decides the tier, the tier decides the bill
  • Tier by access age, not by sentiment. Yesterday serves from NVMe. Last quarter serves from object storage.
  • Price retention per topic per tier monthly. Unpriced history grows until the invoice explains it.
  • Deep dive: 18 Trillion Messages a Day, urgent versus patient by design.

Follow the packet above from hot to cold, then tier your own oldest topic. Retention without tiers is a subscription to urgency. The bill hides in the boundary between hot and cold.

Interview room

Q1. Keep a year of events queryable without paying hot prices all year. Seen at: data plus platform loops.

Q2. Replay demand spikes on six-month-old data. What breaks? Seen at: backend loops.

History tiered. Now promise exactly-once honestly. Lesson 10: Exactly-once →

Lesson 10 · Correctness · A budget, not a feature

Exactly-Once Is a Budget

Assumption: the network delivers once. It delivers at least once, then apologizes with duplicates. One ledger moves $1.4 trillion without double-charging anyone. The trick is not a faster network. It is a key the client sends with every request, stored beside the result, so retries become safe by construction.

producer key per event seen: skip new: apply once, kept
  • Key every effect that money touches. Idempotency keys turn at-least-once networks into exactly-once effects.
  • Store the key beside the result atomically, or the check lies under concurrency.
  • Deep dive: The Ledger That Cannot Blink, double-entry plus keys at $1.4 trillion.

Notice what the diagram never shows. No step asks the network to behave. Reliability comes from assuming retries, duplicates plus reordering, then designing so none of them matter. Exactly-once is not a protocol setting. It is a budget line with a key attached.

Interview room

Q1. Guarantee exactly-once charging over an at-least-once queue. Seen at: Stripe-style plus fintech loops.

Q2. Dedup storage grows forever. When do keys expire? Seen at: backend loops.

Duplicates tamed. Now cap the retry wave itself. Lesson 11: Retry storm →

Lesson 11 · Amplification · The 6x invoice

Retries Without Jitter Are DDoS

Assumption: retries help. Modeled: eleven paused minutes at 60,000 events per second build roughly 39.6 million events of backlog while downstream timeouts fire retries into topics consumers cannot yet drain. At two extra retries per stalled call with no backoff, recovery load lands near six times normal before one new event arrives. Retries without jitter plus backoff are a self-inflicted DDoS with a deploy as the trigger.

stall 11m retry x2 6x load cap 1x backoff plus jitter plus caps, all three
  • Backoff plus jitter plus a hard stop. Miss any of the three and retries amplify instead of healing.
  • Price your policy with the retry storm calculator before the incident instead of after.
  • Deep dive: Kafka Rebalance Storm, the full recovery math.

Idempotency from Lesson 10 is what makes retries safe to attempt. Without it, backoff only spaces out the damage. The two lessons compose. Design them together, review them together.

Interview room

Q1. Convince me your retries cannot DDoS your own brokers. Seen at: Amazon-style loops.

Q2. Lag is extreme plus breakers open. What consumes first on recovery? Seen at: SRE loops.

Amplifier capped. Now slow the producer instead. Lesson 12: Backpressure →

Lesson 12 · Flow · Slow down to speed up

Backpressure Beats Buffering

Assumption: bigger buffers absorb bigger spikes. Buffers only postpone the reckoning while hiding the signal. A payments migration aimed at 10,000 transactions per second survived by separating urgent payments from patient analytics into lanes with independent workers. Backpressure tells producers to slow down. Buffering tells operators everything is fine until it is not.

flood in urgent lane patient lane slow signal steady out
  • Lane traffic by urgency before it shares workers. One surge must never starve everything.
  • Signal slowness upstream explicitly. Dropped plus delayed with notice beats silent plus lost.
  • Deep dive: Razorpay Kafka, lanes that cut latency 25x.

Separate your own urgent from patient traffic this quarter, then load-test the urgent lane alone. Queues that absorb spikes are designed. Queues that merely buffer are postponed incidents with dashboards.

Interview room

Q1. Traffic triples at 6 PM with mixed urgency. Walk through the flow. Seen at: Google-style loops.

Q2. When does backpressure hurt more than buffering? Seen at: backend loops.

Flow steadied. Now envelope the month. Lesson 13: Envelope →

Lesson 13 · Economics · Backlog plus invoice

The Envelope Math

Assumption: throughput is a technical metric. It is a financial one wearing a dashboard. Sixty thousand events per second sustained for a month is 155 billion events. At eleven paused minutes the backlog alone is 39.6 million. Run your shape through the envelope calculator before launch. Envelopes surprise only the unpriced.

60k eps 2.6M sec 155B events priced modeled envelope, labeled as estimates

Present envelopes as ranges with stated assumptions, never as single confident numbers. Precision without assumptions is theater. The bill hides in the boundary between the model and the month.

Interview room

Q1. Size a pipeline for 155 billion monthly events, out loud. Seen at: Uber-style plus data loops.

Q2. Retention doubles plus budget does not. What moves tiers? Seen at: FinOps-flavored loops.

Month priced. Now rehearse the interview. Lesson 14: Interview bench →

Lesson 14 · Interviews · Clarify plus estimate plus break

The Interview Bench

The panel asks for notifications for 100 million users. The candidate who clarifies scope, estimates out loud, plus names three ways the design breaks passes. The candidate who recites broker trivia does not. Clarify first. Estimate second. Break your own design third, before they do.

clarify estimate break it offer
  • Clarify scope before drawing. Urgent versus patient, ordered versus not, retained for how long.
  • Estimate partitions plus consumers plus backlog out loud with the RPS envelope calculator. Panels score math over adjectives.
  • Break your own design unprompted. Hot keys, rolling storms, plus poison events, with fixes attached.

Rehearse three answers this week: a partition sizing with headroom math, a rebalance story with protocol fixes, plus an exactly-once design with key expiry. Bring diagrams. Engineers draw. Tourists describe.

Interview room

Q1. Design notifications for 100 million users end to end in 45 minutes. Seen at: Meta-style plus WhatsApp-style loops.

Q2. Tell me about a pipeline failure you caused plus fixed. Seen at: behavioral plus backend loops.

Answers rehearsed. Now run one pipeline for real. Lesson 15: Run pipeline →

Lesson 15 · Capstone · One topic, three consumers

Run One Pipeline

Start small on purpose. One topic, three consumers, one dashboard. Produce a million events, kill a consumer mid-run, watch the rebalance, read the dead letters, replay the log, price the month. A pipeline you have broken on purpose is the only kind you understand.

1 topic 3 consumers kill drill priced, replayed break it on purpose, then price the month
  • Run the fourteen drills from lessons 1 through 13 against one real topic. Gaps are visible at a glance.
  • Kill one consumer mid-run plus measure the rebalance. Untested failover is a rumor with runbooks.
  • Close with the invoice. Events per month times modeled unit price, labeled as estimates.

Course complete. Interview the whole system next: System Design Course, then browse all courses plus all case studies. New events arrive constantly. So will new backlogs.

The log remembers. The invoice explains. You just built the study guide. Earn the grade.

Sources plus method

This course teaches from published Buildopsy postmortems linked under each lesson: Kafka rebalance cascades, Razorpay partition limits, Roblox tiering, Discord hot partitions, Stripe idempotency, Notion change capture, plus the retry storm model. Incident details follow those accounts. Cost framing is modeled from public list prices.