System designby Learnastra

System-design interview · Extended interviews

Design a distributed message log

By Anup Rai

Design a retained log that preserves acknowledged messages, orders records within each partition and lets independent consumers recover without skipping or repeating business updates.

You will learn to

  • Explain offsets and consumer groups using one retained log.
  • Choose ordering, acknowledgement, and retention guarantees independently.
  • Trace producer retry, consumer replay, and safe external side effects.

Practice in this chapter

8 interview questions with model answers and follow-ups.

Go to interview practice

Useful foundations: Message queues, event logs, delivery guarantees, and backpressure · Replication and durability · Quorums, consensus, leases, and fencing · Databases, data models, and ACID transactions

Workload and timing examples are interview assumptions.

01Problem and scope

A distributed message log retains ordered records so independent consumers can process and replay them at different rates. An offset locates a record within one partition. Each consumer group saves the next offset it needs to read. The broker acknowledging storage, a consumer receiving the message and that consumer completing its database update or external action are separate events. Event M17 for order O51 is an example: a downstream search consumer can resume after downtime because the accepted event remains retained.

On one server, append records to a file and store each consumer’s next offset. A failed consumer resumes from that bookmark. This already explains persistence and replay. Distribution adds partitions, replication, and ownership changes; it does not change why the bookmark exists.

I ask, “Should a completed consumer remove the message, or must several independent services replay it?” We choose a retained partitioned log. Search, fraud analysis, and analytics each own a consumer group and independently advance their bookmarks. A conventional work queue that hides one message while a worker processes it is a different interface and is not silently included.

The ordering requirement is also explicit: events for customer42 should remain ordered within their route, but two unrelated customers do not need one global order. Order O51 is accepted by the shop's database before M17 enters this system. The producer uses an outbox if that business commit and event creation must survive together. Our broker cannot repair an event that the application never durably recorded.

02Functional requirements

  1. Manage topics and partitions. Create topics under an administrative policy; retain independent partitions with compatible routes for each ordering key. Global order across partitions is outside scope.
  2. Append durably. Return a partition and offset only after the broker has stored the record under the promised durability policy. This does not mean the search projection has updated O51.
  3. Read with independent consumer groups. Fetch bounded batches. Within one group, one current owner reads a partition; other groups consume the same retained events independently.
  4. Commit progress. Persist a group's next-offset bookmark. It is not an atomic acknowledgment from every external system that consumed the message.
  5. Replay or reset explicitly. Support seven-day replay and at-least-once delivery. An expired offset returns an explicit error; rebuild from an archive or current-state snapshot when required.
  6. Handle poison events deliberately. Configure bounded retry, stop, or quarantine with an auditable skip for each group.

Group isolation and retention

The producer can stop the search consumer for ten minutes, restart from its saved next offset, and rebuild the missing view. Another group continues unaffected. Deleting a group does not delete topic data, while retention may eventually delete data even if a group has not consumed it. The broker does not pretend an expired interval was processed.

Retained log versus task queue

A task queue often adds per-message visibility leases and acknowledgements so workers can compete for individual jobs. A retained log keeps history and advances group offsets. Specify append durability separately from processing correctness. Provider transaction/idempotency features have defined boundaries, as Kafka's design explains. Kafka design.

03Non-functional requirements

  1. Throughput. Assume 100,000 one-kilobyte messages/s sustained and a twofold short peak; validate storage and network capacity by benchmark.
  2. Latency. Target regional append p95 below 20 ms and healthy fetch visibility within 100 ms of commitment.
  3. Availability and replay. Target 99.95% monthly append availability and seven days of replay. These are exercise assumptions.
  4. Durability. Replicate each partition across three failure domains with durable records on acknowledging replicas. A committed append survives one replica loss.
  5. Partition safety. Without its write quorum, a partition refuses new appends instead of returning an accepted offset that can vanish after failover.
  6. Ordering. Provide committed order within a partition. Parallel workers must preserve the required business order. Move the bookmark only past records whose work has all finished, without skipping a gap.
  7. Bounded resources. Bound producer buffers, broker queues, fetch sizes and disk retention. Apply backpressure or retriable rejection before reserved backlog space runs out.

Guarantee boundaries

The quorum/durable-log protocol is this design's chosen contract, not a claim that every broker configuration has identical acknowledgment semantics. Do not silently delete data inside the promised retention window to keep returning success. Exactly-once effects in an arbitrary external service remain outside the broker's guarantee.

04Capacity estimates

