Distributed Logs and Ordering Guarantees

LESSON

Consensus and Coordination

010 30 min intermediate

Distributed Logs and Ordering Guarantees

By the end of this lesson, you will be able to...

  • Trace how a keyed event becomes an ordered record and a consumer advances through that record's partition.

  • State the scope of an ordering guarantee instead of calling a whole event platform “ordered.”

  • Identify why replay, retries, and cross-partition workflows need rules beyond append order.

Idea in one sentence: A distributed log gives durable ordered histories, but the order belongs to a named boundary—usually one partition—not automatically to every event in the system.

Core Insight

An order service publishes events for order-7: it is placed, inventory is reserved, payment is captured, and fulfillment can begin. Billing, fulfillment, fraud detection, and analytics all read the events.

It is tempting to say, “Put events in a distributed log; the log keeps them in order.” That is a good start for one append-only sequence. It breaks when the service scales. A partitioned log contains several ordered sequences at once, and consumers can make progress through them at different speeds.

The important design question is therefore not “Does this platform support ordering?” It is: which events must share an order, and what key places them inside that same ordered history?

For the order workflow, order_id is a useful key. It lets every event for order-7 occupy one partition and therefore one local sequence. It does not create a total order between order-7 and every other order, and it does not make retries or external side effects disappear.

The Small Situation: Two Orders, Two Histories

Use a topic called order-events with two partitions. The offsets below are illustrative positions, not timestamps.

partition 0                         partition 1

offset 88: OrderPlaced(order-8)     offset 420: OrderPlaced(order-7)
offset 89: PaymentCaptured(order-8) offset 421: InventoryReserved(order-7)
                                    offset 422: PaymentCaptured(order-7)

The producer uses order_id as its partition key. All order-7 records go to partition 1; all order-8 records happen to go to partition 0. Within partition 1, 420 -> 421 -> 422 is a meaningful append order.

The number 422 does not say that PaymentCaptured(order-7) happened after PaymentCaptured(order-8) at offset 89. Offsets belong to different histories. Comparing their integers across partitions is like comparing page 89 in one book with page 422 in another.

This small fact prevents a large class of accidental guarantees. The platform can process many orders in parallel because it does not first establish one global sequence for all of them.

The Initial Model: A Topic Is One Ordered Queue

The first model works for a single unpartitioned log: producers append records, and every consumer reads positions in order.

0 -> 1 -> 2 -> 3 -> 4

It becomes incomplete after partitioning. A topic is now a collection of logs. A consumer group may own one or more partitions, and each partition is assigned to one group member at a time. The group can process partition 0 and partition 1 concurrently. That parallelism is the reason to partition; it also means the group has no one automatic interleaving of their records.

The naive model also hides recovery. A consumer may process a record, crash before recording its new position, and later process the same record again. The log can still be perfectly ordered. Repeated work is a question of consumer progress and side effects, not evidence that append order failed.

The missing model is a log as a set of durable, partition-scoped histories plus a separately managed record of each consumer's progress.

The Better Model: Positions, Keys, and Consumer Progress

Plain meaning: A log assigns each record a durable position in one history. A key chooses the history; a consumer position says how far one reader has progressed through it.

In this scenario: order_id=order-7 routes the workflow to partition 1. Fulfillment reads offsets 420, 421, and 422 in that order, recording progress so it can resume after a failure.

Technical names: The ordered history is a partition. Its record position is an offset. A set of cooperating readers is a consumer group, and its stored position is a committed offset or equivalent checkpoint.

producer
  -> choose partition from the key
  -> append a record at the next partition offset

consumer group member
  -> fetch records from its assigned partition in offset order
  -> process them
  -> record how far it has safely progressed

The phrase “safely progressed” carries a design choice. If processing a record triggers an external action, the consumer must decide how its effect and its checkpoint relate. A log offset is evidence about reading a record; it is not evidence that an email, payment, or database transaction outside the log occurred exactly once.

A Worked Trace: One Order, a Retry, and a Restart

Assume the fulfillment consumer owns partition 1. Its next stored position is offset 420.

Step 1: append the order's local history

The order service appends the three events with the same key:

P1 / 420  OrderPlaced(order-7)
P1 / 421  InventoryReserved(order-7)
P1 / 422  PaymentCaptured(order-7)

The producer's configured acknowledgement and retry behavior determine when it receives a successful response and how it handles an uncertain response. A timeout is not proof that the broker rejected the append. Retrying carelessly can result in a duplicate logical event, even though each physical record still has a well-defined partition offset.

For this lesson, assume one copy of each record was appended. The consumer reads them in the order shown.

Step 2: process and checkpoint

Fulfillment processes offsets 420 and 421. It checks that the order exists and that inventory is reserved. It then records its next position as 422—meaning it has safely finished records before 422 and should resume at 422 after a restart.

processed:          420, 421
next position:      422

This checkpoint is local to the fulfillment consumer group. Analytics may be at offset 420; fraud detection may have replayed from 0. Many consumers can derive independent views from the same partition because a log is a reusable history, not a one-use handoff.

Step 3: an external effect exposes the boundary

At offset 422, fulfillment calls an external warehouse API to create a packing task. The API accepts the request, but the consumer crashes before it commits its next position, 423.

