DEV Community

tejaswipandava
tejaswipandava

Posted on

Kafka Interview Prep (Part 2) — When Things Break in Production

This is Part 2 of a two-part Kafka prep sheet. Part 1 covered the fundamentals — partitions,
brokers, delivery semantics, ISR, KRaft. This part assumes you're comfortable with that and
shifts into the kind of questions that actually separate mid-level from senior/lead candidates:
real-world design reasoning and what happens when production breaks.

📎 Haven't read Part 1 yet? Start there first → Kafka Interview Prep Part 1 — Know Your
Fundamentals Cold
. Everything below assumes you already know partitions, keys, ISR, and the
at-least-once/at-most-once/exactly-once semantics.


Real-World Example: How Zomato Shows a Delivery Partner's Live Location

A great illustration of why Kafka fits a problem, not just that it's used. The question
sounds like simple location tracking, but at scale — millions of delivery partners sending GPS
updates every few seconds — it becomes a genuine streaming/architecture problem.

The Pipeline

📍 Delivery Partner App
   (sends lat, long, timestamp, order_id every few seconds)
        │
        ▼
⚡ Kafka topic: "location-updates"
   (high-throughput event streaming, millions of events reliably)
        │
        ▼
🔄 Stream Processing (consumer services)
   - validate incoming locations
   - filter invalid/duplicate updates
   - enrich with additional data (e.g., ETA, route info)
   - execute business logic
        │
        ▼
⚡ Fast Storage: Redis
   (only the latest location per order/rider — extremely fast reads)
        │
        ▼
📱 Customer App
   (polls or subscribes to the latest location, updates the map in real time)
Enter fullscreen mode Exit fullscreen mode

Why Kafka — Answering the "Why," Not Just the "What"

Interviewers rarely want to hear "Kafka is used for live tracking." They want to hear why it
fits this specific problem
. The strong answers to the questions they're actually probing for:

Q: Why stream events instead of making synchronous API calls?
A: With millions of riders each sending a GPS ping every few seconds, a synchronous
"app → API → DB" call per update would mean every downstream service (map service, ETA service,
analytics, fraud/anomaly detection) has to be called directly and be available at that instant,
for every single ping. That couples the rider app's write path to the uptime and latency of every
consumer, and a single slow consumer would back-pressure the whole ingestion path. Streaming
decouples "accept the update" from "do something with it."

Q: Why introduce a message broker at all?
A: Because there isn't just one consumer of a location update — potentially several: the
"update customer map" service, an ETA-recalculation service, a fraud/anomaly-detection service, an
analytics pipeline. A broker lets all of them consume the same event stream independently,
each at their own pace, without the producer (the rider's app) needing to know or care who's
listening. Add a new consumer later (e.g., a heatmap service) with zero changes to the producer.

Q: How are events processed after Kafka?
A: Consumer services read from the topic and run a validation → dedup → enrichment → business
logic
pipeline — e.g., reject GPS coordinates that are physically impossible (rider "teleporting"
across the city in 2 seconds), drop duplicate pings, enrich with computed ETA, and only then write
the result forward.

Q: Why use Redis for live location storage, not the primary DB?
A: Because the access pattern is "give me the single latest location for this order," read
extremely frequently, overwritten extremely frequently, with no need for history.
A relational DB
optimized for durable, queryable, historical records is the wrong tool for a value that's
overwritten every few seconds and read on every map refresh — Redis's in-memory key-value model
(order_id → {lat, long, timestamp}) gives sub-millisecond reads/writes with none of the overhead
of transactional durability that this specific access pattern doesn't need. (Order history/audit,
if needed, would still be persisted separately, asynchronously, from the same stream.)

Q: How does the customer finally see the moving rider?
A: The customer app either polls the latest-location endpoint (backed by Redis) on an interval,
or — for a more real-time feel — subscribes via WebSockets/long-polling to location updates for
their specific order, with the backend pushing new positions as they land in Redis. Either way, the
customer-facing read path never touches Kafka directly or the heavy stream-processing pipeline — it
only ever reads the cheap, fast, latest-value store.