Assume 100K messages/second at one KB, using decimal units.

Quantity Arithmetic Consequence
Ingest 100K/s × 1 KB = 100 MB/s Sequential batched writes
Daily raw log 100 MB/s × 86,400 = 8.64 TB Retention is substantial
Seven days 8.64 × 7 = 60.48 TB Before compression/indexes
Three copies 60.48 × 3 = 181.44 TB Add recovery headroom
Ten-minute stopped consumer 100K/s × 600s = 60M messages Must outpace new traffic to catch up

One hundred partitions would average 1,000 messages/second each. Benchmark throughput rather than claiming that count is enough. More consumers help only when there are partitions for them to own. Events for one heavily used ordering key must still be processed in order.

Three copies of 100 MB/s create 300 MB/s of aggregate replica write payload, and follower replication adds about 200 MB/s of network traffic beyond producer ingress. Every additional full-rate consumer group reads another 100 MB/s before compression and batching. With five groups, delivery can dominate the client-facing network even if storage writes remain sequential.

A consumer recovering 60 million messages while new traffic continues at 100,000/s must exceed that arrival rate. At 150,000/s, net catch-up is 50,000/s and recovery takes 1,200 seconds, or twenty minutes. At exactly 100,000/s it never catches up. Consumer lag should therefore be reported in time and bytes as well as record count.

If a measured partition can safely sustain 5 MB/s of this workload, 100 MB/s suggests twenty partitions before headroom, skew, and consumer parallelism. One hundred may be reasonable, but the benchmark and workload must justify it. A single key producing 20 MB/s cannot be split across those partitions without changing its ordering contract.

05APIs and contracts

Two identifiers answer different questions here: the producer identity and sequence identify an append attempt, while the group and next offset identify how far a particular consumer has finished. An epoch identifies the current partition-leader term; a generation identifies the current consumer-group assignment. They let the broker reject requests made under obsolete ownership.

Operation Example
Append append(orders,key=customer42,producer=P8,seq=44,payload=M17)
Durable result {partition:P2,offset:117,epoch:9}
Fetch fetch(P2,from=117,maxBytes=65536)
Progress (group=search,P2,nextOffset=118,generation=6)
Business identity M17 → order O51,version=3,status=paid

Store records in immutable closed segments and one active append file. A sparse index points to locations near a requested offset. Checksums detect corruption; a committed watermark marks the last record readers may see. Batching spreads network and disk overhead across records but makes them wait briefly. The broker recognizes repeated producer identities and sequences only while it retains that retry state. A new application resend may look like a new operation.

Fetch returns records no later than the committed watermark and includes the next fetch position. A group offset commit carries membership generation 6 and nextOffset 118. The coordinator rejects a commit from generation 5 after reassignment. Clients cannot commit arbitrary future positions without an explicitly privileged reset operation.

Administrative replay uses a new group or a controlled offset reset with an audit reason. Appending a tombstone for key compaction is different from deleting a historical event immediately; compaction and time retention have separate contracts. Payload schema versions travel with records so replaying an old segment does not require today's consumer to guess its encoding.

06Data model and access patterns

The broker must distinguish records that exist in a log, records committed by replication, and records a consumer group declares finished. These three positions can differ. Producer retry state and controller metadata support those positions, but neither substitutes for a consumer’s own business-result record.

State Key Meaning
Partition log topic, partition, offset Ordered retained record bytes
Replica commitment partition, leader epoch, committed offset Last committed record readers may see
Producer deduplication producer incarnation, partition, sequence Which append retry already exists
Group progress group, partition, next offset, generation Records that group finished without a gap
Controller metadata topic, replicas, routing version Current brokers, replicas and partition routes

Log files are divided into segments, with a sparse offset-to-byte index and checksums. Closed segments are immutable until deletion or compaction produces replacements. The active segment accepts batched appends. Segment replacement and deletion are coordinated with active readers so a fetch either retains a valid file handle or retries against the replacement index.

Producer sequence state must survive the failover policy with the log; keeping it only in a leader's memory would duplicate acknowledged retries after election. Its retention and producer incarnation lifecycle are explicit. A restarted application that changes its identity cannot expect the broker to recognize a semantically identical order event automatically.

Group bookmarks belong in durable replicated metadata. A search view's processed-event table belongs to that consumer's own database. The broker saves group progress; the consumer database saves business results. Because these are separate commits, the bookmark alone cannot prove an update happened once.

07Basic working design