After reassignment or restart, a consumer receives offset 422 again. This is the safe default for avoiding a lost record: the checkpoint did not claim progress that was never durably recorded. But the warehouse might now see the same request twice.

warehouse task created
consumer crashes
offset 423 was not committed
-> offset 422 is replayed

The append order is intact. The application still needs an idempotency key, a deduplication record, or an atomic boundary appropriate to the warehouse call. Those techniques are the next lesson's focus. Here, the lesson is that replay changes delivery work without changing the partition's order.

Step 4: order ends at the partition boundary

Meanwhile, another consumer processes partition 0 quickly and sees PaymentCaptured(order-8) before fulfillment reaches PaymentCaptured(order-7). That scheduling fact does not violate anything. The topic promised each partition's order, not a global completion sequence across the topic.

If a business invariant really spans two orders or two keys—for example, “only the first 100 orders may claim a limited promotion”—keying by order_id is not enough. The service needs an explicit shared authority or coordination boundary for the promotion. A consumer cannot reconstruct a global order from two local offsets.

So far: a partition gives a durable sequence and a recoverable reading position. The key decides which events receive that sequence together. Cross-partition relationships, duplicate effects, and global resource limits remain application or coordination work.

What This Changes in a Design Review

Before this model, a team may say, “All order events go to Kafka, so consumers see the workflow in order.” After it, a precise statement is possible:

All events for one order_id share one partition.
Within that partition, consumers process records in offset order.
Consumers may replay records after failure.
No global order is promised across different order_ids or partitions.

That statement creates better questions:

Choosing order_id is a situated preference when each order can be handled independently. If stock allocation or a promotion quota spans many orders, the boundary must change. “Use a globally ordered topic” is one possible answer, but it is a costly coordination decision, not a default cure.

Trade-offs, Limits, and Signals

Partitioned logs make durable replay and local ordering scalable. The trade-off is that one global story becomes many local stories, and the application must name the relations that cross them.

They help when an entity key captures the unit of sequential work. They cost parallelism when one key becomes hot, and they cost reasoning when a workflow spreads across keys or topics. A log does not protect an external side effect from repetition, does not turn offsets into wall-clock time, and does not create atomic transactions across arbitrary consumers and services.

Useful signals include consumer lag per partition, unequal traffic that makes one partition hot, a rising rate of duplicate-effect suppression, increasing retries from producers, and repeated rebalances that delay consumers. These signals distinguish a throughput problem, a progress problem, and a semantics problem. They do not mean the log has “lost order.”

Retention is another boundary. Replay is only available while the required records remain retained and accessible. A consumer reset is powerful precisely because it can rerun history; it should therefore be an intentional operational action, not a substitute for understanding the consumer's checkpoint and side effects.

Common Confusions

Confusion: “A topic is a globally ordered queue.”

Why it is tempting: A topic has one name, and each partition is ordered.

Better model: A partitioned topic contains several ordered logs. State the partition key and the order scope before relying on it.

Confusion: “A duplicate means the log delivered records out of order.”

Why it is tempting: The same event appears twice in a consumer's work.

Better model: A retry or replay can repeat a record while preserving its position. Ordering and exactly-once effects are separate promises.

Confusion: “A larger offset happened later in the real world.”

Why it is tempting: Offsets look like increasing timestamps.

Better model: An offset orders records only in its own partition. It is not a cross-partition clock; the next lesson introduces tools for reasoning about those relationships.

Check Your Understanding

Check: OrderPlaced(order-7) is at partition 1 offset 420, while PaymentCaptured(order-8) is at partition 0 offset 89. Which happened first according to the log?

Think first, then reveal.

Answer: The log does not define that comparison. The offsets belong to different partitions. You need a separate contract or evidence if the application needs a cross-order relation.

Check: A consumer's next stored position is 422. Which records has it safely completed in this example?

Think first, then reveal.

Answer: It has completed offsets before 422: 420 and 421. On restart it should resume at 422, so 422 must still be safe to replay.

Practice: Review a Partitioning Plan

A team plans to key OrderPlaced, InventoryReserved, and PaymentCaptured by order_id, but it will key a shared promotion counter by campaign_id. It claims that the event log will therefore serialize every order's payment and ensure no more than 100 claims of the campaign.

Review the claim. Name what the chosen keys do guarantee, what they do not guarantee about the campaign limit, and one signal you would monitor after launch.

Model answer: Keying the order events by order_id gives each order a local partition order, so one consumer can process that order's workflow sequentially. The shared campaign counter is a different key and may be in a different partition; the log does not create one atomic, globally ordered limit across all orders. The campaign needs its own shared authority or coordination mechanism. Monitor hot-partition traffic for the campaign key, consumer lag, duplicate suppression, or producer retries; each can reveal pressure on the chosen boundary.

Connections

The previous lesson located authority in primary-backup, multi-leader, and leaderless models. A partitioned log makes one of those boundaries tangible: each partition gives a local append history while partitioning spreads the workload across several histories.

The next lesson adds logical, vector, and hybrid clocks. They help describe causal and concurrent relationships once a workflow no longer fits inside one partition's simple offset order.

Resources

Key Takeaways

PREVIOUS Replication Models: Primary-Backup, Multi-Leader, and Leaderless NEXT Logical, Vector, and Hybrid Logical Clocks