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#
| Behavior | Senior (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 die | Prices 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 storms | Designs 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 DLQ | Writes 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#
| Intent | Constraint | Strategy | Failure Mode | Correctness Bar |
|---|---|---|---|---|
| Event backbone / pub-sub | Many independent consumers, replay, schema evolution | Topic per domain event, keyed, RF=3, 7d+ retention, schema registry | Schema break takes down 12 downstream teams | At-least-once, ordered per key |
| Stream processing input (CDC, analytics, aggregation) | Throughput, ordering per entity, exactly-once aggregates | Debezium CDC or producer events → Flink/Kafka Streams with EOS | Consumer lag grows unbounded; reprocessing double-counts | Exactly-once within the pipeline |
| Work queue (jobs, notifications, tasks) | Per-message retry, delays, priority, independent failure | Kafka can do it with retry topics + DLQ, but SQS/RabbitMQ/share groups fit better | One 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#
| Position | Rationale |
|---|---|
| The partition key is the design | It fixes ordering scope, parallelism, and hot-spot risk. Decide it first, out loud. |
acks=all + min.insync.replicas=2 + RF=3 by default | The only config where an acknowledged write survives a single broker loss without data loss. |
| At-least-once + idempotent consumers | EOS doesn't extend to external side effects; idempotency keys do. |
| Over-partition keyed topics at creation | Adding partitions remaps keys and breaks ordering; size for 2–3× projected peak. |
| Lag is the SLO, not broker CPU | Consumer lag in seconds is what users feel. |
| Every DLQ has a named owner and a replay tool | A DLQ nobody reads is a silent data-loss mechanism. |
| Schemas are contracts with compatibility checks | Schema 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.
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 acks | Acked when | Survives leader crash? | Latency |
|---|---|---|---|
0 | Sent to socket | No — may lose anything | Lowest |
1 | Leader wrote to its log | No — lost if leader dies before followers fetch | Low |
all (-1) | All current ISR members have it | Yes, 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=1and 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."
Rebalancing reassigns partitions when consumers join, leave, or miss heartbeats. Rebalance-related configs every Staff candidate should know:
| Config | Default | Why it matters |
|---|---|---|
session.timeout.ms | 45s (3.0+) | Consumer considered dead after missing heartbeats this long |
max.poll.interval.ms | 300s | If processing one batch takes longer, the consumer is kicked → rebalance → redelivery → loop |
max.poll.records | 500 | Lower it when per-record work is slow |
partition.assignment.strategy | Cooperative sticky (modern clients) | Incremental rebalances move only the partitions that must move |
group.instance.id | unset | Static 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#
| Mode | Config | Keeps | Use for |
|---|---|---|---|
| Delete (time/size) | cleanup.policy=delete, retention.ms=604800000 (7d) | Every record until the segment ages out | Event streams, logs, clickstream |
| Compact | cleanup.policy=compact | Latest record per key forever; null value = tombstone deletes the key | Changelogs, CDC tables, current-state topics, __consumer_offsets |
| Compact + delete | cleanup.policy=compact,delete | Latest per key, but only within the retention window | Bounded state caches |
| Tiered storage | KIP-405, remote.storage.enable=true | Recent segments on broker disk; older segments in S3/GCS | Months 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?"
| Domain | Must be ordered together | Key | Hot-key risk |
|---|---|---|---|
| E-commerce orders | All events of one order (created → paid → shipped) | order_id | Low — orders are small |
| Bank ledger | All postings to one account | account_id | High — a payroll account posts 100K/day |
| Chat | Messages in one conversation | conversation_id | Medium — large public channels |
| Clickstream | Nothing strictly | null or session_id | None with null key |
| CDC from Postgres | All changes to one row | primary key | Low, unless one row is a counter |
| IoT | Readings of one device | device_id | Low 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#
| Semantics | How | Cost | Where it's correct |
|---|---|---|---|
| At-most-once | Commit offset before processing; acks=0/1 | Loses data on crash | Metrics sampling, best-effort telemetry |
| At-least-once | Process, then commit; acks=all, retries on | Duplicates on crash/rebalance | Default 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 interval | Stream 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 design | Anything 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#
| Choice | What Works | What Breaks | Who Pays |
|---|---|---|---|
acks=1 for throughput | ~30–50% lower p99 produce latency | Acked data lost on leader crash | Downstream teams discovering gaps; finance if it's orders |
min.insync.replicas=1 | Writes continue with 2 brokers down | Single copy of acked data | Whoever explains the loss in the post-mortem |
| EOS everywhere | Simpler reasoning inside Kafka | Throughput overhead; false confidence at external sinks | Consumer teams who still double-charge |
| Idempotent consumers | Correct under any redelivery | Dedupe store to operate; key discipline | Each consuming team (correctly) |
| Long retention (30d+) | Replay and backfill after bugs | 3× disk per GB without tiered storage | Platform 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#
| Dimension | Kafka | Redpanda | Apache Pulsar | AWS Kinesis | SQS / RabbitMQ |
|---|---|---|---|---|---|
| Model | Partitioned log | Partitioned log (Kafka API) | Segmented log over BookKeeper, compute/storage split | Sharded log | Queue (per-message ack) |
| Ordering | Per partition | Per partition | Per partition / key-shared | Per shard | FIFO queues only (SQS FIFO ~300–3,000 msg/s per group w/ batching) |
| Throughput | Hundreds of MB/s–GB/s per cluster | Similar or higher per core (C++, thread-per-core) | High; brokers stateless | 1 MB/s in, 2 MB/s out per shard | SQS standard: effectively unlimited; RabbitMQ ~20–50K msg/s per node |
| Replay | Yes, by offset/time | Yes | Yes | Yes, 24h default up to 365d | No (once acked, gone) |
| Consumer scaling | ≤ partitions per group | ≤ partitions | Shared subscriptions scale past partitions | ≤ shards (or enhanced fan-out) | Unlimited competing consumers |
| Ops burden | Medium–high (self-host); low on MSK/Confluent | Lower — single binary, no JVM | High — Pulsar + BookKeeper + metadata store | Zero | Zero (SQS) / medium (RabbitMQ) |
| Pick when | Default event backbone, stream processing, ecosystem (Connect, Debezium, Flink) | Kafka semantics with lower tail latency and simpler ops | Multi-tenancy, geo-replication, huge topic counts | AWS-native, modest scale, zero ops | Task 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#
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)#
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#
| Resource | Rough limit / sizing | Note |
|---|---|---|
| Per-partition produce throughput | ~5–10 MB/s comfortable | Planning number, not a hard cap |
| Broker throughput | ~100–300 MB/s in on modern instances | Network 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 typical | Dominated by linger + consumer poll loop |
| Message size | Default max ~1MB | Put large payloads in S3; send a pointer |
| Consumer group size | ≤ partition count | Extra 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#
- Tune clients — batching, compression,
linger.ms. Often 2–3× for free. - Scale consumers up to the partition count.
- Add brokers + reassign partitions (Cruise Control automates rebalancing). Reassignment copies data; throttle it (
leader.replication.throttled.rate) or it saturates the network. - Add partitions — only for unkeyed topics, or via topic migration for keyed ones.
- 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 orrecords-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_totalrate;join-rate;commit-ratedropping 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_sizeper broker; disk usage > 75% alert; per-client produce-byte quotas exceeded. - Fix: Temporarily lower
retention.mson the largest topics; add brokers and reassign. - Prevention:
retention.bytescaps, per-client produce quotas, tiered storage, capacity alert at 60%.
Operational Reality Matrix#
| Failure | Detection Signal | Blast Radius | Mitigation | Owner |
|---|---|---|---|---|
| Lag spiral | Time-lag > SLO 5m | One consumer group's downstream | Scale consumers, shed load | Consuming team |
| Rebalance storm | Rebalance rate, commit-rate 0 | One consumer group | Static membership, lower batch | Consuming team |
| ISR < min | UnderMinIsrPartitionCount | Every producer to affected partitions | Restore brokers | Kafka platform on-call |
| Poison message | Per-partition lag skew | One partition of one group | DLQ + skip | Consuming team; producer fixes schema |
| Disk full | Disk > 75% | All partitions led by that broker | Lower retention, reassign | Platform |
| Schema break | Deserialization error spike across groups | Every consumer of the topic | Roll back producer | Producing team |
| Controller quorum loss | No active controller metric | Whole cluster metadata frozen | Restore quorum (3→ tolerate 1) | Platform |
When to Use vs. Alternatives#
| Need | Pick | Why |
|---|---|---|
| Event backbone with many consumers and replay | Kafka (MSK/Confluent if you don't want to run it) | Log semantics, ecosystem, throughput |
| Kafka semantics, fewer ops, lower tail latency | Redpanda | No JVM, no separate controller ensemble |
| Per-message retries, delays, visibility timeouts | SQS / RabbitMQ | Queue semantics; no head-of-line blocking |
| Small scale on AWS, zero ops | Kinesis or SQS + SNS | Managed; per-shard pricing is fine below ~50 shards |
| Thousands of tenants, native geo-replication | Pulsar | Topic-count scale, built-in multi-region |
| Request/response between services | gRPC/HTTP | Kafka is not an RPC bus — no reply path, seconds of worst-case latency |
| Cheap log/analytics firehose, latency tolerant | S3-native Kafka-compatible | No 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 LOCKEDor 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.
Operational Concerns#
What the On-Call Actually Does#
- Watches
UnderMinIsrPartitionCountand offline partitions — the only broker metrics that mean "producers are failing right now." - Watches lag in seconds per consumer group against that group's SLO, and routes the page to the consuming team.
- Throttles reassignments so rebalancing data after adding brokers doesn't saturate the network (
leader/follower.replication.throttled.rate). - Performs rolling restarts one broker at a time, waiting for URP = 0 between each. A 30-broker cluster roll takes hours; that's normal.
- Handles offset resets — the most dangerous routine operation. Resetting a shared group to
earliestreplays days of side effects. Require a second approver. - 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#
| Metric | Healthy | Alert |
|---|---|---|
UnderMinIsrPartitionCount | 0 | > 0 for 1m (page) |
OfflinePartitionsCount | 0 | > 0 (page) |
UnderReplicatedPartitions | 0 | > 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 count | exactly 1 | ≠ 1 |
| DLQ depth per topic | 0 growth | Any 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 Study | How Kafka Is Used | Key Pattern |
|---|---|---|
| Message Queues | The log-vs-queue decision itself | Partition = ordering + parallelism unit |
| Stream Processing | Durable input and output for windowed aggregation | EOS Kafka → Flink → Kafka |
| Notification System | Fan-out from domain events to channel workers | Retry topics + DLQ, per-channel consumer groups |
| Payment Processing | Payment state events, ledger postings | Outbox + key by payment/account, idempotent consumers |
| News Feed | Post events driving fan-out-on-write | Key by author; celebrity handling |
| Search Indexing | CDC feed into the indexer | Debezium, compacted changelog |
| Metrics & Monitoring | Buffer between agents and the TSDB | Replay 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_idgives 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 Say | What Interviewers Hear | What 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#
| Scenario | L5 Answer | L6 / Staff Answer | L7 / Principal Answer |
|---|---|---|---|
| "Design order event pipeline" | Kafka topic, services consume | Key=order_id, 96 partitions, RF=3/minISR=2/acks=all, outbox+Debezium, idempotent consumers, per-group lag SLOs | Plus: 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 consumers | Check partition cap, sink throughput, per-partition skew; batch writes; shed analytics group | Ask 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 partitions | Compression, fetch-from-follower, tiered storage, retention audit | Cross-AZ is the bill: evaluate diskless/S3-native tiers for logs, keep low-latency cluster for money, chargeback by produced bytes |
| "Multi-region" | MirrorMaker | Active/passive with MM2, offset translation, RPO = lag; region-prefixed topics for active/active | Decide which domains need regional independence at all; define the RPO per domain with business owners |
The Staff Kafka Checklist#
- Ordering scope → key: "Events for one order must be ordered; key = order_id."
- 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."
- Durability config: "RF=3, min ISR 2, acks=all, idempotent producer, unclean election off."
- Delivery semantics: "At-least-once; consumers dedupe on event_id; EOS only for Kafka-to-Kafka."
- Failure handling: "Retry topics for transient errors, DLQ with a named owner, lag-in-seconds alerts per group."
- 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#
| Dimension | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| Data model | Topics and consumers | Key, partition math, event schema, outbox | Event contract governance across teams |
| Semantics | "Exactly-once" | At-least-once + idempotency; EOS boundary | Org rule: every consumer idempotent; audited |
| Durability | RF=3 | acks/minISR/unclean election, and what fails when | Durability tiers priced per domain |
| Operations | "Monitor lag" | Time-lag SLOs per group, DLQ ownership, quotas | Platform SLOs, chargeback, change control for topic configs |
| Evolution | "Add partitions" | One-way doors identified up front | 3-year cluster topology and migration plan |
Strong hire signals
| Signal | What 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
| Signal | Why It Misses the Bar |
|---|---|
| Kafka in every box | Tool-first thinking, no cost of operation |
| No key discussion | Ordering and hot-spot risk unaddressed |
| Believes EOS covers side effects | Will ship duplicate charges |
| No DLQ owner | Silent 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.
| Option | What Works | What Breaks | Who Pays |
|---|---|---|---|
| One shared cluster | Cross-domain events are trivial; one expert team; best utilization | Blast radius = the company; noisy neighbors; upgrades need everyone's blessing | Every team during the one bad afternoon |
| Per-domain clusters (payments, logistics, analytics) | Isolation; per-domain durability/cost tiers; independent upgrades | Cross-domain events need mirroring; more clusters to run | Platform headcount |
| Per-team clusters | Maximum autonomy | 30 half-run clusters, inconsistent configs, no one expert | Incident responders, security |
| Fully managed (MSK/Confluent) | No broker ops | Cost at scale; vendor limits; cross-AZ still billed | Finance |
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.
| Scale | Ingress | Self-host infra/month | Cross-AZ/month | Headcount | Managed estimate/month |
|---|---|---|---|---|---|
| Startup | 5 MB/s | ~$2K (3 brokers) | ~$0.5–1K | 0.25 FTE | ~$2–5K |
| Growth | 100 MB/s | ~$20–35K (12–20 brokers) | ~$15–30K | 2 FTE (~$40K) | ~$60–120K |
| Enterprise | 2 GB/s | ~$300–500K (multiple clusters) | ~$250K+ unless fetch-from-follower + compression | 6–10 FTE | Negotiated; 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#
One-Way Doors vs Two-Way Doors#
| Decision | Reversibility | Cost to Reverse |
|---|---|---|
| Partition key of a published topic | One-way | New topic, dual-produce, migrate every consumer |
| Partition count on keyed topic | One-way | Same as above |
| Event schema field semantics | One-way once consumers depend on it | Versioned topic (v2) and a deprecation quarter |
| Serialization format (JSON → Avro/Protobuf) | One-way-ish | Every producer and consumer changes |
| Managed vs self-hosted | Two-way | Mirroring cutover over weeks; Kafka protocol is portable |
| Retention, quotas, broker count | Two-way | Config |
| Kafka vs Redpanda (same API) | Two-way | Mirror 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#
| Signal | What 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.v1a 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