Hiring BarSupport

Design with Apache Kafka — Staff-Level Technology Guide

Technology guide38 min read7 diagrams

Why This Matters#

Kafka is not a message queue. It is a replicated, partitioned commit log that happens to be usable as one — and nearly every Kafka mistake in an interview comes from forgetting that. Queues delete a message when it's consumed; Kafka keeps it for a retention window and lets any number of consumers read it at their own offset. Queues give you per-message acknowledgement; Kafka gives you a single committed offset per partition. Queues scale consumers freely; Kafka caps parallelism at the partition count. Every one of those differences is a design decision you must make out loud.

That is why Kafka appears in more Staff system design loops than any other technology. Ad click aggregation, notification fan-out, CDC pipelines, event sourcing, payment state machines, activity feeds, log shipping — the interviewer expects the log. The L5 candidate draws a box labeled "Kafka" between two services. The L6 candidate says "topic orders.v1, 48 partitions keyed by order_id for per-order ordering, RF=3, min.insync.replicas=2, acks=all, idempotent producer, and consumers commit offsets after the side effect with an idempotency key so redelivery is safe." The L7 candidate asks who owns the event schema — because in year three, the topic contract is the API that 40 teams depend on and nobody can change.

The L5 → L6 gap is not knowing what ISR stands for. It is knowing that the partition key is the design — it decides ordering, parallelism, hot spots, and the cost of ever changing your mind.

The L5 → L6 → L7 Contrast#

BehaviorSenior (L5)Staff (L6)Principal (L7)
First move"Put Kafka between the services""What must be ordered relative to what? That's my partition key.""Is this an internal pipe or a published event contract other teams will build on? That decides schema governance and ownership."
Delivery semantics"Kafka has exactly-once""At-least-once delivery + idempotent consumers; EOS only for Kafka-to-Kafka read-process-write."Sets the org standard: every consumer must be idempotent; EOS is an optimization, never a correctness crutch
Durability"RF=3"acks=all, min.insync.replicas=2, unclean election off — and says what gets rejected when 2 brokers diePrices durability: cross-AZ replication is often the largest line on the Kafka bill
Scaling"Add partitions""Partition count is a one-way door for keyed topics — I'll size for 3× peak now."Decides cluster topology for the org: shared multi-tenant vs per-domain clusters, with quotas
Failure"Consumers retry"Lag-based alerting, DLQ with an owner, poison-message handling, rebalance stormsDesigns the org's replay/backfill posture and who is allowed to reset offsets on shared topics
Ownership"Platform team runs Kafka"Platform owns brokers; producing team owns topic schema and DLQWrites the event contract standard: schema registry, compatibility mode, deprecation window
Why "First move" separates levels

"Add Kafka for decoupling" is not wrong. It is incomplete in a way that leaks into every later decision. Kafka guarantees order only within a partition, and the partition is chosen by hashing the key. So "what must be ordered?" determines the key; the key determines hot spots (one celebrity user, one giant merchant); and the partition count determines maximum consumer parallelism. A Staff candidate resolves this in the first minute: "Events for the same order must be processed in order; events across orders need not be. Key = order_id." Everything downstream becomes defensible.

Why "Delivery semantics" separates levels

Kafka's exactly-once semantics (idempotent producer + transactions + read_committed) are real, but they cover Kafka → processor → Kafka. The moment your consumer writes to Postgres, calls Stripe, or sends an email, you are back to at-least-once and your consumer must be idempotent. The Senior answer quotes the feature; the Staff answer draws its boundary; the Principal answer makes idempotent consumers a platform rule so no team rediscovers the boundary during an incident.

The 60-Second Pitch#

"I'd use Kafka as the durable event backbone. It gives us an ordered, replayable log per partition, sustains hundreds of MB/s per cluster on commodity brokers, and lets multiple consumer groups — fulfillment, analytics, search indexing — read the same events independently at their own pace. I'd key the topic by order_id so per-order events stay ordered, run RF=3 with min.insync.replicas=2 and acks=all so an acknowledged write survives a broker loss, and make every consumer idempotent because delivery is at-least-once to anything outside Kafka. The retention window — 7 days by default — is our replay budget for bugs and backfills. I would not use Kafka for per-message task queues with individual retries and delays; that's SQS or a job scheduler."

The Three Intents#

IntentConstraintStrategyFailure ModeCorrectness Bar
Event backbone / pub-subMany independent consumers, replay, schema evolutionTopic per domain event, keyed, RF=3, 7d+ retention, schema registrySchema break takes down 12 downstream teamsAt-least-once, ordered per key
Stream processing input (CDC, analytics, aggregation)Throughput, ordering per entity, exactly-once aggregatesDebezium CDC or producer events → Flink/Kafka Streams with EOSConsumer lag grows unbounded; reprocessing double-countsExactly-once within the pipeline
Work queue (jobs, notifications, tasks)Per-message retry, delays, priority, independent failureKafka can do it with retry topics + DLQ, but SQS/RabbitMQ/share groups fit betterOne poison message blocks a partition (head-of-line)Each task done at least once, individually

🎯 Staff Move: "I'll treat Kafka as the event backbone, not a task queue. If a single slow or poison message must not block the messages behind it, that's a queue-shaped requirement and I'd put SQS or a job scheduler behind the Kafka consumer, rather than bend partition semantics into per-message retries."

The Staff Positions#

PositionRationale
The partition key is the designIt fixes ordering scope, parallelism, and hot-spot risk. Decide it first, out loud.
acks=all + min.insync.replicas=2 + RF=3 by defaultThe only config where an acknowledged write survives a single broker loss without data loss.
At-least-once + idempotent consumersEOS doesn't extend to external side effects; idempotency keys do.
Over-partition keyed topics at creationAdding partitions remaps keys and breaks ordering; size for 2–3× projected peak.
Lag is the SLO, not broker CPUConsumer lag in seconds is what users feel.
Every DLQ has a named owner and a replay toolA DLQ nobody reads is a silent data-loss mechanism.
Schemas are contracts with compatibility checksSchema registry with BACKWARD (or FULL) compatibility enforced at produce time.

Architecture & Internals#

Only five internals change design decisions: partitions and the log, replication and ISR, the controller (KRaft), consumer groups and offsets, and retention/compaction.

The Log, Partitions, and Segments#