Why Kafka Specifically (Summary)

  • ✅ Handles millions of location updates per second.
  • ✅ Preserves event ordering within partitions (e.g., partition by order_id so a given rider's location updates are processed in the order they were sent).
  • ✅ Decouples producers from consumers — the rider app doesn't know or care who's downstream.
  • ✅ Multiple independent downstream consumers off the same stream (map updates, ETA, fraud detection, analytics) without touching the producer.
  • ✅ Fault tolerance and horizontal scalability — a consumer going down doesn't lose events; they sit durably in the topic until it recovers or a replacement consumer picks up the partition.

The interview lesson: System Design interviews aren't about memorizing tools — they're about
understanding why each component exists and how they fit together into a coherent, scalable
system. Knowing "Kafka handles high throughput" is knowing a fact; being able to explain why
synchronous calls break down at this scale, why a broker (not just any queue) is the right shape,
and why Redis (not the primary DB) is the right store for the "latest value" access pattern is
what demonstrates actual system design understanding.


Senior/Lead Interview: Production Failure Scenarios

Kafka interviews at the senior/lead level aren't really about Kafka — they're about what happens
when things break
. Knowing "what is a topic / partition / broker" is table stakes. The real
signal comes from how you reason about production incidents:

❌ "What is Kafka?" / "What is a topic?" / "What is a partition?"
✅ "Something broke in production — what happens, and how do you fix it?"

Scenario 1: Consumer crashes after processing but before committing the offset

Setup: The consumer successfully processes a message (e.g., writes to a DB, calls a downstream
API), but crashes before the offset commit is acknowledged back to Kafka.

What happens on restart? Since the offset was never committed, Kafka still considers that
message "unread" for this consumer group. On restart (or after a rebalance hands the partition to
another consumer), the same message is delivered again — the side effect (DB write, API call)
runs a second time. This is duplicate processing, and it's not a bug in Kafka — it's the
expected behavior of at-least-once delivery, which is the default and most common consumer
configuration (commit-after-process).

How do you design for safe processing?

  • Make the processing itself idempotent. This is the real fix, not a workaround — if reprocessing the same message twice produces the same end state (e.g., UPSERT instead of INSERT, or a dedup check keyed on a message ID/business key before applying the side effect), duplicate delivery becomes harmless.
  • Track a processed-message ID (idempotency key) in the same transaction as the side effect — e.g., write "processed message X" and the actual business change atomically (same DB transaction), so a retry can check "have I already handled X?" before reapplying it.
  • Use Kafka transactions / exactly-once semantics (EOS) for read-process-write pipelines (consume from one topic, produce to another) — this makes the consume-offset-commit and the produce atomic as a unit, but note this only covers the Kafka-internal hop; it does not make an external side effect (a REST call, an external DB write) automatically exactly-once — you still need idempotency at that boundary.
  • Never rely on "commit-then-process" (at-most-once) as the fix — that flips the risk from duplicates to silent message loss if the crash happens after commit but before processing completes, which is usually worse for anything business-critical.

The senior-level answer, stated plainly: "Kafka's at-least-once guarantee means duplicates are
an expected, designed-for outcome, not a bug to eliminate at the Kafka layer — the fix belongs in
the consumer's processing logic being idempotent, not in trying to make delivery magically
exactly-once end-to-end."

Scenario 2: Consumer lag keeps increasing

Setup: Producer publishes at 10K msgs/sec, consumer processes at 7K msgs/sec. The gap
compounds over time.

What is consumer lag? The difference between the latest offset produced to a partition and
the consumer's committed offset — i.e., how many messages are sitting unprocessed, waiting for
the consumer to catch up. Rising lag means the consumer is falling behind the producer's rate.

How do you find the bottleneck? Don't jump straight to "add more consumers" — diagnose first:

  • Is it the consumer's processing logic that's slow (e.g., a slow downstream DB write, a slow external API call per message, unnecessary serialization overhead)? Profile per-message processing time.
  • Is it partition count limiting parallelism? Max active consumers in a group = number of partitions. If you only have, say, 4 partitions, adding a 5th consumer does nothing — you're capped regardless of how many consumer instances you run.
  • Is it a downstream dependency that's the actual bottleneck (e.g., the DB the consumer writes to is saturated)? Scaling consumers in that case just shifts load onto an already-struggling downstream system and doesn't fix the lag — it might even make it worse (more concurrent connections hammering a slow DB).
  • Is it skewed partition load (a hot key sending a disproportionate share of traffic to one partition)? Average consumer throughput can look fine while one partition's consumer is drowning — check per-partition lag, not just the aggregate.

Do you just scale consumers? Only if the diagnosis actually points there:

  • If processing is CPU/IO-bound per message and partitions allow it → yes, add consumers (up to the partition count) to parallelize.
  • If partitions are the ceiling → you have to increase partition count first (with the operational caveat that this reshuffles key-to-partition mapping, affecting ordering guarantees for existing keys — plan this deliberately, not reactively mid-incident).
  • If the downstream dependency is the real bottleneck → scaling consumers doesn't help; you need to fix/scale the downstream system, or decouple further (e.g., buffer into an intermediate queue, batch writes, or scale the DB).
  • If it's a transient spike (e.g., a one-off burst) → sometimes the correct answer is "let it drain" rather than permanently over-provisioning consumer capacity for a rare spike.

The senior-level answer: "Consumer lag is a symptom, not a diagnosis. I'd check per-partition
lag, processing time per message, partition count versus consumer count, and downstream dependency
health before deciding whether the fix is more consumers, more partitions, faster processing, or a
downstream fix — 'just add consumers' without that diagnosis can make things worse."

Scenario 3: Multiple event types in one topic

Setup: A single topic (e.g., order-events) now carries OrderCreated, OrderCancelled, and
OrderUpdated events, mixed together.

How do you design the consumer to handle this safely?

  • Include an explicit event-type field in the message (e.g., a header, or a type field in the payload/envelope) — never infer the type from payload shape alone; that's fragile as schemas evolve.
  • Use a schema per event type, all registered under a shared/compatible schema strategy (e.g., a Schema Registry with a union type, or a consistent envelope: {event_type, version, payload}) so the consumer can deserialize correctly before dispatching.
  • Dispatch to type-specific handlers — a single consumer loop reads the envelope, checks event_type, and routes to the appropriate handler (handleOrderCreated, handleOrderCancelled, etc.) rather than one handler trying to branch on ad-hoc payload inspection.
  • Preserve ordering guarantees where they matter — if OrderCreated must be processed before OrderCancelled for the same order, make sure the partition key is still the order_id (not, say, event type) so all events for that order land in the same partition and are processed in send order, regardless of type.
  • Consider whether one topic is even the right call — mixing event types in one topic is fine when consumers need a unified, ordered stream per entity (e.g., a full order lifecycle in order); it's the wrong call if different event types have very different consumers, retention needs, or throughput profiles — in that case, separate topics per event type (with the order_id still as the partition key in each) avoids one high-volume event type's consumers being forced to also filter through low-volume noise, and lets each event type scale/retain independently.

The senior-level answer: "I wouldn't just deserialize-and-branch blindly — I'd standardize on a
typed envelope, dispatch to type-specific handlers, and specifically ask why these event types
share a topic: is it because ordering across types for the same entity matters (keep them
together), or is it just historical accident (in which case, splitting into separate topics per
event type is usually cleaner)."


📎 Haven't covered the basics yet? Go back to Kafka Interview Prep Part 1 — Know Your
Fundamentals Cold
for architecture, partitioning, brokers, and delivery semantics.

Top comments (0)