Intermediate to senior

Backend Interview Prep

Fourteen chapters on HTTP and API design, SQL, indexing and transactions, NoSQL, authentication, caching, concurrency, messaging, resilience, deployment and observability, with tested SQL and Python.

Chapter 9 of 14Concurrency and distributed systems · Messaging, Queues and Event-Driven Design

Messaging, Queues and Event-Driven Design

Queues and event streams let services communicate without waiting for each other, smooth out traffic spikes, and recover from failures. They also introduce new problems: duplicates, ordering, poison messages and data that is eventually, not immediately, consistent. Interviewers test whether you understand delivery guarantees, idempotent consumers, the outbox pattern and sagas. This chapter covers each with runnable models.

1. Why asynchronous messaging

BenefitExample
Decouplingthe order service does not need the email service to be up
Load levellinga flash sale queues work instead of crashing the database
Reliabilitya message persists until it is processed; consumers retry
Scalabilityadd consumers to drain a backlog
Fan-outone event, many independent reactions (email, analytics, search index)

Costs: eventual consistency, harder debugging (no single stack trace), duplicates and ordering issues, and operating the broker. Do not make a call asynchronous merely because you can; keep the user's critical path synchronous when they need an immediate answer, and move side effects (emails, thumbnails, notifications, analytics) to the background.

2. Two models

Message queue (work queue)Log / event stream
ExamplesRabbitMQ, Amazon SQS, ActiveMQApache Kafka, Amazon Kinesis, Redpanda, Pulsar
Consumptioneach message goes to one consumer and is removed after acknowledgementmessages stay in an ordered, retained log; each consumer group tracks its own offset
Replayno (once acknowledged, gone)yes: reset the offset and reprocess
Orderingper queue, weakened by multiple consumers and retriesper partition, strong
Typical usebackground jobs, task distributionevent sourcing, data pipelines, many independent subscribers, audit
Scalingadd competing consumersadd partitions and consumers in a group (parallelism ≤ partitions)

Publish-subscribe delivers a copy to every subscriber (topics, fan-out exchanges). A log gives pub/sub with replay: every consumer group sees every message.

Kafka vocabulary

A topic is split into partitions; messages with the same key go to the same partition, so they are ordered relative to each other (per-key ordering). A consumer group shares the partitions among its consumers (each partition is read by one consumer in the group). Consumers commit offsets to record progress. Brokers replicate partitions for durability (leader and in-sync replicas), and acks=all makes a write wait for them.

import hashlib

def partition_for(key, partitions):
    return int(hashlib.md5(key.encode()).hexdigest(), 16) % partitions

events = [("order-1", "created"), ("order-2", "created"), ("order-1", "paid"), ("order-1", "shipped"), ("order-2", "paid")]
parts = {p: [] for p in range(3)}
for key, ev in events:
    parts[partition_for(key, 3)].append((key, ev))

# every event of one order lives in a single partition, in the order it was produced
for key in ("order-1", "order-2"):
    owner = partition_for(key, 3)
    sequence = [ev for k, ev in parts[owner] if k == key]
    assert sequence == [ev for k, ev in events if k == key]
assert [ev for k, ev in parts[partition_for("order-1", 3)] if k == "order-1"] == ["created", "paid", "shipped"]

Global ordering across partitions does not exist; design so only per-key order matters (choose the key as the entity whose events must be ordered, such as the order ID or the account ID).

3. Delivery guarantees

