Mid to Senior Engineer

System Design Interview Prep

A structured path from the interview framework through core concepts, key technologies and patterns to eighteen full problem breakdowns, each with diagrams and weak, solid and excellent answers to every deep dive.

Chapter 5 of 36Core concepts · Queues, Streams and Asynchronous Processing

Queues, Streams and Asynchronous Processing

Synchronous calls tie a request to every step of the work. Asynchronous processing breaks that tie: the request records what must happen, answers quickly, and workers do the work later. This chapter covers the components that make that possible and the questions that follow.

1. Why go asynchronous

Reasons that justify a queue in a design, each tied to a requirement:

  • Absorb spikes. A burst of writes is buffered and processed at a steady rate, instead of overloading the database.
  • Decouple components. The sender does not need the receiver to be up, fast or even known.
  • Shorten the user-facing path. Send the email, resize the image or update the search index after the response, not during it.
  • Retry failures independently. A failed step is retried without making the user wait or resubmit.
  • Fan out. One event can feed many consumers.

The costs are equally real: the system becomes eventually consistent, ordering and duplicates need thought, debugging follows a message across components, and the queue itself must be run and monitored.

2. Two models: queues and logs

A message queue delivers each message to one consumer among a pool of workers. Once a worker acknowledges it, the message is gone. It models task distribution: "process this job".

A publish-subscribe topic delivers each message to every subscribed consumer group. It models events: "an order was placed", of interest to billing, shipping and analytics.

A distributed log (streaming platform) stores messages in order, durably, for a retention period. Consumers track their own offset, so they can read at their own speed and replay history. The same log serves as a queue (one consumer group) and as pub-sub (many groups).

The log model has become the default for high-volume event pipelines, because replay makes recovery, backfills and new consumers easy.

3. How a partitioned log works

A topic is split into partitions. Within a partition, messages are strictly ordered; across partitions there is no ordering guarantee. A producer chooses the partition by hashing a key, so all messages with the same key land in the same partition and stay in order.

topic "orders", 4 partitions
key = order_id  ->  hash(order_id) mod 4  ->  partition

A consumer group assigns each partition to exactly one consumer in the group. So:

  • parallelism is limited by the number of partitions, because extra consumers beyond that sit idle;
  • per-key ordering is preserved, since one consumer reads a partition in order;
  • adding consumers rebalances partitions across them.

Choose the key to give both ordering where it matters and an even spread. A key with a heavy skew, such as one enormous customer, overloads a single partition.

<!--fig:partitions-->
Topic "orders" - 4 partitions, ordered within each partition P0 0 1 2 3 4 5 6 7 P1 0 1 2 3 4 5 6 7 P2 0 1 2 3 4 5 6 7 P3 0 1 2 3 4 5 6 7 Producerskey = order_id Consumer 1reads P0, P1 Consumer 2reads P2 Consumer 3reads P3 Consumer group "billing" Same key goes to the same partition, so one order's events stay in order. Parallelism is capped by the partition count. Figure 1. A partitioned log: order within a partition, parallelism across partitions, one consumer per partition within a group.

Partitions are replicated across brokers, with one leader per partition and followers that stay in sync. A write is acknowledged after the configured set of in-sync replicas has it, which is the same durability-versus-latency trade-off as database replication.

4. Delivery semantics

GuaranteeMeaningHow it happens
At-most-onceNever duplicated, may be lostAcknowledge before processing
At-least-onceNever lost, may be duplicatedAcknowledge after processing; retry on failure
Exactly-once effectProcessed once in outcomeAt-least-once delivery plus idempotent or transactional processing

If a worker crashes after doing the work but before acknowledging, the message is redelivered and processed again. So consumers must be idempotent: record processed message identifiers, use natural idempotency such as "set" operations, or write the result and the offset in one transaction.

5. Failure handling

  • Retries with backoff. Retry transient errors with exponentially increasing delays and random jitter, so many consumers do not retry in lockstep.
  • Dead-letter queue (DLQ). After a set number of failures, move the message to a separate queue for inspection. Without one, a single bad message (a "poison pill") blocks a partition or loops forever.
  • Visibility timeout. In queue systems, a message being processed is hidden from others for a period. If the worker does not acknowledge in time, it reappears. Set the timeout longer than the work takes.
  • Ordering versus retries. Retrying one failed message in an ordered partition blocks everything behind it. Choose to block, to skip with a DLQ, or to retry from a side queue, and say which and why.
<!--fig:retry-->
success transient error poison message redeliver Main queue Worker idempotent handler Ack: done Retry queue delay 1s, 2s, 4s + jitter Dead-letter queue after N failures Alert + replayafter fix Figure 2. Failure handling: retry with backoff and jitter, then park the message in a dead-letter queue instead of blocking everything behind it.

6. Backpressure

If producers outpace consumers, the queue grows without bound, latency rises and storage fills. Handle it deliberately:

  • Monitor lag: the gap between the newest offset and the consumer's offset.
  • Scale consumers up, within the partition limit.
  • Bound the queue and reject or shed low-priority work when it is full.
  • Slow the producer, for example by returning a retry-after response.
  • Prioritise: separate queues for urgent and bulk work, so a flood of one does not delay the other.

A queue that hides overload does not remove it, and the interviewer wants to hear what you do when lag keeps growing.

7. Patterns you can name

Work queue. Many workers drain one queue. Add workers to increase throughput.

Fan-out. One event goes to many consumers, each with its own queue or consumer group, such as order placed leading to email, inventory and analytics.

Outbox. Write the state change and the outgoing event in one database transaction, and publish from the outbox table afterwards. Avoids losing events between a database commit and a broker send.

Change data capture (CDC). Stream the database's change log to other systems, keeping search indexes, caches and warehouses in sync without dual writes.