A topic is a set of partitions. Each partition is an append-only, totally ordered sequence of records identified by a monotonically increasing offset. On disk, a partition is a directory of segment files (default 1GB each) plus sparse offset and time indexes.

  • Writes are sequential appends — this is why a broker can sustain hundreds of MB/s on spinning disk or modest SSDs.
  • Reads are sequential too, and hot reads are served from the OS page cache, often via zero-copy sendfile. A consumer that is caught up reads from memory; a consumer 3 days behind reads from disk and can evict the page cache for everyone else.
  • Retention deletes whole segments by age (retention.ms, default 7 days) or size (retention.bytes). There is no per-message delete.
Diagram: The Log, Partitions, and Segments

The partitioner: the default producer partitioner computes murmur2(key) mod num_partitions. Same key → same partition → ordered. Null key → sticky partitioning across partitions (batches fill one partition, then move on), which maximizes batching but gives no ordering.

Replication, ISR, and What acks Actually Means#

Each partition has one leader and RF−1 followers on different brokers (and, with rack awareness, different AZs). Followers fetch from the leader. The In-Sync Replica set (ISR) is the set of replicas caught up within replica.lag.time.max.ms (default 30s).

Producer acksAcked whenSurvives leader crash?Latency
0Sent to socketNo — may lose anythingLowest
1Leader wrote to its logNo — lost if leader dies before followers fetchLow
all (-1)All current ISR members have itYes, if ISR ≥ min.insync.replicas+1 replication RTT (~1–5ms intra-region)

The durability formula:

RF = 3, min.insync.replicas = 2, acks = all
  -> write acked only when >= 2 replicas have it
  -> survives loss of any 1 broker with zero acknowledged-data loss
  -> if ISR shrinks to 1: producers get NotEnoughReplicas -> writes REJECTED
     (availability sacrificed for durability — this is the point)

unclean.leader.election.enable = false  (default)
  -> an out-of-sync replica is never elected leader
  -> partition stays offline rather than silently truncating acked data

🎯 Staff Insight: "With RF=3 and min ISR 2, losing two brokers that share a partition makes that partition reject writes. That's the correct behavior for orders. For clickstream I might accept min.insync.replicas=1 and keep writing — but I'd say that out loud and get the analytics owner to sign off that losing a few seconds of clicks is fine."

The Controller — From ZooKeeper to KRaft#

Historically, Kafka stored cluster metadata (brokers, topics, partition leaders, ISR) in ZooKeeper, and one broker acted as controller. KRaft (KIP-500) replaces ZooKeeper with a Raft-replicated metadata log inside Kafka itself, run by a quorum of 3 or 5 controller nodes. KRaft became production-ready in Kafka 3.3, and Kafka 4.0 removed ZooKeeper mode entirely.

Why it matters in design:

  • Partition ceiling. ZooKeeper-era guidance kept clusters to roughly ~4K partitions per broker and ~200K per cluster, because controller failover had to reload all metadata from ZooKeeper. KRaft's event-sourced metadata log makes failover near-instant and raises the ceiling substantially.
  • One fewer system to operate — no separate ZooKeeper ensemble, no ZooKeeper-specific failure modes.
  • Controller quorum sizing follows Raft: 3 controllers tolerate 1 failure, 5 tolerate 2.

Consumer Groups, Offsets, and Rebalancing#

A consumer group shares the partitions of a topic: each partition is assigned to exactly one consumer in the group at a time. Therefore max useful consumers = partition count. A 12-partition topic with 20 consumers leaves 8 idle.

Progress is a committed offset per (group, partition), stored in the internal __consumer_offsets compacted topic. Committing offset N means "everything before N is done."

Diagram: Consumer Groups, Offsets, and Rebalancing

Rebalancing reassigns partitions when consumers join, leave, or miss heartbeats. Rebalance-related configs every Staff candidate should know:

ConfigDefaultWhy it matters
session.timeout.ms45s (3.0+)Consumer considered dead after missing heartbeats this long
max.poll.interval.ms300sIf processing one batch takes longer, the consumer is kicked → rebalance → redelivery → loop
max.poll.records500Lower it when per-record work is slow
partition.assignment.strategyCooperative sticky (modern clients)Incremental rebalances move only the partitions that must move
group.instance.idunsetStatic membership: a restarting pod reclaims its partitions without a rebalance

Kafka 4.0 ships the new consumer group protocol (KIP-848), which moves assignment to the broker and makes rebalances incremental by default — the fix for the "stop-the-world rebalance" storms that plagued large consumer groups.

Retention, Compaction, and Tiered Storage#

ModeConfigKeepsUse for
Delete (time/size)cleanup.policy=delete, retention.ms=604800000 (7d)Every record until the segment ages outEvent streams, logs, clickstream
Compactcleanup.policy=compactLatest record per key forever; null value = tombstone deletes the keyChangelogs, CDC tables, current-state topics, __consumer_offsets
Compact + deletecleanup.policy=compact,deleteLatest per key, but only within the retention windowBounded state caches
Tiered storageKIP-405, remote.storage.enable=trueRecent segments on broker disk; older segments in S3/GCSMonths of replay without buying broker disks

Retention is your replay budget. If a consumer bug corrupts data for 4 days and retention is 3 days, you cannot fix it from Kafka. Tiered storage changes the economics: broker-attached SSD at ~$0.08–0.10/GB-month × RF=3 vs object storage at ~$0.023/GB-month × 1 copy — roughly a 10× cost difference per retained GB.

🎯 Staff Insight: "Retention isn't a storage setting, it's a recovery-time objective. I'd set it to at least the longest plausible time-to-detect a consumer bug plus the time to ship a fix — for most teams that's 7 days minimum, 30 with tiered storage for anything feeding a ledger or search index."


Data Modeling / Core Usage — "The Entire Game"#

In Cassandra the game is the partition key. In Kafka it is the same game with worse consequences: the message key picks the partition, the partition is the unit of ordering and parallelism, and partition count is nearly impossible to change for keyed topics.

Step 1: Decide the Ordering Scope#

Ask: "Which events must be seen in the order they happened, by the same consumer?"

DomainMust be ordered togetherKeyHot-key risk
E-commerce ordersAll events of one order (created → paid → shipped)order_idLow — orders are small
Bank ledgerAll postings to one accountaccount_idHigh — a payroll account posts 100K/day
ChatMessages in one conversationconversation_idMedium — large public channels
ClickstreamNothing strictlynull or session_idNone with null key
CDC from PostgresAll changes to one rowprimary keyLow, unless one row is a counter
IoTReadings of one devicedevice_idLow per device, high if one gateway fronts 10K devices