One broker accepts a batch, validates its checksum and sequence, appends records to a local file, durably flushes under its acknowledgment policy, and returns the assigned offsets. A consumer fetches from offset 117, applies M17, and stores nextOffset 118. Search can pause while producers continue appending, then resume from that bookmark.

For the initial example, one partition is enough. It preserves every order event's log order and supports multiple groups by keeping separate bookmarks. A small metadata database or durable internal log stores those positions. The server rejects records exceeding its configured batch size rather than allocating an unbounded buffer.

This design establishes retention, batching, fetch, and replay before distribution. Batching several records into one write reduces syscall and flush overhead, but a short maximum batching delay, called linger, prevents a quiet producer waiting indefinitely for a full batch. A consumer fetch similarly waits only up to a bounded timeout or byte threshold.

The broker acknowledgment says that local disk has accepted the record. It does not yet survive that machine's permanent loss. That is a visible limitation of the baseline and the reason to add replication, not an excuse to describe every disk write as automatically durable everywhere.

architecture · baselineOne retained file and group bookmarks

A bookmark belongs to a consumer group; reading does not delete the record.

One retained file and group bookmarksA bookmark belongs to a consumer group; reading does not delete the record. p to b: Append M17; b to disk: Durable append / bookmark; c to b: Fetch from offset 117; c to sink: Apply M17; c to b: Commit next offset 118Append M17Durable append / bookmarkFetch from offset 117Apply M17Commit next offset 118ACTOROrder producerSERVICESingle log brokerSTORELog segments +group offsetsWORKERSearch consumerSTORESearch projectiondatabasesync
Read each connection in order
  1. syncAppend M17Order producer → Single log broker
  2. syncDurable append / bookmarkSingle log broker → Log segments + group offsets
  3. syncFetch from offset 117Search consumer → Single log broker
  4. syncApply M17Search consumer → Search projection database
  5. syncCommit next offset 118Search consumer → Single log broker

08Find the baseline flaws

At 100 MB/s, seven days of raw history require 60.48 TB on the baseline machine. Even if its sequential write bandwidth is sufficient, disk capacity, consumer reads, and rebuild time make it a single operational bottleneck. A permanent disk loss destroys locally acknowledged history. Adding a second read process changes neither fact.

The most tempting consumer mistake is to commit nextOffset 118 before updating O51. If the consumer crashes between those operations, its replacement starts after M17 and search never sees the paid state. Reversing the order avoids that loss but allows repetition: update O51, crash, replay M17. This is why at-least-once consumption requires a duplicate-safe sink: the destination database or service must recognize repeated work and avoid applying the same business change twice.

A producer has a similar uncertainty. The broker commits sequence 44 but its reply disappears. Retrying as sequence 45 creates another append; retrying 44 under retained producer state can return offset 117. Application event M17 still needs its own identity because a later business resend may occur under a different producer incarnation.

Finally, adding partitions with a naive modulo hash can move customer42 while old events remain on P2. New events on P7 could be consumed before older ones. Partition expansion is therefore a routing transition with ordering consequences, not a harmless capacity toggle.

09Improve the design, step by step

First, replicate each partition. We need accepted messages to survive one broker failure. A leader orders appends, a quorum durably replicates them, and election preserves the committed prefix. This adds fault tolerance at the cost of replication bandwidth and acknowledgment latency. A stale or isolated leader cannot commit new records without quorum. Local-only acknowledgment remains an option for explicitly disposable telemetry, not for the stated order-event contract.

Second, introduce keyed partitions and routing metadata. The trigger is one broker's disk, network, or consumer parallelism limit. Assign customer keys to stable partition routes and distribute leaders across brokers. This raises aggregate throughput while keeping one key serial. The cost is more routing metadata, pauses when ownership changes, and uneven load across partitions. A single partition is preferable when low traffic and true total ordering matter. Expansion either preserves existing routes or stops the key's new traffic until its old-partition events have drained, then switches routing at a coordinated barrier.

Third, coordinate consumer ownership by generation. The trigger is worker failure and elastic group size. A group coordinator assigns partitions and increments the generation on reassignment. Fetches, offset commits and consumer work carry that assignment generation. This allows takeover, but rebalances can pause work and stale workers may still reach external sinks. The broker rejects old owners’ offset commits. The destination database or service must separately reject stale versions or repeated business updates. Manual static assignment remains simpler for a small fixed pipeline.

