System-design interview · Extended interviews
Design a distributed message log
Design a retained partitioned event log with durable producer acceptance, independent consumer groups and replay-safe database effects, while keeping broker and application guarantees separate.
You will learn to
- Explain partitions, retained offsets and independent consumer progress using one order event.
- Trace safe producer retries and the consumer effect-before-offset boundary.
- Size retention, partition skew and catch-up without promising global order or arbitrary exactly-once effects.
Practice in this chapter
8 interview questions with model answers and follow-ups.
Go to interview practiceUseful 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.
Dotted concept links open the relevant explanation in a new tab.
01Choose a retained log rather than an unspecified queue
The order service emits event M17 when order O51 becomes paid. Search, fraud analysis and analytics must process it independently, including after downtime. Choose a retained event log: records remain for a defined period while each consumer group maintains its own position. A worker finishing a record does not delete it for other groups.
This differs from a task queue where workers compete for individual jobs with visibility deadlines and per-message acknowledgments. Both are useful, but combining their terms without a contract makes ordering, replay and failure behavior unclear. This design offers topics, partitions, producer append, group fetch, progress commits and seven-day replay.
Require order for related events routed by an ordering key, such as customer42. Unrelated customers need no global order. The broker's acceptance means the record was saved under the configured durability policy; it does not mean search applied the order update or a downstream email was sent.
Use at-least-once consumer processing: a failure may cause a record to be read again. Consumers must make their relevant effects safe to repeat. The order service also needs its own outbox or equivalent reliable publication so a committed order change cannot disappear before an event ever reaches this broker.
First ask whether consumers need independent replay of an event history or competing ownership of individual tasks. Confirm the ordering key, retention period and acknowledgment failure model. Choose a seven-day retained log with independent groups and per-key order here; consumer side effects remain a separate contract.
02Functional requirements
Publish records. Accept bounded producer batches into topics/partitions and return committed record positions under a documented retry identity.
Consume independently. Let each consumer group fetch records and durably record its next required offset without deleting data for other groups.
Replay history. Allow authorized groups to resume or reset within seven days of retained full history and report an explicit gap when history has expired.
Manage consumers and access. Assign partition ownership within groups, recover progress after reassignment and enforce topic/group permissions and administrative replay controls.
03Non-functional requirements
These are illustrative interview assumptions, not product facts or measured benchmarks. Confirm them before choosing components, then validate the completed design under the stated workload. Here p95 means the 95th-percentile latency: 95% of measured requests take no longer than that value. Report errors and rejected work alongside latency; a fast failure is not a successful outcome.
Workload. Plan for 100,000 one-kilobyte messages/s, or 100 MB/s ingest. Each full-rate consumer group adds roughly 100 MB/s of reads; the example’s search, fraud and analytics groups must be budgeted separately.
Broker response time. Target durable append acknowledgment within 50 ms p95 and bounded fetch responses within 100 ms p95 for already available records under admitted same-region load. Include broker queuing and batch wait; empty long polls and downstream processing have separate timing contracts.
Ordering and processing. Preserve committed order within a stable key-to-partition route. Provide at-least-once consumer processing; producer retry suppression and consumer effect deduplication protect different boundaries, with no global order or arbitrary exactly-once external-effect promise.
Durability and availability. Place three replicas across three zones and acknowledge only a durably committed majority. Preserve accepted history after one replica or zone failure; stop unsafe appends without a majority and treat regional loss separately.
Retention and bounded resources. Keep full history for seven days, even if a consumer falls behind. Reserve storage/repair capacity, reject excess new work before risking accepted history, and authorize append, fetch and offset resets.
04Append a file and remember separate bookmarks
Start with one broker, one append-only log file and durable consumer bookmarks. A producer submits a bounded batch. The broker validates it, appends records, flushes according to its local acknowledgment policy and returns their offsets. An offset is a position in that partition's history, not a global event identifier.
The search group fetches M17 at offset 117. After processing it, search saves nextOffset 118. Fraud can remain at offset 90 and analytics at 120 without changing search's progress. Restarting search uses its saved bookmark to resume. Retention, not one consumer's acknowledgment, eventually removes old records.
Split the file into segments so expired old history can be removed without rewriting the active tail. A sparse offset index locates nearby bytes, and checksums detect damaged records. Batch appends and fetches to reduce network and storage overhead, with bounded waiting so a quiet stream does not wait forever for a full batch.
This baseline explains the entire producer/read/replay interface. Its acknowledged records survive a process restart if local recovery succeeds, but not permanent loss of its only disk. That explicit limit motivates replication; distributing readers alone would not protect the stored history.
Finishing in one group leaves the event available to the others until retention removes it.
Read each connection in order
- syncM17 with stable identitiesOrder outbox producer → Append broker
- syncAppend and acknowledge policyAppend broker → Retained segments
- syncFetch from search offsetSearch group → Append broker
- syncFetch from fraud offsetFraud group → Append broker
- syncNext completed offsetSearch group → Durable group bookmarks
- syncIndependent next offsetFraud group → Durable group bookmarks
05Price retained history and independent readers
Assume 100,000 messages/s at one decimal kilobyte each: 100 MB/s of raw ingest. One day holds 8.64 TB and seven days 60.48 TB. Three copies require 181.44 TB before indexes, compression, spare space and repair traffic. Retention is a substantial storage commitment, not a free property of using a log.
Three copies write about 300 MB/s of payload across replicas; followers add roughly 200 MB/s of replication traffic beyond the producer stream. Each full-rate consumer group reads another 100 MB/s. With several groups, delivery bandwidth can dominate client-facing traffic even when sequential append is efficient.
A consumer stopped for ten minutes accumulates 60 million messages. At 150,000 processed messages/s while 100,000/s new messages continue, net catch-up is 50,000/s and recovery takes twenty minutes. At exactly the arrival rate, it never catches up. Track lag age and bytes as well as record count.
If a tested partition handles 5 MB/s safely, ingest alone suggests at least twenty partitions before skew and headroom. More partitions can enable more consumers, but one key generating 20 MB/s cannot be divided among them while retaining its original strict order. Benchmark the actual payload and destination processing before choosing partition counts.
06Separate append identity from processing progress
Broker operations
Append a record
append(
topic,
key,
producer,
sequence,
event
)
Proposes a record under the producer retry contract.
Fetch a bounded portion of committed history
fetch(
partition,
fromOffset,
maxBytes
)
Reads a bounded committed portion of retained history.
Returned position and stored records
| Record | Fields | Meaning |
|---|---|---|
| Append result | partition, offset |
Identifies the committed record position. |
| Group bookmark | group, partition, nextOffset, generation |
The next position that group still needs to process. |
| Application event | eventId, aggregateId, version, payload |
Stable business meaning that survives a new producer session. |
| Partition metadata | Current leader, replicas and routing ownership | Identifies the current partition owner and copies. |
Worked consumer progress after M17 is processed
group: search
record: M17
offset: 117
nextOffset: 118
Producer identity and sequence distinguish an uncertain retry from a new append within the supported lifecycle. A lost reply should be retried with the same identity; assigning a new sequence can append a second record. The retry state must survive supported failover, not live only in one leader's memory.
M17 also keeps an application event ID because an outbox may resend it from a restarted producer. Broker retry suppression cannot automatically recognize the same business event under a new producer identity.
Include payload schema versions for consumers replaying older data. Authenticate topic append/fetch permissions and group access. Administrative offset reset is a separate privileged operation; an ordinary consumer must not conceal unfinished work by jumping to an arbitrary future position.
07Replicate a committed partition history
For this design, place three replicas of each partition across three zones and use a proven consensus log. One current leader orders appends. Return an accepted offset only after a majority has durably recorded the committed append; safe election must preserve that committed history. This tolerates loss of one replica or zone. Without a majority, stop appends. An isolated previous leader must not acknowledge a separate history. Regional loss needs a separate recovery plan.
Check that the selected storage and deployment configuration actually supports this durability promise. Replica count alone does not identify when data reached durable storage or which replica can become leader. For example, a broker setting that waits for replica acknowledgment is not automatically a guarantee that every replica synchronously flushed each record to persistent media.
Readers consume only the protocol's committed visible records. Transactional broker features can add a further distinction between replicated data and records belonging to completed rather than aborted transactions. Explain the chosen product's supported visibility rule instead of treating every local tail offset as readable success.
Partition leaders and controller metadata have different jobs: the controller identifies placement, while the partition replicas hold event bytes. Replicate the control state too, but do not mistake a healthy directory for retained user data. During insufficient safe replica participation, reject or delay appends rather than issue an accepted offset that may disappear after failover.
08Commit a business effect before advancing its bookmark
Search reads M17 and must update order O51. If it commits nextOffset 118 first, then crashes before changing search state, its replacement skips M17 forever. Reversing the order avoids that loss but introduces repetition: update O51, crash before offset commit, then read M17 again.
Make repetition safe at the destination:
In one search-database transaction, insert the unique processed event identity and update O51. If M17 is already recorded with the same payload, return the existing processing result.
Commit that transaction.
Only then advance the broker bookmark.
A crash between those commits causes harmless replay rather than a missing or repeated business change.
For this example, M17 carries complete order state at version 3. A version guard prevents delayed version 2 from replacing newer state. That guard is not a substitute for all event semantics: an unapplied additive delta cannot simply be discarded because a later numbered delta arrived first. Delta consumers need contiguous sequence handling or a rebuild from complete state.
The broker and search database do not share a transaction here. Stable event identity makes the gap safe. Kafka-style transactional processing can cover cooperating broker/state/output boundaries, but an arbitrary external database or provider still needs an appropriate integration contract.
A processed-event record and search update share one transaction; the broker bookmark follows.
Read each connection in order
- syncCommit M17 identity + O51 version 3Search consumer → Search database
- blockedCrash before saving next offsetSearch consumer → Broker group progress
- syncResume at offset 117Replacement consumer → Broker group progress
- syncM17 already applied; safe no-opReplacement consumer → Search database
- syncCommit nextOffset 118Replacement consumer → Broker group progress
09Partition keys and coordinate group ownership
Route each ordering key to a stable partition and spread partition leaders across brokers. Independent partitions add aggregate throughput and consumer parallelism. Within one consumer group, assign each partition to one current owner; different groups read the same partition independently. More group members than partitions do not create additional ownership slots for that group.
A coordinator increments the assignment generation during reassignment. Offset commits carry that generation so a paused old consumer cannot overwrite current progress after takeover. The old process may still write to an external database or service, so that destination must also deduplicate events or check versions.
Process a partition sequentially in the simplest consumer. If later parallelizing independent keys, commit only a fully completed prefix. Finishing offset 119 while 118 is still running does not permit nextOffset 120; a crash would skip 118. Track completion gaps explicitly or retain sequential processing where throughput allows.
Changing partition count can change a key's route. New customer42 events on a different partition might overtake older retained ones. Preserve existing routes or stop new events on the old route, finish processing its earlier events, then switch routes. Adding partitions is therefore an ordering decision as well as a capacity change, not an invisible modulo adjustment.
Producers discover the current partition leader and append using a stable ordering key. The leader acknowledges only after the chosen durable-majority commit. Consumer groups read committed history independently, apply sink effects safely and commit generation-checked progress through the coordinator; the control plane does not contain the event bytes.
Read each connection in order
- syncDiscover partition ownerProducers → Replicated group / placement control
- syncAppend by ordering keyProducers → Partition leaders
- replicationReplicate partition historyPartition leaders → Two followers per partition
- returnAccepted committed offsetPartition leaders → Producers
- syncAssignment / fenced offsetsIndependent consumer groups → Replicated group / placement control
- syncFetch committed recordsIndependent consumer groups → Partition leaders
- syncReplay-safe business effectIndependent consumer groups → Consumer-owned sinks
10Make replay and poison-event policies visible
Retain the chosen full event history for seven days. A group's slow progress does not automatically extend that horizon. Alert before its oldest required segment expires so operators can add net catch-up capacity or preserve an archive. Once the bookmark predates available history, return an explicit gap and choose a rebuild or restoration path. Silently starting at the newest offset would conceal data loss in the projection.
Compaction is a different policy: it can retain the latest value for a key while removing intermediate history. That may suit rebuilding current state but cannot provide the same full sequence of transitions. A deletion marker under compaction also has retention rules; it is not an immediate erasure of every historical copy.
A malformed or repeatedly failing event needs bounded retry and an explicit quarantine/skip policy. Skipping it may violate later ordering for that key. Pause the key or partition when necessary rather than calling every dead-letter action harmless. Keep event identity, reason and replay controls so repair does not manufacture new unrelated business events.
At larger retention sizes, remote immutable segments can reduce local disk cost but add historical-read latency and object lifecycle dependencies. The first design can remain on local replicated storage when the replay window and measured capacity permit it.
11Trace the two independent timeout boundaries
| Failure | Recovery |
|---|---|
| Append commits, producer reply disappears | Retry the same supported producer identity/sequence and recover its position. |
| Consumer effect commits, bookmark reply disappears | Replay M17; the sink's processed-event transaction returns its existing outcome. |
| Consumer owner changes | New owner starts at durable group progress; stale generation commits are rejected. |
| Partition leader fails | Elect a valid successor under the chosen replication protocol and refresh client metadata. |
| Disk reserve runs low | Backpressure or reject before accepting more data than the retention promise can preserve. |
These are separate guarantees. Producer retry suppression limits duplicate appends; consumer sink deduplication limits duplicate effects. Neither removes the other boundary. An email or payment consumer needs a stable external action identity or reconciliation after uncertain responses, because a local processed-event row cannot force a remote provider to behave transactionally.
Bound producer buffers, fetch sizes and in-flight batches. Retry with backoff and jitter rather than sending unlimited uncertain appends to every replica. Repair and consumer catch-up compete with normal traffic, so reserve throughput and storage space for them. Keep accepted retained records readable during overload instead of deleting promised history merely to improve new-write availability.
12Measure lag against the remaining replay window
Measure admitted durable appends against 50 ms p95 and available-record fetches against 100 ms p95, including broker queuing and batch wait. Separately monitor empty long polls, committed-to-consumed event age, replica lag, unavailable partitions, disk reserve, group churn and time until an old required segment expires. A global average can hide one hot key whose downstream view is hours behind.
Test leader loss before and after commitment, uncertain producer replies, consumer failure after sink commit, stale group commits and retention gaps. Replay old schema versions during consumer upgrades. A partition-count migration should test per-key event order across the old and new routes, not only whether an administrative call succeeded.
Use topic and group permissions, encrypted transport and per-tenant append/fetch limits. Bound both compressed and decompressed sizes so a small compressed batch cannot exhaust broker memory. Administrative replay and reset operations need audit records because they can deliberately repeat or skip business work.
The main tradeoff is retained replay at substantial storage cost. Compression may help, but measure it for the actual payload; encrypted or already compressed records may shrink little. One strictly ordered hot key and external non-idempotent effects remain real limits. If the requirement instead becomes millions of unrelated long jobs with individual retries, choose a task-queue contract rather than forcing them into a retained-log bookmark model.
13Check the design against its requirements
Use the numbered requirements to check the final design. FR refers to the functional list; NFR refers to the non-functional list. Performance rows specify tests still required, not achieved benchmark results.
| Requirement | Design mechanism | Verification and remaining limit |
|---|---|---|
| FR1; NFR2–4: committed append | Stable producer identity/sequence, ordered consensus log and durable-majority acknowledgment. | Lose the append response and then a zone. Recover the same accepted record; load-test 50 ms p95 with replication and batching enabled. |
| FR2,4; NFR3: independent safe progress | Per-group bookmarks, current assignment generations and effect-before-offset processing. | Crash search after its sink commit and reassign the partition. Replay safely and ensure another group’s progress does not change. |
| FR2–3; NFR1–2,5: readable retained history | Partitioned segment storage, bounded fetches and independent group bandwidth budgets. | Load-test three full-rate groups and 100 ms p95 fetches for available records; test seven-day replay and an explicit expired-offset gap. |
| NFR3–5: honest limits during overload | Stable key routes, quotas, storage reserve and net catch-up headroom. | Inject one hot key, majority loss and a ten-minute consumer outage. Refuse unsafe appends and calculate catch-up; adding partitions cannot split one key’s required order. |
14Rapid revision
Remember: Save the consumer effect before advancing its offset, and save the event ID with that effect so replay is harmless.
| Concern | Complete mechanism |
|---|---|
| Product | Retain events so independent groups can read them; one worker’s completion does not delete them. |
| Ordering | Commit events in order within each partition; keep related keys routed to the same partition. |
| Producer retry | Reuse producer identity/sequence within its supported lifetime; retain the business event ID if a restarted producer resends it. |
| Durability | Acknowledge only at the chosen durable commit point; elect leaders that preserve those writes after supported failures. |
| Consumer effect | Save the processed event ID and database update in one transaction, then advance the group’s position. |
| Parallel work | Advance only through consecutive completed events; stop at the first unfinished event. |
| Reassignment | The broker rejects former owners’ progress updates; the destination must separately reject stale writes. |
| Retention | Set the replay deadline and how older readers recover; compaction keeps latest keyed state rather than every event. |
| Capacity | Budget replicated copies and consumer reads; catching up requires capacity beyond new arrivals. |
Close by losing M17's producer reply and then crashing search after its database commit. Explain which system recognizes each retry: the broker for appends, the destination database for applied effects. This is more precise than promising “exactly once delivery” without naming which storage or external effect the claim covers.
Practise the interview questions
Say your answer aloud before opening the model answer. Then answer the follow-up and compare the reasoning.
Why choose a retained log for search, fraud and analytics?
Reveal a model answer
Each service needs an independent position and replay history. One consumer finishing does not remove the record for the others.
Interviewer follow-up
When is a task queue more suitable?
Reveal the follow-up answer
Independent jobs needing per-message leases and retries rather than shared retained history.
What the answer must demonstrate: Choose independent replay history or per-message task semantics before architecture.
An append reply is lost. Why reuse the sequence?
Reveal a model answer
The append may already be committed. Supported retry identity lets the broker recover the same position instead of treating it as a new append.
Interviewer follow-up
Why retain an application event ID too?
Reveal the follow-up answer
A later resend of the same event may come from a restarted producer with a new broker identity.
What the answer must demonstrate: Reuse supported append identity and retain business identity across producer incarnations.
What if search commits its update and then crashes before its offset?
Reveal a model answer
The replacement replays the event. A unique processed-event identity committed with the update makes that replay harmless.
Interviewer follow-up
Why not commit the offset first?
Reveal the follow-up answer
A crash before the effect would then skip it permanently.
What the answer must demonstrate: Commit sink deduplication and effect together before advancing broker progress.
Offset 119 finishes before 118. Can the group commit 120?
Reveal a model answer
No. The bookmark must represent a completed prefix; otherwise recovery skips unfinished 118.
Interviewer follow-up
What is the simpler alternative?
Reveal the follow-up answer
Process the partition sequentially until throughput justifies explicit gap tracking.
What the answer must demonstrate: Advance only through the fully completed partition prefix.
What does a consumer generation protect?
Reveal a model answer
It lets the coordinator reject progress commits from a superseded assignment.
Interviewer follow-up
Does it stop an old process writing a database?
Reveal the follow-up answer
Not by itself. The sink still needs idempotency, versions or another enforced ownership contract.
What the answer must demonstrate: Distinguish broker-generation fencing from external sink protection.
Why can increasing partition count change ordering?
Reveal a model answer
A key may move while older records remain on its previous partition, allowing new events to overtake old ones.
Interviewer follow-up
How can it be handled?
Reveal the follow-up answer
Preserve stable routes or coordinate a transition that respects the old ordering boundary.
What the answer must demonstrate: Preserve per-key ordering across a route or partition-count change.
A consumer is eight days behind a seven-day log. What should happen?
Reveal a model answer
Return an explicit retention gap and use an agreed archive or projection rebuild path. Do not silently skip to the newest event.
Interviewer follow-up
Does compaction preserve every transition instead?
Reveal the follow-up answer
No. It can preserve latest state while removing intermediate history.
What the answer must demonstrate: Expose retention gaps and distinguish time history from latest-key compaction.
Does broker transactional processing make an arbitrary payment exactly once?
Reveal a model answer
No. Its guarantee covers only the storage and outputs that participate in its documented transaction protocol; a payment provider needs its own stable identity and uncertain-outcome recovery.
Interviewer follow-up
Does an outbox remove all uncertainty?
Reveal the follow-up answer
It preserves the intent, but a remote effect can still happen before its reply is recorded.
What the answer must demonstrate: Limit exactly-once claims to cooperating storage/effect domains.
Blank-page exercise · 45 minutes
Build the answer yourself
Design a retained order-event log, then lose an append acknowledgment and crash a consumer after updating its database but before committing its offset.
- Agree numbered functional and non-functional requirements, including retained-log semantics, ordering, broker latency and replay durability. Then explain partitions, retained offsets and independent consumer progress using one order event.
- Trace safe producer retries and the consumer effect-before-offset boundary.
- Size retention, partition skew and catch-up without promising global order or arbitrary exactly-once effects.
- Trace a timeout and a concurrent request using the actual durable records.
- Review the final architecture against every numbered requirement, including measured bottlenecks, targets still needing validation and remaining failure limits.
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 logSearch saved M17 at offset 117, then crashed before saving nextOffset 118. What happens on restart?Recall first, then reveal
The group bookmark causes M17 to be read again. The processed-event record saved with the database update prevents applying it twice. An offset identifies position, not completed external work.
Save the effect before progress; make replay harmless with a saved event ID.
Return to lessonDesign a distributed message logWhat prevents a replayed event from updating the database twice?Recall first, then reveal
Save the processed event ID and its database change in one transaction; skip an ID already recorded.
Effect and dedupe together.
Return to lessonDesign a distributed message logCan extra partitions split the work for one strictly ordered key?Recall first, then reveal
No. Keeping that key’s order requires one sequence; extra partitions spread other keys.
More partitions spread different keys; one strictly ordered key still has one sequence.
Return to lessonFinal revision
Summary and interview notes
Keep events in a replicated partitioned log so independent consumer groups can read and replay them. Save each consumer’s database update safely before advancing its position; broker acceptance alone does not complete that update.
Remember these points
- Keep broker acceptance separate from consumer effects.
- Reuse producer identity for uncertain appends and event identity for repeated consumer processing.
- Keep each key ordered and advance group progress only past records whose processing finished.
- Budget retained copies and net catch-up capacity.
Interview tips
- Separate producer acceptance, consumer progress and destination effects.
- Lose M17’s append reply, then crash search after saving O51 but before its bookmark.
Important qualifications
- Traffic and latency figures are interview assumptions, not claims about a named company's deployment.
Continue after the core interview
Explore the advanced version
The advanced lesson keeps the full detailed design. Use these sections when you want to examine the stronger requirements and failure cases.
- Detailed replica and Kafka configuration boundaries
Exact acknowledgment, storage and transaction visibility depend on the selected implementation.
- Consumer takeover effect proof
Study sink uniqueness, full-state versions and additive-delta differences in depth.
- Routing expansion and tiered segments
Large deployments need deliberate key-route migration and historical-segment lifetime rules.
- Retention and failure drills
Explore poison events, stale workers and resource exhaustion beyond the main example.
Technical references
- Apache Kafka designPrimary discussion of partitioned logs, replication, consumer progress, and delivery semantics.
- Kafka 4.3 producer API documentationVerified 4.3.1 Javadoc: producer retries, idempotence and Kafka transaction boundaries; not an external-effect guarantee.
- Kafka 4.3 KRaft operationsOfficial current metadata-controller deployment model; distinct from data-partition replication.
- Kafka 4.3.1 release announcementVerified June 25, 2026 release; used to establish that the 4.3 API example is released, not merely a draft documentation selector.
Practice marks stay in this browser.