If nothing needs ordering, use a null key — sticky partitioning gives the best batching and even load.

Step 2: Size the Partition Count#

partitions >= max( target_throughput / per_partition_producer_throughput,
                   target_throughput / per_partition_consumer_throughput,
                   max_consumer_parallelism_needed )

Rule-of-thumb per-partition throughput: ~5-10 MB/s producer, and whatever
your consumer can process (often far less — a consumer doing a 5ms DB write
per record handles ~200 rec/s per thread).

Example: 50K events/s, 1KB each = 50 MB/s
  producer side:   50 / 10       = 5 partitions
  consumer side:   50,000 / 1,500 rec/s per consumer = 34 consumers
  headroom 2-3x:   ~100 partitions -> pick 96 or 128

The consumer side almost always dominates. Candidates size partitions for broker throughput and then discover their consumers do a database write per record.

Why partition count is a one-way door for keyed topics: hash(key) mod N changes when N changes. Going from 48 to 64 partitions sends order_123 to a new partition while its older events sit on the old one; a consumer can now see "shipped" before "paid." The fix is a new topic + migration, not an ALTER.

🎯 Staff Move: "I'll create the topic with 96 partitions even though today's load needs 20. Partitions are cheap until you're in the tens of thousands per cluster; re-keying a live ordered topic is a quarter-long migration."

Step 3: Design the Event (the Contract)#

{
  "event_id": "01J8Z6Q2V3X9KX7G2M4T3Y8B1C",
  "event_type": "order.paid",
  "schema_version": 3,
  "occurred_at": "2026-03-14T09:26:53.589Z",
  "producer": "payments-svc",
  "key": "ord_7Hq2",
  "trace_id": "4bf92f3577b34da6a3ce929d0e0e4736",
  "data": { "order_id": "ord_7Hq2", "amount_minor": 4599, "currency": "USD" }
}
  • event_id — globally unique (ULID/UUIDv7); the consumer's idempotency key.
  • occurred_at — event time, not processing time; stream processors window on it.
  • Schema registry (Avro/Protobuf/JSON Schema) with BACKWARD compatibility: new consumers can read old events; fields are added as optional, never renamed or retyped.
  • Fat vs thin events: thin (order.paid {order_id}) forces every consumer to call back to the order service — coupling and a thundering herd. Fat (full state) decouples but makes the event a larger contract. Staff default: fat enough that the top 3 consumers never call back.

Step 4: Topic Naming and Topology#

<domain>.<entity>.<event-or-version>      orders.order.v1
<domain>.<entity>.changelog              inventory.sku.changelog   (compacted)
<consumer-group>.<topic>.retry.<n>       fulfillment.orders.v1.retry.1
<consumer-group>.<topic>.dlq             fulfillment.orders.v1.dlq

One topic per event stream with one ordering scope and one schema family. Don't create a topic per customer (topic explosion) and don't put every event in one events topic (no independent retention, ACLs, or schema).

Producer Configuration That Matters#

acks=all
enable.idempotence=true            # default true since Kafka 3.0; dedups retries per partition
max.in.flight.requests.per.connection=5   # safe with idempotence; ordering preserved
retries=2147483647
delivery.timeout.ms=120000
linger.ms=10                       # wait up to 10ms to fill a batch
batch.size=131072                  # 128KB batches
compression.type=zstd              # 3-5x on JSON; lz4 if CPU-bound

linger.ms is the latency-for-throughput knob: 0ms sends immediately with small batches; 10–20ms typically doubles or triples throughput and cuts broker request rate by an order of magnitude.

The Transactional Outbox — Getting Events Out of Your Database#

The dual-write problem: "Update the orders table, then publish to Kafka" — if the process dies between the two, the DB and the event stream disagree forever.

BEGIN;
  UPDATE orders SET status = 'PAID' WHERE id = 'ord_7Hq2';
  INSERT INTO outbox (event_id, topic, key, payload)
  VALUES ('01J8Z6...', 'orders.order.v1', 'ord_7Hq2', '{...}');
COMMIT;
-- Debezium tails the WAL, publishes outbox rows to Kafka (at-least-once).
-- Consumers dedupe on event_id.

The outbox (or CDC directly from the table) is the Staff default for any event that must agree with a database. See PostgreSQL.


The Tunable Tradeoff — Delivery Semantics#

SemanticsHowCostWhere it's correct
At-most-onceCommit offset before processing; acks=0/1Loses data on crashMetrics sampling, best-effort telemetry
At-least-onceProcess, then commit; acks=all, retries onDuplicates on crash/rebalanceDefault for everything, with idempotent consumers
Exactly-once (Kafka-to-Kafka)Idempotent producer + transactions + isolation.level=read_committed; offsets committed inside the transaction~3–10% throughput overhead; latency = transaction commit intervalStream processing (Kafka Streams, Flink) producing to Kafka
Effectively-once (end-to-end)At-least-once + idempotent sink (upsert by key, dedupe table, conditional write)A dedupe store and a key designAnything with external side effects
Exactly-once inside Kafka:
  consume(offsets) -> process -> produce(outputs) + sendOffsetsToTransaction
  -> commitTransaction()   (atomic: outputs visible AND offsets advanced, or neither)

The boundary:
  consume -> charge credit card -> commit offset
  crash after charge, before commit -> redelivered -> charged twice
  EOS cannot help. Only an idempotency key at the payment provider can.

🎯 Staff Move: "I'll design for at-least-once and make the sink idempotent — upsert keyed on event_id, or a processed-events table in the same transaction as the side effect. I'd only turn on Kafka transactions for the Flink aggregation that writes back to Kafka, because that's the one hop where EOS actually closes the gap."

Who Pays for Each Choice#

ChoiceWhat WorksWhat BreaksWho Pays
acks=1 for throughput~30–50% lower p99 produce latencyAcked data lost on leader crashDownstream teams discovering gaps; finance if it's orders
min.insync.replicas=1Writes continue with 2 brokers downSingle copy of acked dataWhoever explains the loss in the post-mortem
EOS everywhereSimpler reasoning inside KafkaThroughput overhead; false confidence at external sinksConsumer teams who still double-charge
Idempotent consumersCorrect under any redeliveryDedupe store to operate; key disciplineEach consuming team (correctly)
Long retention (30d+)Replay and backfill after bugs3× disk per GB without tiered storagePlatform budget