Fourth, tier immutable segments and enforce resource budgets. The trigger is seven-day retention and multiple independent readers competing for hot disks. Keep active/recent segments local, move verified closed segments to object storage, and apply per-tenant append/fetch budgets. This reduces local capacity pressure but adds remote fetch latency, object lifecycle, and restore dependencies. Local-only storage remains preferable for a short replay window with tight historical-fetch latency.

None of these changes makes a remote charge atomic with a group offset. A broker transaction can cover only its documented transaction domain; a database view or payment service needs an explicit integration contract.

10Detailed architecture

Commit and election contract

Choose a failure model, such as surviving one broker loss with three replicas. “Three copies” alone does not say when to acknowledge. The append policy must specify durable replica participation, and leader election/log reconciliation must preserve acknowledged history. An out-of-date replica cannot simply become leader and erase acknowledged records.

Partition leaders carry epochs; old leaders and stale producers are fenced, meaning their outdated authority is rejected. After failure, choose an eligible leader, reconcile records not yet committed, and repair replicas. Controllers must also agree durably on which broker owns each partition. Multi-region replication adds latency or a declared asynchronous recovery-point gap; it does not automatically preserve the local acknowledgement guarantee everywhere.

Producer and consumer paths

Producers authenticate to broker endpoints and use controller metadata to find partition leaders. The controller is itself a replicated metadata authority; partition data replication and controller agreement are distinct paths. Each partition's leader and followers retain independent log copies in different failure domains.

Consumers authenticate their group, obtain current assignments, and fetch committed records from the assigned partitions. The group coordinator stores next offsets and generations durably. Search then writes to its own projection database, which includes a processed-event identity table. The broker has no direct authority over that database's transaction.

Retention and regional recovery

Background segment managers verify checksums, archive eligible closed segments, and delete only under the configured retention/compaction policy. A historical fetch may read a remote segment through the broker rather than granting arbitrary clients bucket credentials. The producer waits for the promised commit before receiving success. Replica repair, archiving, consumer work and offset updates happen separately.

A three-copy placement across zones reduces correlated failure exposure but does not eliminate region loss or operator deletion. Remote replication or backup adds a separately declared disaster-recovery contract.

Kafka implementation boundaries

For a concrete retained-log implementation, Apache Kafka 4.x uses KRaft for controller metadata; ZooKeeper mode was removed in Kafka 4.0. The verified 4.3 documentation and 4.3.1 release provide a current example, without implying that this interview protocol is Kafka’s exact implementation. Kafka topic replication uses in-sync replicas (ISR), distinct from the controller’s Raft quorum. With replication factor 3, acks=all and min.insync.replicas=2, a successful append requires the documented ISR acknowledgment condition and rejects when too few replicas remain. It does not simply mean any two of three replicas, and acks=all does not by itself require each replica to fsync every append. Define the actual process, disk, power-loss and zone failure assumptions before equating that configuration with the exercise’s explicit durable-replica protocol. Keep unsafe leader-election choices out of a no-acknowledged-loss contract.

Kafka producer idempotence suppresses supported protocol retries, while transactions can coordinate Kafka records and consumed offsets within their documented domain. A database projection still needs the separate sink transaction shown here. For Kafka transactional producers, consumers requiring only committed transactions select the documented read_committed behavior; the broker replication watermark alone is not the visibility rule for aborted or still-open transactions.

architecture · finalReplicated partitions and independent consumers

The broker commits records, the group coordinator stores consumer progress and the destination database commits business updates; none of those commits automatically performs the others.

