Kafka and Event-Driven Architecture
On this page 17
Part IV — System design & architecture · Interview reference
Kafka is a distributed durable log used for high-throughput event streaming. Seniors must explain partitions, consumer groups, delivery semantics, and how event-driven design changes consistency — including the myths around “exactly once.”
Core Kafka mental model
Producer → Topic (partitions: ordered logs) → Consumer group (each partition → ≤1 consumer in group)
↑
retention by time/size (not delete-on-ack like classic queues)
| Concept | Meaning |
|---|---|
| Topic | Named stream of events |
| Partition | Ordered, append-only segment of a topic; unit of parallelism |
| Offset | Position in a partition |
| Consumer group | Competing consumers; each partition assigned to one member |
| Broker / ISR | Replication; durability depends on acks + min ISR |
| Retention | Time/size based; consumers can rewind (unlike many queues) |
Ordering guarantees (state precisely)
- Total order: only within a partition
- Cross-partition: no global order
- Keyed producers: same key → same partition (unless custom partitioner / topic change) → per-key order
Design implication: choose partition key for the ordering you need (e.g. orderId).
Production case study (high volume)
Context: Ride-hailing trip events keyed poorly (cityId) → one partition owns a mega-city; LinkedIn-scale activity stream originally inspired Kafka’s design goals.
Why seniors care: Ordering is per-partition; hot keys destroy consumer parallelism; seniors pick keys for both order and balance.
Failure / symptom: One consumer CPU pegged; lag on partition 7 only; producers see elevated request-latency.
Resolution: Key by tripId/userId with documented order needs; salt extreme hot keys; add partitions with migration plan; monitor per-partition lag.
Seen at / similar to: LinkedIn Kafka origin story; Uber/Lyft event buses; Netflix Keystone; Shopify Kafka at Black Friday.
Delivery semantics
| Claim | Reality |
|---|---|
| At-most-once | Commit before processing → possible loss |
| At-least-once | Process then commit → duplicates on failure |
| Exactly-once | End-to-end EOS is hard; Kafka transactions + idempotent producer cover specific pipelines (consume-process-produce) under constraints — not magic for all side effects (emails, external charge APIs) |
Senior line: Design for at-least-once + idempotent consumers unless your pipeline fits Kafka transactions end-to-end.
Production case study (high volume)
Context: Payments ledger consumer charged cards then committed offsets; crash between charge and commit → double charge on redelivery.
Why seniors care: “Exactly-once” marketing doesn’t cover Stripe side effects; seniors design idempotent consumers first.
Failure / symptom: Duplicate captures for same eventId; finance reconciles mismatches; lag looks healthy.
Resolution: Idempotency store keyed by event/order id; process-then-commit with transactional outbox for downstream; Kafka transactions only for Kafka↔Kafka pipelines.
Seen at / similar to: Stripe + Kafka patterns; banking core processors; Confluent EOS docs caveats.
Producer essentials
| Setting / idea | Why it matters |
|---|---|
acks=all | Wait for ISR — durability |
| Idempotent producer | Avoid dupes on producer retry (enable.idempotence) |
| Transactions | Atomic writes across partitions/topics (with limits) |
| Batching / linger / compression | Throughput vs latency |
| Schema (Avro/Protobuf + Registry) | Compatibility evolution |
Consumer essentials
| Topic | Notes |
|---|---|
| Commit strategy | Auto-commit risks; manual commit after successful side effects |
| Rebalance | Stop-the-world or cooperative; causes pause / duplicate window |
| Lag | Business KPI; alert on consumer lag |
| Poison messages | DLQ / quarantine; don’t block partition forever |
| Parallelism ceiling | Max useful consumers ≈ partition count |
Production case study (high volume)
Context: Ad-click stream (~billions events/day) consumers doing sync HTTP to a feature store inside poll loop; session timeouts → rebalance storms.
Why seniors care: Rebalances pause processing and widen duplicate windows; lag is a product freshness KPI; poison pills stall a partition forever.
Failure / symptom: Frequent rebalances; lag spikes; one bad message blocks partition; ISR shrink under disk pressure.
Resolution: Decouple processing (queue handoff); cooperative rebalance; DLQ; raise timeouts carefully; alert lag vs retention; don’t block heartbeat thread.
Seen at / similar to: LinkedIn/Netflix consumer lag culture; Datadog/Confluent monitoring guides; Uber H3/geo event pipelines.
Event-driven architecture (EDA) patterns
| Pattern | Purpose |
|---|---|
| Domain events | “OrderPlaced” — facts that happened |
| Event notification | Thin event; consumers fetch details via API |
| Event-carried state transfer | Fat event; reduces read-back coupling |
| CQRS + events | Update read models asynchronously |
| Saga choreography | Each service reacts to events; compensations |
| Outbox | Write DB + outbox row atomically; publisher relays to Kafka |
| CDC | DB change stream (Debezium) → Kafka |
Production case study (high volume)
Context: Order service dual-wrote to Postgres then Kafka; under disk pressure Kafka produce failed after DB commit → missing “OrderPlaced” for downstream fulfillment. Why seniors care: Dual-write is a classic SEV; outbox/CDC is the senior fix; blast radius is silent data loss, not loud errors. Failure / symptom: Fulfillment gaps; support “order paid but not shipped”; no producer error on request path. Resolution: Transactional outbox; Debezium CDC; reconcile jobs; metrics on outbox age; contract tests for events. Seen at / similar to: Shopify/Amazon outbox narratives; Debezium at Walmart/Uber-scale adopters; Nubank/banking event cores.
When Kafka / EDA fits
Use when
- Multiple independent consumers need the same facts
- High throughput event ingestion / activity streams
- Decoupling producers from evolving consumers
- Replay / rebuild read models matters
Avoid / reconsider when
- Simple request/response with one consumer and low scale (SQS/queue may suffice)
- Strong sync consistency required for UX without careful design
- Team cannot operate brokers, schemas, lag, rebalances
Failure modes interviewers love
- Hot partition — skewed key (e.g.
nullor popular user) - Rebalance storm — session timeout / slow processing
- Dual writes (DB then Kafka without outbox) → lost events
- Non-idempotent consumer + at-least-once → double charge
- Schema break → poison pill for all consumers
- Lag + retention → data expired before consume
- Sync RPC inside consumer → couples throughput to dependency latency
Java under the hood (kafka-clients)
| Component | Internals seniors cite |
|---|---|
KafkaProducer | RecordAccumulator batches by partition; background Sender thread; linger.ms / batch size / compression; acks, idempotent producer (enable.idempotence), transactions |
KafkaConsumer | User poll loop + background heartbeat; partition assignment; commitSync / commitAsync after side effects |
| Rebalance | Cooperative / eager; pause processing window → duplicate risk with at-least-once |
| Spring | @KafkaListener on a pool of consumer threads — size vs partition count; error handlers / DLQ |
| Serialization | Serializer/Deserializer; Confluent Avro/Protobuf + Schema Registry common |
Thread model: do not block the Sender/heartbeat paths with heavy work; bound consumer concurrency and downstream pools. See thread pools.
EDA tradeoffs vs sync microservices
| Sync | Event-driven | |
|---|---|---|
| Coupling | Temporal coupling | Temporal decoupling |
| Consistency | Easier immediate | Eventual |
| Debugging | Stack traces / request logs | Trace across async hops |
| Evolution | API versioning | Event compatibility |
| Fan-out | Extra calls | Natural multi-subscribe |
What interviewers probe
- How does a consumer group work? What happens when you scale consumers past partitions?
- Ordering — how to preserve per-user order.
- Exactly-once — push until they admit caveats / idempotency.
- Outbox vs dual write
- Design notification system / order pipeline with Kafka.
- Lag diagnosis approach.
- Compacted topics — changelog / latest state per key (e.g. KTables intuition).
- Hot partition + lag diagnosis using per-partition metrics.
Senior-level expectation: Correct semantics; keying/partitioning design; operational metrics (lag, ISR, rebalances); idempotent side effects.
Pitfalls
- Using Kafka as RPC (request-response over topics) without need
- Too few partitions early (hard to raise without thinking through keys)
- Giant payloads in topics (store blob in object storage; event points to it)
- Ignoring schema compatibility
- Treating consumer lag as “infra only” rather than product freshness SLA
- Dual-write without outbox
- Believing EOS covers external payment APIs
Cross-references
- Microservices rules
- Database joins & query optimisations — outbox/CDC
- Thread pools — consumer parallelism
- AWS services — MSK / alternatives (SQS, Kinesis)
- Spec-driven development — event schemas as contracts