Anti-Patterns — What Kills Kafka Deployments#

1. Kafka as a Per-Message Task Queue#

A slow or failing message at offset N blocks N+1…N+10,000 on the same partition (head-of-line blocking). Teams bolt on retry topics, delay topics, and in-memory skip lists, and rebuild a worse SQS. If you need per-message visibility timeouts, delays, and priorities, use a queue — or Kafka 4.x share groups (KIP-932, "queues for Kafka", early access), with eyes open about maturity.

2. Hot Partitions from a Skewed Key#

Keying ledger postings by merchant_id when one merchant is 30% of volume. One partition, one consumer, one broker's disk at 100% while 95 partitions idle. Fix: a finer key (merchant_id:account_id), or a salted key (merchant_id#(hash % 8)) when strict per-merchant ordering is not needed, with an aggregation step downstream.

3. Adding Partitions to a Keyed Topic in Production#

Breaks per-key ordering silently. Nothing errors. Consumers process "cancelled" before "created" for a fraction of keys during the transition. Fix: new topic with the new count, dual-produce, drain, cut over.

4. Topic-per-Tenant or Topic-per-User#

10K customers → 10K topics × 12 partitions × RF 3 = 360K partition replicas. Controller metadata, file handles, and recovery time all scale with partitions. Fix: shared topic keyed by tenant, with quotas.

5. Processing Longer Than max.poll.interval.ms#

A consumer that calls a slow API for each of 500 records exceeds 300s, gets evicted, triggers a rebalance, its partitions go to another consumer that hits the same slow records — a rebalance storm with 0 progress and duplicates everywhere. Fix: lower max.poll.records, move slow work to an async pool with bounded concurrency, or pause partitions while processing.

6. Unowned DLQ#

Messages fail, go to *.dlq, and nobody looks. After 7 days retention deletes them. It is data loss with extra steps. Fix: DLQ depth alert → named owning team, a replay tool, and a weekly review.

7. Consumers Reading From the Beginning in Production#

A new consumer group with auto.offset.reset=earliest on a 30-day, 20TB topic floods broker disks and evicts the page cache that the latency-sensitive consumers rely on. Fix: quotas on consumer fetch bytes, backfills from tiered storage or a separate replica cluster.

8. Using Kafka as the System of Record Without Compaction or Backup#

"The log is the database" works — with compacted topics, infinite retention, and a tested restore. Without those it's a 7-day buffer that someone treated as permanent.


The Technology Landscape — Head-to-Head Comparison#

DimensionKafkaRedpandaApache PulsarAWS KinesisSQS / RabbitMQ
ModelPartitioned logPartitioned log (Kafka API)Segmented log over BookKeeper, compute/storage splitSharded logQueue (per-message ack)
OrderingPer partitionPer partitionPer partition / key-sharedPer shardFIFO queues only (SQS FIFO ~300–3,000 msg/s per group w/ batching)
ThroughputHundreds of MB/s–GB/s per clusterSimilar or higher per core (C++, thread-per-core)High; brokers stateless1 MB/s in, 2 MB/s out per shardSQS standard: effectively unlimited; RabbitMQ ~20–50K msg/s per node
ReplayYes, by offset/timeYesYesYes, 24h default up to 365dNo (once acked, gone)
Consumer scaling≤ partitions per group≤ partitionsShared subscriptions scale past partitions≤ shards (or enhanced fan-out)Unlimited competing consumers
Ops burdenMedium–high (self-host); low on MSK/ConfluentLower — single binary, no JVMHigh — Pulsar + BookKeeper + metadata storeZeroZero (SQS) / medium (RabbitMQ)
Pick whenDefault event backbone, stream processing, ecosystem (Connect, Debezium, Flink)Kafka semantics with lower tail latency and simpler opsMulti-tenancy, geo-replication, huge topic countsAWS-native, modest scale, zero opsTask queues, per-message retry/delay

Emerging: S3-native "diskless" Kafka-compatible systems (WarpStream, and diskless-topic proposals in upstream Kafka) trade ~100–500ms produce latency for eliminating broker disks and cross-AZ replication cost. That's the right trade for logs and analytics, the wrong one for a payments state machine.


Patterns#

Pattern 1: Event Backbone with Independent Consumer Groups#

Diagram: Pattern 1: Event Backbone with Independent Consumer Groups

Each group has its own lag SLO — analytics at 15 minutes behind is fine; fulfillment at 15 minutes behind is an incident. Alert per group, not per topic.

Pattern 2: CDC + Outbox (Database → Kafka)#

Debezium reads the PostgreSQL WAL / MySQL binlog and emits row changes to Kafka, keyed by primary key. Use for search indexing, cache invalidation, and data-lake ingestion without dual writes. Watch: replication slot growth on the source DB if the connector stalls — a stuck Debezium connector can fill the Postgres disk.

Pattern 3: Retry Topics + DLQ (Non-Blocking Retries)#

Diagram: Pattern 3: Retry Topics + DLQ (Non-Blocking Retries)

Unblocks the main partition, at the cost of losing per-key ordering for retried messages. Acceptable for notifications; unacceptable for a ledger — there, you block and page instead.

Pattern 4: Compacted Topic as a Replicated Table#

A compacted inventory.sku.changelog holds the latest state per SKU. Services bootstrap a local cache (RocksDB/in-memory) by reading from offset 0, then follow the tail. This is Kafka Streams' KTable and the basis of event-sourced read models.

Pattern 5: Stream Processing (Kafka → Flink/Kafka Streams → Kafka)#

Windowed aggregation, joins, enrichment, with exactly-once via transactions. See Flink & Stream Processing and the Stream Processing case study.

Pattern 6: Multi-Region — Active/Passive Mirroring#

MirrorMaker 2 (or Confluent Cluster Linking) replicates topics asynchronously to a DR region. RPO = replication lag (seconds to minutes); consumer offsets must be translated because offsets differ between clusters. Active/active requires region-prefixed topics (us.orders, eu.orders) and consumers that read both — and a conflict story for anything keyed globally.


Scaling#

The Numbers#