Replicated partitions and independent consumersThe broker commits records, the group coordinator stores consumer progress and the destination database commits business updates; none of those commits automatically performs the others. p to meta: 1. Resolve leader / routing version; p to lead: 2. Append stable producer sequence; lead to local: Append local ordered records; lead to f1: Replicate log; acknowledge durable append; lead to f2: Replicate log; acknowledge durable append; c to group: 3. Join / receive generation; group to offset: Store assignment and next offset; c to lead: 4. Fetch committed prefix; c to sink: 5. Atomic event ID + projection; c to group: 6. Commit completed prefix; other to group: Independent group assignment; other to lead: Independent replay; lead to archive: Archive verified segments1. Resolve leader / routingversion2. Append stable producersequenceAppend local ordered recordsReplicate log; acknowledgedurable appendReplicate log; acknowledgedurable append3. Join / receive generationStore assignment and nextoffset4. Fetch committed prefix5. Atomic event ID + projection6. Commit completed prefixIndependent group assignmentIndependent replayArchive verified segmentsACTORProducers + outboxrelayG1STOREReplicated controllermetadataG1SERVICEPartition leadersG2STOREFollower logs: zone BG2STOREFollower logs: zone CG2STORELeader logs: zone AG2SERVICEGroup coordinatorG3STOREReplicated groupprogressG3WORKERSearch consumergroupG3WORKERIndependentanalytics groupG3STORESearch DB +processed-event IDsG4STOREVerifiedclosed-segmentarchiveG2syncreplicationasyncG1 Production and metadataG2 Partition durabilityG3 Group coordinationG4 External business authority
Read each connection in order
  1. sync1. Resolve leader / routing versionProducers + outbox relay → Replicated controller metadata
  2. sync2. Append stable producer sequenceProducers + outbox relay → Partition leaders
  3. syncAppend local ordered recordsPartition leaders → Leader logs: zone A
  4. replicationReplicate log; acknowledge durable appendPartition leaders → Follower logs: zone B
  5. replicationReplicate log; acknowledge durable appendPartition leaders → Follower logs: zone C
  6. sync3. Join / receive generationSearch consumer group → Group coordinator
  7. syncStore assignment and next offsetGroup coordinator → Replicated group progress
  8. sync4. Fetch committed prefixSearch consumer group → Partition leaders
  9. sync5. Atomic event ID + projectionSearch consumer group → Search DB + processed-event IDs
  10. sync6. Commit completed prefixSearch consumer group → Group coordinator
  11. syncIndependent group assignmentIndependent analytics group → Group coordinator
  12. syncIndependent replayIndependent analytics group → Partition leaders
  13. asyncArchive verified segmentsPartition leaders → Verified closed-segment archive

11Write path and acknowledgement

The append protocol orders records within a partition and preserves acknowledged entries through supported leader changes. Producer retries use stable identities within a stated horizon.

  1. P8 sends M17/sequence 44 for customer42, routed to P2.
  2. Leader epoch 9 appends offset 117 and replicates according to the chosen durable acknowledgement policy.
  3. The producer receives offset 117 after commitment; a lost reply is retried with the same identity.
  4. C12 in group search fetches 117 and updates O51 to version 3 in its projection database, atomically recording processed event M17.
  5. Only after that transaction succeeds does C12 commit nextOffset 118.
  6. Another group can still consume 117 independently; one group’s progress does not delete the event.

An idempotent producer lets the broker recognize supported append retries. Application resends and external updates still need their own safeguards, as Kafka documents. Kafka producer API.

  1. If the producer connection breaks before the result, P8 resends sequence 44. The current leader consults replicated producer state and returns the same committed position when it is still within the supported identity lifetime. It rejects conflicting bytes for the same sequence.
  2. If commitment did not happen before leader failure, the new leader resolves the uncommitted tail under the replication protocol. The producer's retry may now create the one committed record. The application should not infer a result from an old leader's local offset alone.
  3. The order service's own outbox relay marks M17 as published only after it has a durable append result. A crash before that bookkeeping can cause a semantic resend, so the consumer still records M17 even when producer retry suppression is enabled.

An append timeout therefore remains uncertain until retry or lookup resolves it. Search visibility is a later event and is measured independently from producer acknowledgment latency.

12Read and delivery path

Each consumer group reads from its own next offset. Advance the bookmark only after the required business update is safe: either commit them together or make repeating the update harmless.

  1. Consumer C12 joins group search and receives P2 with generation 6 and nextOffset 117. It verifies the assignment before starting a bounded fetch.
  2. The broker locates the segment using its sparse index, validates record framing/checksums, and returns committed records starting at 117. If 117 predates retention, it returns an explicit gap condition rather than offset 118 as if nothing were missing.
  3. C12 validates the payload schema and applies records under the required key ordering. Parallel processing may dispatch independent keys, but it tracks which offsets are completed.
  4. If 119 finishes while 118 is still pending, the completed prefix cannot advance beyond 118. Committing 120 would skip unfinished work after a crash.
  5. After its sink transaction succeeds, C12 commits the next completed prefix using generation 6. If a rebalance has moved the partition to generation 7, the coordinator rejects C12's stale commit.
  6. The replacement starts from durable group progress. Repeated records are normal; the destination uses saved event identities to avoid repeating their updates. A fresh group can independently replay all retained records without disturbing search.

Lag includes the difference between committed broker position and completed consumer position, but byte size and event age matter. Ten large image-reference messages and ten thousand tiny state events do not imply the same recovery cost or freshness.

13Correctness deep dive