Event sourcing. Store the sequence of events as the source of truth and derive current state by replaying them. It gives a perfect audit trail and replay, at the price of complexity in querying and evolving event formats.

CQRS. Separate the write model from one or more read models built from events, so each is optimised for its job. Read models are eventually consistent.

Scheduled and delayed jobs. A delay queue or a scheduler holds work until its time, which is how reminders and retries with delays are done.

8. Batch and stream processing

Batch processing computes over a bounded dataset, such as the day's logs. It favours throughput and is simple to reason about.

Stream processing computes continuously over unbounded events and gives low latency. It brings new questions:

  • Event time versus processing time. Events arrive late and out of order, so windows are usually defined on event time.
  • Windows: tumbling (fixed, non-overlapping), sliding (overlapping) and session (gap-based).
  • Watermarks estimate how far event time has progressed so that a window can close even though late data may still arrive.
  • State such as counters must be stored durably and checkpointed so that processing resumes correctly after a crash.

A lambda architecture runs a batch layer and a speed layer in parallel and merges them. A kappa architecture uses a single streaming path with replay for reprocessing, which is simpler to maintain.

Potential deep dives

Deep dive 1: How do you guarantee ordering?

The challenge. Events about one entity must be processed in order, yet you want parallelism.

Weak: one global queue and one consumer. Ordering is preserved and throughput is capped at one worker.

Solid: partition by key. Route events with the same key to the same partition, and read each partition with one consumer. Order holds per key, and parallelism scales with the number of partitions.

Excellent: also handle retries, skew and rebalances. A failing message blocks its partition, so decide between blocking, skipping to a dead-letter queue, or moving it to a retry queue (and accept that order for that key is relaxed). Choose a key with an even spread to avoid hot partitions, and attach sequence numbers so consumers can detect gaps or duplicates. Plan partition counts for growth, because changing them changes key placement. Handle rebalances so a consumer losing a partition does not process a message twice at the same time.

Deep dive 2: Exactly-once, honestly

The challenge. Duplicates cause double effects, and loss causes missing effects.

Weak: claim exactly-once delivery. It is not achievable in general over an unreliable network, and the claim signals a gap in understanding.

Solid: at-least-once delivery and idempotent consumers. Acknowledge after processing, so a crash causes redelivery, and make the handler safe to repeat with a processed-message identifier or a naturally idempotent operation.

Excellent: be precise about the boundary. Use transactional features where the whole flow stays inside the log (read, process and write back with the offset committed atomically). For effects outside it, write the result and the message identifier in one database transaction, or use the outbox on the producer side. State the guarantee as exactly-once effect, built from at-least-once delivery plus idempotency.

Deep dive 4: What do you do when consumers fall behind?

The challenge. Producers outpace consumers, the backlog grows and latency rises.

Weak: add a bigger queue. It hides the problem until storage fills or the oldest message is hours old.

Solid: scale consumers and watch lag. Add consumers up to the partition count, alert on lag, and optimise the slow handler.

Excellent: shed, prioritise and degrade deliberately. Separate queues by priority so urgent work is not stuck behind bulk work. Bound queue size or age and drop or defer low-value messages. Apply backpressure to producers where possible. Make slow downstream calls asynchronous and batched. Decide in advance what the system does when the backlog is permanent: reject, degrade or add capacity, and communicate it.

Deep dive 5: Queue or log?

The challenge. Both move messages. They suit different problems.

Weak: use whichever is familiar. Misusing a log as a task queue with per-message delays is awkward, and misusing a queue for replayable event history loses the history.

Solid: a queue for tasks, a log for events. A queue distributes tasks among competing workers, with acknowledgement, delay and priority. A log retains ordered events that many consumers read independently and replay.

Excellent: say what each gives and costs. A log gives replay, independent consumer groups and high throughput, and its unit of progress is an offset, so one slow message blocks its partition. A queue gives per-message acknowledgement, delays and dead-lettering, and typically gives up long retention and replay. Many systems use both: a log as the backbone of events and queues for task execution fed from it.

What is expected at each level

Mid-level. You know why a queue helps, the difference between a queue and a topic, and that messages can be delivered more than once.

Senior. You explain partitions and ordering by key, at-least-once delivery with idempotency, dead-letter queues, backpressure and lag, and choose between a log and a queue with reasons.

Staff. You reason about the operational model: capacity planning for retention and replication, schema evolution, replay safely at scale, multi-region replication of streams and cost.

Interview questions and model answers

Q: When would you introduce a queue into a design? When a requirement asks for it: spikes the database cannot take, work that need not finish before the response, or several consumers of the same event. I would state which of those applies, not add it by reflex.

Q: How do you guarantee ordering? Only within a partition, so I use a key such as the order identifier so that all events for one order go to one partition. I do not rely on global ordering, since it prevents parallelism.

Q: A consumer processed the same message twice. Why and what do you do? At-least-once delivery means a crash before acknowledgement causes redelivery. I make the handler idempotent, for example by storing the message identifier with the result in one transaction and skipping known identifiers.

Q: What if one message keeps failing? Retry a few times with exponential backoff and jitter, then send it to a dead-letter queue with its error, alert on the DLQ depth and replay after the fix. Otherwise it blocks the partition.

Q: The queue is growing faster than consumers drain it. What next? Check consumer lag and throughput, scale consumers up to the partition count, increase partitions if needed, optimise the slow handler, and if the load is a sustained excess, apply backpressure or shed low-priority work.

Common mistakes

  • Adding a queue with no reason tied to a requirement.
  • Assuming global ordering.
  • Forgetting idempotency when delivery is at-least-once.
  • No dead-letter strategy, so one bad message stalls everything.
  • Using more consumers than partitions and expecting more throughput.
Header Logo