ResourceRough limit / sizingNote
Per-partition produce throughput~5–10 MB/s comfortablePlanning number, not a hard cap
Broker throughput~100–300 MB/s in on modern instancesNetwork and disk bound; replication multiplies traffic by RF
Partitions per broker~4K (ZooKeeper-era guidance)KRaft raises cluster-wide ceilings substantially
Partitions per cluster~200K (ZooKeeper-era)More partitions = slower recovery and more open files
Produce latency p99 (acks=all, same region)~5–20ms+linger.ms
End-to-end latency p99~10–50ms typicalDominated by linger + consumer poll loop
Message sizeDefault max ~1MBPut large payloads in S3; send a pointer
Consumer group size≤ partition countExtra consumers sit idle

Replication Traffic Math (the Hidden Cost)#

Ingress 100 MB/s, RF=3, consumers read 3x (three groups)
  broker write:      100 MB/s leader + 200 MB/s follower fetch = 300 MB/s disk writes
  network out:       200 MB/s replication + 300 MB/s consumers   = 500 MB/s
  cross-AZ (3 AZs):  ~2/3 of producer traffic + 200 MB/s replication
                     + ~2/3 of consumer traffic (unless fetch-from-follower)
  at ~$0.02/GB cross-AZ: 100 MB/s ~= 8.6 TB/day ingress
                     -> tens of TB/day cross-AZ -> $500-1,000+/day

Mitigations: rack-aware fetch-from-follower (KIP-392) so consumers read from a replica in their own AZ, producer-side compression (zstd cuts bytes 3–5× for JSON), and larger batches.

Scaling Moves in Order#

  1. Tune clients — batching, compression, linger.ms. Often 2–3× for free.
  2. Scale consumers up to the partition count.
  3. Add brokers + reassign partitions (Cruise Control automates rebalancing). Reassignment copies data; throttle it (leader.replication.throttled.rate) or it saturates the network.
  4. Add partitions — only for unkeyed topics, or via topic migration for keyed ones.
  5. Split clusters by domain or criticality once one cluster's blast radius is too large.

Failure Modes & Recovery#

1. Consumer Lag Spiral#

  • Symptom: Lag grows linearly; downstream data minutes-to-hours stale.
  • Root cause: Consumer throughput < produce rate — slow sink, a downstream DB degraded, or a traffic spike beyond partition parallelism.
  • Detection: kafka_consumergroup_lag (records) and lag in seconds (via Burrow or records-lag-max + timestamp); alert when time-lag > SLO for 5 min.
  • Fix: Scale consumers to partition count; batch sink writes; shed non-critical consumers.
  • Prevention: Size partitions for 3× peak consumer parallelism; load test the sink, not just Kafka.

2. Rebalance Storm#

  • Symptom: Group constantly rebalancing; throughput near zero; duplicate processing.
  • Root cause: Processing time > max.poll.interval.ms, rolling deploys of 200 pods each triggering an eager rebalance, or GC pauses > session timeout.
  • Detection: kafka_consumer_coordinator_rebalance_total rate; join-rate; commit-rate dropping to 0.
  • Fix: Lower max.poll.records; cooperative-sticky assignor; static membership (group.instance.id).
  • Prevention: KIP-848 protocol on 4.x; bound per-record work; deploy with surge limits.

3. Under-Replicated Partitions → Write Rejection#

  • Symptom: Producers get NotEnoughReplicasException; order creation fails.
  • Root cause: Two brokers down or a slow disk pushing followers out of ISR, so ISR < min.insync.replicas.
  • Detection: kafka_server_replicamanager_underreplicatedpartitions > 0; UnderMinIsrPartitionCount > 0 (page).
  • Fix: Restore brokers; replace the failing disk; do not enable unclean leader election on money topics.
  • Prevention: Rack-aware placement across 3 AZs; never do rolling restarts while URP > 0.

4. Poison Message / Head-of-Line Block#

  • Symptom: One partition's lag climbs while others are at zero; consumer logs repeat the same exception.
  • Root cause: A malformed or schema-incompatible record the consumer can't process and keeps retrying.
  • Detection: Per-partition lag skew; records-consumed-rate = 0 for one partition; error log rate.
  • Fix: Route to DLQ after N attempts, skip, alert the owner.
  • Prevention: Schema registry compatibility at produce time; consumer deserialization errors → DLQ, never infinite retry.

5. Disk Full on a Broker#

  • Symptom: Broker goes offline; partitions it leads fail over; if widespread, producers see errors.
  • Root cause: Retention sized by time not bytes plus a traffic spike, or a stuck compaction, or a new high-volume producer.
  • Detection: kafka_log_log_size per broker; disk usage > 75% alert; per-client produce-byte quotas exceeded.
  • Fix: Temporarily lower retention.ms on the largest topics; add brokers and reassign.
  • Prevention: retention.bytes caps, per-client produce quotas, tiered storage, capacity alert at 60%.
Diagram: 5. Disk Full on a Broker

Operational Reality Matrix#

FailureDetection SignalBlast RadiusMitigationOwner
Lag spiralTime-lag > SLO 5mOne consumer group's downstreamScale consumers, shed loadConsuming team
Rebalance stormRebalance rate, commit-rate 0One consumer groupStatic membership, lower batchConsuming team
ISR < minUnderMinIsrPartitionCountEvery producer to affected partitionsRestore brokersKafka platform on-call
Poison messagePer-partition lag skewOne partition of one groupDLQ + skipConsuming team; producer fixes schema
Disk fullDisk > 75%All partitions led by that brokerLower retention, reassignPlatform
Schema breakDeserialization error spike across groupsEvery consumer of the topicRoll back producerProducing team
Controller quorum lossNo active controller metricWhole cluster metadata frozenRestore quorum (3→ tolerate 1)Platform

When to Use vs. Alternatives#