The search database makes M17's business effect and its deduplication record atomic. It does not share a transaction with the broker:

apply(event M17, order O51, version 3, status paid):
  begin search database transaction
  insert processed(tenant, stream, M17, fingerprint) if absent
  if already present:
      verify the original fingerprint; return the saved processing result
  update O51 only if incoming version > stored version
  commit
then commit group nextOffset = 118 with current generation

The database commits the event identity and order-view update together. If the transaction fails, neither is saved. An already recorded M17 is a safe no-op even if a new consumer owns the partition. The order-version check also prevents a delayed older event from overwriting a newer projection; it is distinct from event-ID deduplication.

Interleaving Search state Group progress
C12 applies M17 O51 paid/v3; M17 recorded Still 117
C12 crashes before offset commit Same durable state Still 117
C13 takes generation 7 and replays M17 Unique insert finds prior result Still 117
C13 commits completed prefix No second business change 118

If C12 resumes, its generation-6 offset commit is rejected. It might still call the database, so the sink's checks remain necessary. A broker generation check cannot reach into arbitrary external services.

sequence · consumer-raceCrash after effect, before bookmark

The search database recognizes an already processed event and skips its duplicate update. Separately, the group coordinator rejects offset commits from an obsolete consumer generation.

Crash after effect, before bookmarkThe search database recognizes an already processed event and skips its duplicate update. Separately, the group coordinator rejects offset commits from an obsolete consumer generation. a to db: Atomically record M17 + O51/v3; db to a: Committed; a to a: Crash before offset commit; b to g: Acquire P2 at next offset 117; b to db: Replay M17; db to b: Already applied; no second effect; b to g: Commit next offset 118 / generation 7; a to g: Late commit from generation 6; g to a: Reject obsolete generationPARTICIPANTConsumer C12 /generation 6PARTICIPANTSearch databasePARTICIPANTGroup coordinatorPARTICIPANTConsumer C13 /generation 71. Atomically record M17 +O51/v32. Committed3. Crash before offsetcommit4. Acquire P2 at next offset1175. Replay M176. Already applied; no second effect7. Commit next offset 118 /generation 78. Late commit from generation 69. Reject obsolete generationsyncreturnblocked
Read each connection in order
  1. syncAtomically record M17 + O51/v3Consumer C12 / generation 6 → Search database
  2. returnCommittedSearch database → Consumer C12 / generation 6
  3. syncCrash before offset commitConsumer C12 / generation 6 → Consumer C12 / generation 6
  4. syncAcquire P2 at next offset 117Consumer C13 / generation 7 → Group coordinator
  5. syncReplay M17Consumer C13 / generation 7 → Search database
  6. returnAlready applied; no second effectSearch database → Consumer C13 / generation 7
  7. syncCommit next offset 118 / generation 7Consumer C13 / generation 7 → Group coordinator
  8. syncLate commit from generation 6Consumer C12 / generation 6 → Group coordinator
  9. blockedReject obsolete generationGroup coordinator → Consumer C12 / generation 6

14Failure and recovery

Failure or race Required response and boundary
Consumer crashes after sink commit C12 commits the projection update and then crashes before advancing group progress. C13 takes over and reads M17 again. Its database transaction finds M17 already recorded, leaves O51 unchanged, and safely advances to 118. Without that check, an increment or external charge could happen twice.
Parallel work and poison events Parallel consumers must commit only the completed prefix: finishing 119 while 118 is unfinished does not permit nextOffset 120. Assignment generations reject old workers’ offset commits, but the destination must still check versions or deduplicate external updates. Poison events require bounded retries and an explicit quarantine/dead-letter decision. Skipping a failed event can violate later per-key business ordering; document whether to pause that key/partition.
Leader partition During a leader partition, the quorum side may elect an eligible replacement while the isolated leader stops committing. Clients refresh metadata after errors; they do not write independently to any reachable replica. On rejoin, the former leader reconciles its uncommitted suffix and repairs from the authoritative history.
Disk reserve or retention exhausted When disk reserve falls below a safe threshold, producers receive bounded backpressure or rejection before new acceptance. Existing retained data remains readable. A slow consumer nearing the seven-day boundary triggers an alert early enough to add net catch-up capacity or export a replay archive. After retention actually passes its bookmark, recovery requires a rebuild decision; silently resetting to the newest offset would conceal data loss in the view.

15Operations, security, and cost