GuaranteeMeaningHow it happens
At most oncea message may be lost, never duplicatedack before processing; no retry
At least oncenever lost, may be duplicatedack after processing; retry on failure
Exactly onceprocessed once, effect oncenot generally achievable end to end across systems; approximated with at-least-once delivery plus idempotent processing (or transactions within one system such as Kafka's transactional producers and consumers)

The practical stance: assume at-least-once, and make consumers idempotent.

class Broker:
    """A tiny at-least-once queue: unacknowledged messages are redelivered."""
    def __init__(self):
        self.queue, self.inflight = [], {}
        self.next_id = 0

    def publish(self, body):
        self.next_id += 1
        self.queue.append((self.next_id, body))

    def receive(self):
        if not self.queue:
            return None
        msg = self.queue.pop(0)
        self.inflight[msg[0]] = msg
        return msg

    def ack(self, msg_id):
        self.inflight.pop(msg_id, None)

    def redeliver_unacked(self):                       # the visibility timeout expired
        self.queue = list(self.inflight.values()) + self.queue
        self.inflight = {}

broker = Broker()
broker.publish({"order": 7, "amount": 500})

balance = {"charged": 0}
def naive_handler(body):
    balance["charged"] += body["amount"]               # not idempotent

msg = broker.receive()
naive_handler(msg[1])                                  # processed...
# ...but the consumer crashes before acknowledging
broker.redeliver_unacked()
msg = broker.receive()
naive_handler(msg[1]); broker.ack(msg[0])
assert balance["charged"] == 1000                      # the customer was charged twice: the cost of at-least-once without idempotency

4. Idempotent consumers

Make processing safe to repeat:

  • Natural idempotency: "set status to shipped" is idempotent; "add 500 to the balance" is not.
  • Deduplication table: record processed message IDs (or a business key) in the same transaction as the effect, with a unique constraint.
  • Conditional updates (WHERE status = 'pending') and version checks.
  • Idempotency keys on outgoing calls to other services.
import sqlite3

db = sqlite3.connect(":memory:")
db.executescript("""
CREATE TABLE processed (message_id TEXT PRIMARY KEY);
CREATE TABLE ledger (order_id INTEGER, amount INTEGER);
""")

def handle(message_id, order_id, amount):
    try:
        with db:                                                   # one transaction: dedupe record and effect commit together
            db.execute("INSERT INTO processed VALUES (?)", (message_id,))   # a duplicate fails the unique constraint
            db.execute("INSERT INTO ledger VALUES (?, ?)", (order_id, amount))
        return "processed"
    except sqlite3.IntegrityError:
        return "duplicate ignored"

assert handle("msg-1", 7, 500) == "processed"
assert handle("msg-1", 7, 500) == "duplicate ignored"              # the redelivered message changes nothing
assert db.execute("SELECT SUM(amount) FROM ledger").fetchone()[0] == 500
assert handle("msg-2", 8, 100) == "processed"

5. Failure handling: retries and dead-letter queues

  • Retry transient failures with exponential backoff and jitter; cap the attempts.
  • Poison messages (always fail, such as malformed data) must not block the queue forever: after attempts, move them to a dead-letter queue (DLQ) with the error context, alert on its depth, and provide a way to inspect and redrive them.
  • Visibility timeout (SQS): a received message becomes invisible for a period; if not deleted in time it reappears. Set it longer than your processing time and extend it for long jobs.
  • Ordering and retries conflict: a retried message can arrive after a later one. If order matters, block the partition on failure (head-of-line blocking) or design handlers to tolerate reordering using version numbers.
def process_with_dlq(messages, handler, max_attempts=3):
    done, dead = [], []
    for m in messages:
        for attempt in range(1, max_attempts + 1):
            try:
                handler(m); done.append(m); break
            except Exception as e:
                if attempt == max_attempts:
                    dead.append((m, str(e)))                      # parked with the reason, for inspection and redrive
    return done, dead

def handler(m):
    if m == "bad":
        raise ValueError("cannot parse")

done, dead = process_with_dlq(["a", "bad", "b"], handler)
assert done == ["a", "b"] and dead == [("bad", "cannot parse")]     # one poison message did not block the others

6. The dual-write problem and the outbox pattern

<!--fig:outbox-->
Serviceone database transaction orders row outbox row Relaypoll or CDC, marks sent BrokerKafka / queue Consumersidempotent The order and its event commit together or not at all. Publishing is at-least-once, so consumers must tolerate duplicates. Figure 1. The transactional outbox removes the dual-write problem.

A service must update its database and publish an event. Doing both separately can fail halfway: the database commits but the publish fails (the event is lost), or the publish succeeds but the transaction rolls back (a phantom event). You cannot make a database and a broker commit atomically without distributed transactions.

Transactional outbox: write the business change and an event row into an outbox table in the same database transaction. A separate relay (a poller, or change-data-capture such as Debezium reading the database log) publishes outbox rows to the broker and marks them sent. Delivery to the broker is at-least-once, so consumers must be idempotent.

db.executescript("""
CREATE TABLE orders (id INTEGER PRIMARY KEY, status TEXT);
CREATE TABLE outbox (id INTEGER PRIMARY KEY AUTOINCREMENT, topic TEXT, payload TEXT, published INTEGER DEFAULT 0);
""")
published = []

def place_order(order_id, fail_after_insert=False):
    try:
        with db:                                                    # a single atomic transaction
            db.execute("INSERT INTO orders VALUES (?, 'placed')", (order_id,))
            db.execute("INSERT INTO outbox(topic, payload) VALUES ('orders.placed', ?)", (str(order_id),))
            if fail_after_insert:
                raise RuntimeError("crash before commit")
    except RuntimeError:
        pass

def relay():
    for row_id, topic, payload in db.execute("SELECT id, topic, payload FROM outbox WHERE published = 0").fetchall():
        published.append((topic, payload))                          # send to the broker (may be retried: at least once)
        with db:
            db.execute("UPDATE outbox SET published = 1 WHERE id = ?", (row_id,))

place_order(1)
place_order(2, fail_after_insert=True)                              # rolled back: neither the order nor the event exists
relay(); relay()                                                    # running the relay twice does not republish
assert published == [("orders.placed", "1")]
assert db.execute("SELECT COUNT(*) FROM orders").fetchone()[0] == 1

7. Sagas: distributed workflows without distributed transactions

An order spans services: reserve stock, charge payment, schedule shipping. No single ACID transaction covers them, so use a saga: a sequence of local transactions, each with a compensating action that undoes it if a later step fails.

  • Orchestration: a coordinator calls each step and triggers compensations on failure. Easier to follow and monitor; the orchestrator is a central component.
  • Choreography: services react to each other's events. Looser coupling; the overall flow is implicit and harder to trace.
def run_saga(steps):
    """steps: list of (name, action, compensation). On failure, undo the completed steps in reverse order."""
    completed, log = [], []
    for name, action, compensation in steps:
        try:
            action(); completed.append((name, compensation)); log.append(f"{name}: done")
        except Exception as e:
            log.append(f"{name}: failed ({e})")
            for done_name, comp in reversed(completed):
                comp(); log.append(f"{done_name}: compensated")
            return False, log
    return True, log

state = {"stock": 5, "charged": 0}
def reserve(): state["stock"] -= 1
def release(): state["stock"] += 1
def charge():  raise RuntimeError("card declined")
def refund():  state["charged"] = 0

ok, log = run_saga([("reserve stock", reserve, release), ("charge card", charge, refund)])
assert ok is False and state["stock"] == 5                         # the reservation was rolled back
assert log == ["reserve stock: done", "charge card: failed (card declined)", "reserve stock: compensated"]

Compensations are business-level undo steps (a refund, not a database rollback), may themselves fail (retry them), and must be idempotent. Between steps, other users can observe intermediate states, which you accept or hide with a "pending" status.

8. Event-driven architecture patterns

  • Event notification: a thin event ("order 7 placed") tells others something happened; they fetch details if needed.
  • Event-carried state transfer: the event contains the data so consumers keep local copies without calling back.
  • Event sourcing: the events are the source of truth; current state is a fold of the event history. Gives a complete audit trail and the ability to rebuild or replay, at the cost of complexity (schema evolution, snapshots, eventual consistency in read models).
  • CQRS: separate the write model from one or more read models optimised for queries, updated by events.
  • Change data capture (CDC): stream database changes into a log for caches, search indexes and analytics.
# event sourcing in miniature: state is derived by replaying events
events = [("deposited", 100), ("withdrew", 30), ("deposited", 50), ("withdrew", 20)]
def balance_after(events, n=None):
    total = 0
    for kind, amount in events[:n]:
        total += amount if kind == "deposited" else -amount
    return total
assert balance_after(events) == 100
assert balance_after(events, 2) == 70                              # time travel: the balance after the second event

Schema evolution

Events live a long time and have many readers. Use a schema registry or versioned schemas (Avro, Protobuf, JSON Schema), only make backward-compatible changes (add optional fields, never repurpose or remove them), and include an event ID, type, version, timestamp and correlation ID in an envelope.

9. Backpressure and consumer lag

When consumers fall behind, the backlog grows. Monitor consumer lag (messages or seconds behind), queue depth and age of the oldest message; scale consumers (up to the partition count in Kafka), batch for throughput, and shed or defer low-priority work. Bound queue sizes where possible, and design producers to slow down or fail fast rather than grow memory.

10. Choosing

NeedChoice
Simple background jobs with retries and delaySQS, RabbitMQ, a Redis-based job queue (Sidekiq, BullMQ, Celery)
High-throughput ordered streams with replay and many consumersKafka, Pulsar, Kinesis
Request-reply with low latencydirect RPC (gRPC, HTTP), not a queue
Scheduled and delayed worka scheduler with a durable store, or delayed queue features
Workflow with human steps and long timersa workflow engine (Temporal, Step Functions)

11. Common mistakes

  • Assuming exactly-once delivery and writing non-idempotent consumers.
  • Dual writes to a database and a broker with no outbox.
  • No dead-letter queue, so one bad message stalls everything or loops forever.
  • Ignoring ordering: processing related events out of order, or expecting global order.
  • Unbounded retries with no backoff, amplifying an outage.
  • A giant, shared event schema with no versioning.
  • Using a queue for request-reply where the caller really needs an immediate answer.
  • Large payloads in messages instead of references to object storage.
  • No monitoring of lag, DLQ depth or message age.

12. Practice questions

  1. Compare a message queue and a log such as Kafka. When would you pick each?
  2. What are at-most-once, at-least-once and exactly-once delivery, and how do you get the effect of exactly-once?
  3. How do you make a payment consumer idempotent?
  4. Explain the dual-write problem and the transactional outbox pattern.
  5. What is a saga? Orchestration versus choreography?
  6. How does Kafka preserve per-key ordering and scale consumers?
  7. What do you do with a message that fails repeatedly?
  8. What is event sourcing, and what are its costs?
Header Logo