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)
ConceptMeaning
TopicNamed stream of events
PartitionOrdered, append-only segment of a topic; unit of parallelism
OffsetPosition in a partition
Consumer groupCompeting consumers; each partition assigned to one member
Broker / ISRReplication; durability depends on acks + min ISR
RetentionTime/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

ClaimReality
At-most-onceCommit before processing → possible loss
At-least-onceProcess then commit → duplicates on failure
Exactly-onceEnd-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 / ideaWhy it matters
acks=allWait for ISR — durability
Idempotent producerAvoid dupes on producer retry (enable.idempotence)
TransactionsAtomic writes across partitions/topics (with limits)
Batching / linger / compressionThroughput vs latency
Schema (Avro/Protobuf + Registry)Compatibility evolution

Consumer essentials

TopicNotes
Commit strategyAuto-commit risks; manual commit after successful side effects
RebalanceStop-the-world or cooperative; causes pause / duplicate window
LagBusiness KPI; alert on consumer lag
Poison messagesDLQ / quarantine; don’t block partition forever
Parallelism ceilingMax 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

PatternPurpose
Domain events“OrderPlaced” — facts that happened
Event notificationThin event; consumers fetch details via API
Event-carried state transferFat event; reduces read-back coupling
CQRS + eventsUpdate read models asynchronously
Saga choreographyEach service reacts to events; compensations
OutboxWrite DB + outbox row atomically; publisher relays to Kafka
CDCDB 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

  1. Hot partition — skewed key (e.g. null or popular user)
  2. Rebalance storm — session timeout / slow processing
  3. Dual writes (DB then Kafka without outbox) → lost events
  4. Non-idempotent consumer + at-least-once → double charge
  5. Schema break → poison pill for all consumers
  6. Lag + retention → data expired before consume
  7. Sync RPC inside consumer → couples throughput to dependency latency

Java under the hood (kafka-clients)

ComponentInternals seniors cite
KafkaProducerRecordAccumulator batches by partition; background Sender thread; linger.ms / batch size / compression; acks, idempotent producer (enable.idempotence), transactions
KafkaConsumerUser poll loop + background heartbeat; partition assignment; commitSync / commitAsync after side effects
RebalanceCooperative / 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
SerializationSerializer/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

SyncEvent-driven
CouplingTemporal couplingTemporal decoupling
ConsistencyEasier immediateEventual
DebuggingStack traces / request logsTrace across async hops
EvolutionAPI versioningEvent compatibility
Fan-outExtra callsNatural multi-subscribe

What interviewers probe

  1. How does a consumer group work? What happens when you scale consumers past partitions?
  2. Ordering — how to preserve per-user order.
  3. Exactly-once — push until they admit caveats / idempotency.
  4. Outbox vs dual write
  5. Design notification system / order pipeline with Kafka.
  6. Lag diagnosis approach.
  7. Compacted topics — changelog / latest state per key (e.g. KTables intuition).
  8. 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

Interview reference — explanation quality and judgment, not syntax memorization.