Encrypt transport and authenticate producer and group identities. Topic permissions separate append, fetch, and administrative reset; replaying another tenant's stream must not be a side effect of guessing a topic name. Limit batch size, decompressed size, connections, append bandwidth, and fetch bandwidth. Compression savings are useful only with bounded decompression memory and CPU.

Monitor append p99, committed-to-consumed event age, replica catch-up lag, unavailable partitions, disk reserve, group churn, and time remaining before a consumer loses its oldest required segment. A healthy average append latency can hide one hot partition whose customers are hours behind.

Test leader death before and after quorum commitment, lost append replies, consumer death after sink commit, stale group commits, corrupt segments, and controller restoration. A partition-count migration should include a key-ordering test across the old and new route, not only a successful administrative API call. Before changing a consumer’s schema, test it by replaying retained events in the older formats.

At the assumed workload, seven days of three-copy storage is 181.44 TB before indexes and headroom. Reserving an additional day for repair adds 25.92 TB of replicated payload. That makes retention and compression material capacity decisions. A measured 2:1 compression ratio would halve payload bytes, but cannot be assumed for already compressed or encrypted content.

16Decision ledger and limitations

Choice Benefit Limitation
Retain all events for a time Replay and independent groups Slow consumers can fall behind deletion
Compact by key Retain latest state economically Historical event sequence is incomplete
More partitions More independent work Metadata/rebalance cost and routing transitions
Producer backpressure Protect finite disk/buffers Higher upstream latency or rejections

The design preserves each key’s order and lets groups replay independently. Consumers may repeat business updates unless the destination makes those retries safe. More partitions improve aggregate parallelism but cannot accelerate one strictly ordered hot key. A task queue with visibility leases would better fit unrelated long jobs needing individual rescheduling; it would not replace the independent replay history requirement.

The chosen order-event topic retains the full sequence for seven days and does not compact away intermediate events within that window. Compaction is an alternative for topics whose consumers need the latest state per key, but it cannot preserve the full sequence of order transitions for an audit. Time retention offers a clear replay horizon at substantial storage cost. Remote archival extends that horizon with slower reads and another recovery dependency.

Our next scale trigger is measured partition skew or consumer recovery time, not a desire for a larger partition count. Our next correctness trigger is a sink that cannot make repeat processing safe; that requires a business workflow change rather than a broker setting.

17Interview closing

“I designed a retained event log so order acceptance does not depend on the search service being online. Producers route an ordering key to one partition, receive a committed offset under a stated replicated durability policy, and retry uncertain appends with the same identity. Each consumer group tracks how far it has finished without gaps and can replay seven days of history.

“The hard boundary is between consuming a record and changing an external system. A projection consumer commits the unique event identity and its business-state update in one sink transaction, then advances its group offset. A crash between those commits causes replay, but the unique event identity and order version prevent another business change. Group generations reject obsolete offset commits; they do not magically fence arbitrary external writes.

“The workload creates 60.48 TB of raw history per week and 181.44 TB across three copies. I would benchmark partition throughput and reserve disk and net catch-up capacity before choosing a partition count. The first recovery drill kills a consumer after its sink commit and a leader after quorum commitment, then checks both the retained history and search result.”

If the interviewer demands global order, I would begin with one ordered partition or a sequencer and quantify the throughput and availability cost. If they need millions of independent long tasks instead, I would redesign around task identities, leases, and per-message retry state rather than pretending a log offset provides that interface.

Practise the interview questions

Say your answer aloud before opening the model answer. Then answer the follow-up and compare the reasoning.

Foundation · Question 1

What does committing nextOffset 118 mean?

Reveal a model answer

It means that this consumer group has safely completed the required work through 117 and should resume at 118. It does not delete 117 for other groups, and it is not a global time value. The partition identifies the ordered sequence in which that number has meaning.

What the answer must demonstrate: Define position and ownership explicitly.

Applied · Question 2

A consumer applies event M17 to its database projection but crashes before committing the broker offset. How can the replacement consumer replay safely?

Reveal a model answer

The replacement rereads the event because the saved group bookmark did not advance. The projection transaction checks the same event identity M17 and finds the effect already applied, so it performs no duplicate mutation. It can then commit the next completed offset. This makes replay safe at the business sink.

What the answer must demonstrate: The broker cannot atomically control an arbitrary external side effect.

Applied · Question 3

The append reply was lost. Should P8 send a new event ID?

Reveal a model answer

No. It retries the uncertain append with the same logical identity and supported producer sequence. A new identity can become a second valid event even if the original append succeeded. I would keep producer-retry semantics distinct from a user genuinely creating another order.