NeedPickWhy
Event backbone with many consumers and replayKafka (MSK/Confluent if you don't want to run it)Log semantics, ecosystem, throughput
Kafka semantics, fewer ops, lower tail latencyRedpandaNo JVM, no separate controller ensemble
Per-message retries, delays, visibility timeoutsSQS / RabbitMQQueue semantics; no head-of-line blocking
Small scale on AWS, zero opsKinesis or SQS + SNSManaged; per-shard pricing is fine below ~50 shards
Thousands of tenants, native geo-replicationPulsarTopic-count scale, built-in multi-region
Request/response between servicesgRPC/HTTPKafka is not an RPC bus — no reply path, seconds of worst-case latency
Cheap log/analytics firehose, latency tolerantS3-native Kafka-compatibleNo broker disks, no cross-AZ replication bill

When NOT to Use Kafka#

  • Low volume (< ~1K msg/s) and one consumer. A Postgres table with SELECT … FOR UPDATE SKIP LOCKED or SQS is simpler and cheaper.
  • You need per-message priority or delay. Kafka has neither.
  • Synchronous user-facing request paths that need a reply in < 50ms.
  • Large payloads (> 1MB) — store in object storage, send a pointer.
  • The team cannot staff an on-call for it — then buy it managed or don't adopt it.
Diagram: When NOT to Use Kafka

Operational Concerns#

What the On-Call Actually Does#

  1. Watches UnderMinIsrPartitionCount and offline partitions — the only broker metrics that mean "producers are failing right now."
  2. Watches lag in seconds per consumer group against that group's SLO, and routes the page to the consuming team.
  3. Throttles reassignments so rebalancing data after adding brokers doesn't saturate the network (leader/follower.replication.throttled.rate).
  4. Performs rolling restarts one broker at a time, waiting for URP = 0 between each. A 30-broker cluster roll takes hours; that's normal.
  5. Handles offset resets — the most dangerous routine operation. Resetting a shared group to earliest replays days of side effects. Require a second approver.
  6. Manages quotas (producer_byte_rate, consumer_byte_rate, request percentage per client ID) so one team's backfill can't starve everyone.

Key Metrics & Alerts#

MetricHealthyAlert
UnderMinIsrPartitionCount0> 0 for 1m (page)
OfflinePartitionsCount0> 0 (page)
UnderReplicatedPartitions0> 0 for 10m
Consumer lag (seconds) per group< SLO> SLO for 5m (page consuming team)
Produce request latency p99< 20ms> 100ms
Request handler idle ratio> 30%< 20%
Disk usage per broker< 60%> 75%
Active controller countexactly 1≠ 1
DLQ depth per topic0 growthAny growth (ticket to owner)

Upgrades and Config Changes#

  • Client compatibility: brokers support older clients; upgrade brokers first, clients later.
  • Topic config changes (retention.ms, min.insync.replicas) are live — and therefore dangerous. Put topic configs in Git (Terraform / a topic operator) with review.
  • Schema changes go through the registry's compatibility check in CI, before the producer deploys.

Interview Application — Staff-Level Plays#

Which Case Studies Use Kafka#

Case StudyHow Kafka Is UsedKey Pattern
Message QueuesThe log-vs-queue decision itselfPartition = ordering + parallelism unit
Stream ProcessingDurable input and output for windowed aggregationEOS Kafka → Flink → Kafka
Notification SystemFan-out from domain events to channel workersRetry topics + DLQ, per-channel consumer groups
Payment ProcessingPayment state events, ledger postingsOutbox + key by payment/account, idempotent consumers
News FeedPost events driving fan-out-on-writeKey by author; celebrity handling
Search IndexingCDC feed into the indexerDebezium, compacted changelog
Metrics & MonitoringBuffer between agents and the TSDBReplay after TSDB outage

Every System Design Question Has a Kafka Moment#

  • URL shortener: "Click events go to Kafka with a null key for even spread; a Flink job rolls them up per link per minute. The redirect path never waits on Kafka."
  • Chat: "Kafka keyed by conversation_id gives per-conversation ordering for persistence and fan-out — but the real-time delivery path is WebSockets, not Kafka."
  • Ride hailing: "Driver location pings go to Kafka keyed by geo cell so the matching consumer for a cell sees them in order; 1M drivers × 1 ping/4s = 250K msg/s — ~64 partitions with headroom."
  • Flash sale: "Kafka buffers orders behind the reservation step; it is not the inventory lock — that's a conditional write in the DB."

What Interviewers Probe#

After You Say...They Will Ask...What They're Evaluating
"Kafka for decoupling""What's the key? How many partitions?"Ordering scope and sizing
"Exactly-once""What about the email you send?"Knowing the EOS boundary
"Consumers retry on failure""What happens to messages behind it?"Head-of-line blocking awareness
"RF=3""What's min ISR? What happens when two brokers die?"Durability vs availability decision
"We'll add partitions later""What happens to ordering?"One-way door recognition
"Write to DB then publish""What if you crash in between?"Dual-write / outbox

Common Interview Mistakes#

What Candidates SayWhat Interviewers HearWhat Staff Engineers Say
"Kafka guarantees ordering"Missed the per-partition scope"Ordered per partition; I key by order_id so each order is ordered."
"Kafka gives exactly-once"Will double-charge customers"EOS inside Kafka; idempotent sinks everywhere else."
"Just add more consumers"Doesn't know the partition cap"Consumers ≤ partitions; I sized 96 for 3× peak."
"Kafka as the job queue"Head-of-line blocking incoming"Kafka for events; SQS behind it for per-task retries."
"Publish after the DB commit"Dual-write inconsistency"Transactional outbox + CDC."
"Use acks=1 for speed"Silent data loss accepted"acks=all costs ~2–5ms; I'd pay it for orders."

L5 vs L6 vs L7 Responses#

ScenarioL5 AnswerL6 / Staff AnswerL7 / Principal Answer
"Design order event pipeline"Kafka topic, services consumeKey=order_id, 96 partitions, RF=3/minISR=2/acks=all, outbox+Debezium, idempotent consumers, per-group lag SLOsPlus: orders.v1 is a published contract with schema registry, BACKWARD compatibility, a deprecation policy, and a named owning team
"Consumer is 2 hours behind"Add consumersCheck partition cap, sink throughput, per-partition skew; batch writes; shed analytics groupAsk why a lag SLO breach wasn't paged to the consuming team in 5 minutes — fix the ownership routing org-wide
"Kafka bill doubled"Fewer partitionsCompression, fetch-from-follower, tiered storage, retention auditCross-AZ is the bill: evaluate diskless/S3-native tiers for logs, keep low-latency cluster for money, chargeback by produced bytes
"Multi-region"MirrorMakerActive/passive with MM2, offset translation, RPO = lag; region-prefixed topics for active/activeDecide which domains need regional independence at all; define the RPO per domain with business owners

The Staff Kafka Checklist#

  1. Ordering scope → key: "Events for one order must be ordered; key = order_id."
  2. Partition count with math: "50 MB/s, consumers at 1.5K rec/s → 34 minimum, 96 with headroom — and I'll never add partitions to this topic."
  3. Durability config: "RF=3, min ISR 2, acks=all, idempotent producer, unclean election off."
  4. Delivery semantics: "At-least-once; consumers dedupe on event_id; EOS only for Kafka-to-Kafka."
  5. Failure handling: "Retry topics for transient errors, DLQ with a named owner, lag-in-seconds alerts per group."
  6. Contract: "Schema registry, BACKWARD compatible, producing team owns the schema."

🎯 Staff Insight: Don't use Kafka as an RPC bus, a per-task queue with priorities, or a database without compaction and backup. The strongest Kafka signal is naming what must not be ordered together — that's what unlocks parallelism.

Evaluation Rubric#

DimensionSenior (L5)Staff (L6)Principal (L7)
Data modelTopics and consumersKey, partition math, event schema, outboxEvent contract governance across teams
Semantics"Exactly-once"At-least-once + idempotency; EOS boundaryOrg rule: every consumer idempotent; audited
DurabilityRF=3acks/minISR/unclean election, and what fails whenDurability tiers priced per domain
Operations"Monitor lag"Time-lag SLOs per group, DLQ ownership, quotasPlatform SLOs, chargeback, change control for topic configs
Evolution"Add partitions"One-way doors identified up front3-year cluster topology and migration plan

Strong hire signals

SignalWhat It Sounds Like
Key-first design"What must be ordered? That's the key."
Draws the EOS boundary"Transactions stop at the Kafka edge."
Names the one-way door"Partition count on keyed topics is permanent."
Routes ownership"Lag pages the consuming team; schema breaks page the producer."
Knows when not to"This is a task queue; SQS fits better."

Lean no-hire signals

SignalWhy It Misses the Bar
Kafka in every boxTool-first thinking, no cost of operation
No key discussionOrdering and hot-spot risk unaddressed
Believes EOS covers side effectsWill ship duplicate charges
No DLQ ownerSilent data loss

Common false positives

  • Deep ISR/controller internals ≠ event design. Knowing the fetch protocol doesn't tell you the key.
  • "We ran 500 brokers" ≠ judgment — ask what they'd do with 5K msg/s (answer: probably not Kafka).
  • Kafka Streams fluency ≠ semantics — can they explain where exactly-once ends?

The Principal Lens#

Why L7 Sees This Problem Differently#

At Staff level Kafka is a pipe you configure correctly. At Principal level it is the org's integration layer — the place where team boundaries become data contracts. Once 40 services consume orders.v1, its schema is harder to change than any REST API, because consumers are invisible to the producer and replay old data forever. The L7 question is not "how many partitions?" but "Which events are published contracts, who owns them, and how do we evolve them without a 40-team coordination tax?"

The Org-Level Fault Line#

One shared multi-tenant Kafka platform vs. per-domain clusters.

OptionWhat WorksWhat BreaksWho Pays
One shared clusterCross-domain events are trivial; one expert team; best utilizationBlast radius = the company; noisy neighbors; upgrades need everyone's blessingEvery team during the one bad afternoon
Per-domain clusters (payments, logistics, analytics)Isolation; per-domain durability/cost tiers; independent upgradesCross-domain events need mirroring; more clusters to runPlatform headcount
Per-team clustersMaximum autonomy30 half-run clusters, inconsistent configs, no one expertIncident responders, security
Fully managed (MSK/Confluent)No broker opsCost at scale; vendor limits; cross-AZ still billedFinance

The Principal default: a small number of tiered clusters — a low-latency "critical" cluster for money and orders (acks=all, minISR 2, strict change control) and a high-throughput "firehose" cluster for logs and analytics (cheaper storage, possibly S3-native) — run by one platform team, with topic-level ownership and quotas.

Cost Model#

Assumptions: 3 AZs, RF=3, 7-day retention, zstd ~4× on JSON, compute ~$0.04/vCPU-hr, EBS/SSD ~$0.08/GB-month, cross-AZ $0.02/GB round trip, loaded engineer ~$250K/yr. Directional only.

ScaleIngressSelf-host infra/monthCross-AZ/monthHeadcountManaged estimate/month
Startup5 MB/s~$2K (3 brokers)~$0.5–1K0.25 FTE~$2–5K
Growth100 MB/s~$20–35K (12–20 brokers)~$15–30K2 FTE (~$40K)~$60–120K
Enterprise2 GB/s~$300–500K (multiple clusters)~$250K+ unless fetch-from-follower + compression6–10 FTENegotiated; often $1M+

At growth scale, cross-AZ transfer rivals compute — a fact most designs never mention. Fetch-from-follower, compression, and S3-native tiers for firehose data are the three biggest levers.

The 3-Year Evolution Path#

Diagram: The 3-Year Evolution Path

One-Way Doors vs Two-Way Doors#

DecisionReversibilityCost to Reverse
Partition key of a published topicOne-wayNew topic, dual-produce, migrate every consumer
Partition count on keyed topicOne-waySame as above
Event schema field semanticsOne-way once consumers depend on itVersioned topic (v2) and a deprecation quarter
Serialization format (JSON → Avro/Protobuf)One-way-ishEvery producer and consumer changes
Managed vs self-hostedTwo-wayMirroring cutover over weeks; Kafka protocol is portable
Retention, quotas, broker countTwo-wayConfig
Kafka vs Redpanda (same API)Two-wayMirror and cut over

🧭 Principal Move: "I'll let teams pick their broker vendor and cluster size — those are two-way doors. I will not let them publish an unkeyed, unregistered event that other teams consume, because that's the door we can't walk back through."

The Standard I'd Write#

RFC-EVT-002: Published Event Contracts on Kafka

Scope: Any topic consumed by a team other than its producer.

MUST
  1. Register the schema (Avro or Protobuf) with BACKWARD compatibility enforced
     in CI before deploy.
  2. Include event_id (ULID/UUIDv7), occurred_at, producer, schema_version.
  3. Declare the partition key and ordering guarantee in the topic's catalog entry.
  4. Produce with acks=all and idempotence on "critical" tier clusters.
  5. Every consumer MUST be idempotent on event_id; offsets committed after effects.
  6. Every DLQ has an owning team, a depth alert, and a replay runbook.
SHOULD
  7. Use the transactional outbox or CDC for events derived from DB state.
  8. Carry enough state that the top three consumers never call back.

Exceptions: Filed with the event platform team; approved by platform lead +
consuming teams' leads; time-boxed to one quarter.

Deprecation: Breaking changes ship as a new versioned topic; old topic kept
>= 90 days with consumer-migration tracking.

Success metrics: zero schema-break incidents per quarter; 100% of cross-team
topics in the catalog; median consumer onboarding < 1 day.

What I'd Tell the VP#

"Kafka is now how our teams talk to each other, which makes event formats as important as our public API — and today nobody owns them. Last quarter a field rename broke six teams for four hours. I'm proposing a schema registry with automatic compatibility checks, a catalog of who owns each event, and splitting our one cluster into a critical tier for payments and a cheaper tier for logs. That cuts incident risk on the revenue path and trims an estimated 30–40% of the Kafka bill, mostly network transfer. The cost is one quarter of platform work and a policy engineering directors agree to."

Principal Interview Signals#

SignalWhat It Sounds Like
Treats events as contracts"A published topic is an API with invisible consumers."
Prices the network"Cross-AZ replication is our biggest Kafka line, not brokers."
Tiers the platform"Critical and firehose have different durability, latency, and cost."
Guards one-way doors"Key and schema are governed; broker vendor is not."
Knows when not to centralize"Teams with one producer and one consumer don't need the shared cluster — a queue is fine."

Staff answers that L7 interviewers find insufficient:

  • "We'll add a schema registry" — without a compatibility mode, CI enforcement, or a deprecation process.
  • "Separate clusters for isolation" — without the cost or the cross-cluster event story.
  • "Make consumers idempotent" — as advice, rather than a standard with an audit.

In the Wild#

LinkedIn — Where Kafka Came From#

LinkedIn built Kafka around 2010 to unify activity tracking and operational data pipelines that had been point-to-point integrations, then open-sourced it via Apache. By 2019 LinkedIn publicly described processing over 7 trillion messages per day across many clusters. The design bet — a replicated, sequential-disk log with consumer-managed offsets — is what made that throughput possible on commodity hardware.

Staff insight: Kafka was invented to replace N×M point-to-point pipes with one log. When you propose it, say what N×M integrations it replaces — that's the justification, not "decoupling."

Uber — Kafka as the Nervous System#

Uber runs one of the largest Kafka deployments publicly documented, carrying trillions of messages per day for trip events, logs, and data-lake ingestion, and has written about building consumer proxies and cross-region replication tooling (uReplicator) around it. Their consumer proxy pattern exists precisely because many teams wanted queue semantics — per-message retries and DLQs — on top of a log.

Staff insight: When a large org builds a queue-semantics layer over Kafka, it's evidence that log and queue are different shapes. Name the shape your problem needs.

Netflix — Keystone#

Netflix's Keystone pipeline routes event data from producers through Kafka into stream-processing jobs and sinks (S3, Elasticsearch, other Kafka topics), with Flink handling routing and processing. They've written about running Kafka with a tiered design: front-line clusters for producers and consumer-facing clusters downstream, isolating producer availability from consumer behavior.

Staff insight: Split clusters by who can hurt whom. Producers of critical events should never be affected by a consumer backfilling a month of data.


Practice Drill#

Prompt: "Design the event pipeline for a food-delivery platform: order lifecycle events feed restaurant tablets, courier dispatch, customer notifications, and analytics. 20K orders/minute at peak, ~15 events per order."

Staff Answer

Peak is 20K × 15 / 60 = 5K events/s, ~1KB each = 5 MB/s — modest for Kafka, so the design is about semantics, not throughput. Ordering scope is per order (created → accepted → picked up → delivered must not invert), so topic orders.order.v1 keyed by order_id, 48 partitions — not for bytes but for consumer parallelism: dispatch does a DB write + geo query per event at ~300/s per consumer, so ~17 consumers at peak, 48 gives 3× headroom and I will never repartition it. RF=3 across 3 AZs, min.insync.replicas=2, acks=all, idempotent producers, events emitted via transactional outbox from the order DB with Debezium. Consumer groups: dispatch (lag SLO 2s, pages dispatch team), restaurant-push (3s), notifications (10s, retry topics 30s/5m then DLQ owned by notifications team — out-of-order retries acceptable), analytics (15m, reads via fetch-from-follower). Every consumer dedupes on event_id. Schema registry with BACKWARD compatibility; order service owns the schema. Retention 7 days; analytics lands raw events in the lake for anything older.

Why this is L6:

  • Derives partitions from consumer cost, not broker throughput, and names the one-way door.
  • Separates ordering-critical consumers (dispatch: block and page) from retry-tolerant ones (notifications: retry topics).
  • Per-group lag SLOs routed to the owning team.
  • Outbox eliminates the DB/Kafka dual write.

What L7 adds:

  • Declares orders.order.v1 a published contract with a deprecation policy, because courier and restaurant teams will build on it for years.
  • Places analytics on a cheaper firehose tier mirrored from the critical cluster, so a backfill can't hurt dispatch.
  • Prices cross-AZ traffic and turns on fetch-from-follower and zstd from day one.

Quick Reference Card#

Ordering:          per partition only; key -> murmur2(key) mod N
Parallelism:       consumers per group <= partitions
Durability:        RF=3, min.insync.replicas=2, acks=all, unclean election OFF
Producer:          enable.idempotence=true, linger.ms 5-20, zstd/lz4, 128KB batches
Retention:         default 7d; = your replay budget; tiered storage for months
Compaction:        latest value per key; null = tombstone
Per partition:     ~5-10 MB/s planning number
Max msg size:      ~1MB default; big payloads -> S3 + pointer
Rebalance knobs:   max.poll.interval.ms 300s, session.timeout.ms 45s,
                   cooperative-sticky, static membership, KIP-848 in 4.x
Metadata:          KRaft (3.3+ prod ready, ZooKeeper removed in 4.0)
Semantics:         at-least-once + idempotent consumers; EOS only Kafka->Kafka
Dual write fix:    transactional outbox + CDC (Debezium)

RED FLAGS
  - Adding partitions to a keyed topic
  - "Exactly-once" for side effects outside Kafka
  - Kafka as per-task queue with priorities/delays
  - DLQ with no owner
  - acks=1 on money topics
  - Topic per customer
  1. Loading the index…