What the answer must demonstrate: Do not extend provider guarantees beyond their stated scope.

Foundation · Question 4

Why can’t you freely spread customer42 across every partition?

Reveal a model answer

Its events may race and be observed in different orders because partitions have independent sequences. If customer ordering matters, route the key consistently or add an explicit sequencer/reassembly protocol. One key cannot use independent workers freely while still requiring all its events to stay in order.

What the answer must demonstrate: Distribution changes ordering semantics.

Follow-up · Question 5

Workers finished 117 and 119, but 118 is still running. What can they commit?

Reveal a model answer

Only nextOffset 118, representing the contiguous completed prefix through 117. Advancing to 120 would skip unfinished 118 after a crash. I track gaps and advance the bookmark when the earliest outstanding work becomes safely complete.

What the answer must demonstrate: Parallel execution is not permission to skip progress gaps.

Follow-up · Question 6

A consumer is eight days behind a seven-day log. What do you promise?

Reveal a model answer

Its required records may already be deleted. I would alert well before the retention margin is exhausted and provide a documented recovery path, such as archived replay or a fresh application snapshot. I cannot claim retained-log durability means unlimited history.

What the answer must demonstrate: Use time/byte lag and catch-up capacity, not message count alone.

Applied · Question 7

Why can adding partitions break ordering for customer42?

Reveal a model answer

A modulo partitioner may route new events to a different partition while older events remain on the previous one. Independent consumers can then process new before old. I keep existing key routes, or pause new traffic for the key, finish its old-partition records, then switch its routing version.

What the answer must demonstrate: Partition count changes can change semantics, not only capacity.

Follow-up · Question 8

A consumer is ten minutes behind at 100,000 messages/s. It can process 150,000/s. How long to recover?

Reveal a model answer

It has 60 million messages of backlog and a net drain rate of 50,000/s after current arrivals, so it needs twenty minutes. I verify fetch, sink, and retention budgets sustain that excess rate.

What the answer must demonstrate: Subtract ongoing arrivals when calculating recovery.

Blank-page exercise · 45 minutes

Build the answer yourself

Design the producer’s retained order-event log. Trace M17/offset 117, lose the producer reply, then crash the consumer after updating its database but before committing 118.

  • Explain one-server append and independent group bookmarks.
  • Calculate seven-day retained storage and catch-up work.
  • Specify per-key order and durable acknowledgement policy.
  • Resolve uncertain producer retries and consumer replay.
  • Demonstrate a progress gap and retention overrun.

Check that each component and design decision follows from your requirements and workload.

Recall the key ideas

Answer from memory before opening each card. Explain why the choice works and what it costs. Revisit missed cards tomorrow.

Design a distributed message logWhat is an offset?Recall first, then reveal

A position within one partition’s ordered retained records, not a universal timestamp.

Partition + position.

Return to lesson
Design a distributed message logWhen may progress advance?Recall first, then reveal

After the required updates are durable and every earlier record in that partition has also completed; never advance past unfinished work.

Finish safely, then move the bookmark.

Return to lesson
Design a distributed message logDoes replay mean a business action repeats?Recall first, then reveal

It may unless the destination database or service deduplicates the event or atomically records the update with the input event's identity.

Delivery can repeat; effects must have a rule.

Return to lesson

Final revision

Summary and interview notes

A retained log stores messages separately from each consumer’s progress and business updates. Each partition has its own order. Safe replay needs stable retry identities at the broker and destination, and bookmarks that never skip unfinished work.

Remember these points

  • An offset belongs to one partition, and a committed group offset does not delete data for other groups.
  • Retry an uncertain append under the same supported producer identity; business resends may need longer-lived event deduplication.
  • Commit sink effect and processed-event identity together, then advance only the contiguous completed prefix.
  • Routing changes can reorder a key across partitions; preserve routes or perform an explicit handover.
  • Replication, transaction visibility and fsync are separate guarantees; inspect the selected broker’s exact policy.

Interview tips

  • Kill the consumer between sink commit and offset commit, then show the replacement replay.
  • Calculate weekly replicated bytes and net catch-up rate before selecting partition count.
  • Explain why a version guard works for replacement state but can lose unapplied deltas.

Important qualifications

  • Kafka 4.x uses KRaft metadata; its ISR-based data acknowledgment is not a generic majority-fsync algorithm.
  • Transactions inside a broker do not make arbitrary HTTP or database effects atomic.

Technical references

Practice marks stay in this browser.