Model a partitioned event stream with offsets, consumer-group ownership, scoped ordering, lag, replay, poison events, rebalancing, and backpressure.

Event Streams, Consumer Groups, Replay, Retention, and Backpressure

Publication is only the beginning: retained streams need explicit partitions, offsets, replay, ownership, lag, and backpressure behavior.

Advanced120–155 minutesStream mechanics labPython 3.13+ · standard library / sqlite3 where notedVendor-neutral · free/local mandatory pathLast reviewed: August 2026
01

Define stream partitions, offsets, consumer groups, retention, replay, lag, rebalancing, and backpressure in a vendor-neutral model.

02

Explain why ordering is usually per partition/key and why increasing partitions trades ordering scope for parallelism.

03

Diagnose slow consumers, poison events, group ownership changes, and replay behavior using observable offsets.

04

Select retention and recovery policies that preserve rebuildability without allowing lag to exhaust storage or memory.

1. A stream is a retained ordered log split into parallel lanes

AtlasMart's event stream is not merely a queue of transient messages. In the model used here, a partition is an append-only ordered sequence of records; an offset identifies a record's position within that partition. A consumer group coordinates workers so one group instance owns a partition at a time, allowing parallel processing across partitions while preserving the partition's order for that group.

Retention determines how long records remain replayable independent of whether one consumer processed them. Replay means reading earlier retained offsets again. Lag is the distance between the stream's end and the consumer's committed progress. Backpressure is the system's response when production outpaces safe consumption.

2. Ordering is scoped

If all events for o-1 use the same partition key, OrderCreated can precede OrderPaid for that order in one partition. An event for o-2 in another partition can be processed earlier or later according to independent scheduling. A statement such as “Kafka preserves order” is incomplete without specifying the topic/partition/key and producer/consumer assumptions.

Design consequence

More partitions increase parallelism and failure isolation but reduce the scope over which a simple stream order exists. Cross-key invariants need a separate coordination/modeling strategy.

3. Consumer-group ownership and rebalancing

When a worker joins, leaves, fails, or the partition count changes, group ownership can be reassigned. A safe consumer commits progress only after the application state it represents is durable—or uses a store/transaction protocol that couples the two. Otherwise, committing first risks loss; processing first and checkpointing later risks duplicates.

During a rebalance, observe assignment generation/epoch, current owners, last committed offset, in-flight work, and restart position. Idempotency remains important because group coordination does not erase a crash between external side effect and offset commit.

4. Deliberately wrong approach: retry one poison event forever on the partition thread

A malformed event at partition 1 offset 1 can block every later record in that partition if the consumer retries it synchronously without a bound. Lag rises, retention pressure increases, and unrelated orders sharing that partition stop advancing. A safer policy validates schema, retries transient failures with bounded backoff, and quarantines permanently bad records in a dead-letter workflow that preserves event identity and diagnostic context.

Backpressure needs an explicit strategy: reduce fetch/concurrency, pause partitions, shed nonessential work, expand consumer capacity when safe, or slow upstream production through application-level admission. Unbounded in-memory buffering simply moves the outage to memory exhaustion.

5. AtlasMart lab: offsets, lag, poison record, and rebalance

Mandatory lab environment

Python 3.13+ standard library only. No Kafka broker, Docker, JVM, network, or cloud service is required. The records and partition ownership are deterministic teaching inputs.

python · AtlasMart deterministic simulation
from collections import defaultdict

records = [
    {"partition":0,"offset":0,"key":"o-1","type":"OrderCreated"},
    {"partition":1,"offset":0,"key":"o-2","type":"OrderCreated"},
    {"partition":0,"offset":1,"key":"o-1","type":"OrderPaid"},
    {"partition":1,"offset":1,"key":"o-3","type":"POISON"},
    {"partition":0,"offset":2,"key":"o-4","type":"OrderCreated"},
    {"partition":1,"offset":2,"key":"o-2","type":"OrderPaid"},
]
end_offset={0:2,1:2}
owners={0:"consumer-A",1:"consumer-B"}
committed={0:-1,1:-1}
dlq=[]
processed=[]

print("INITIAL OWNERSHIP:", owners)
for r in records:
    owner=owners[r["partition"]]
    if r["type"] == "POISON":
        dlq.append(r)
        print(owner, "routes poison event to DLQ", (r["partition"],r["offset"]))
        committed[r["partition"]]=r["offset"]
        continue
    processed.append((owner,r["partition"],r["offset"],r["key"],r["type"]))
    committed[r["partition"]]=r["offset"]
print("processed:", processed)
print("committed offsets:", committed, "DLQ:", [(x['partition'],x['offset']) for x in dlq])

print("\nLAG / BACKPRESSURE")
slow_committed={0:2,1:0}
lag={p:end_offset[p]-slow_committed[p] for p in end_offset}
print("slow consumer committed:", slow_committed, "lag:", lag)
print("policy: pause upstream-adjacent work or reduce fetch/concurrency before memory/latency explodes")

print("\nREBALANCE")
owners={0:"consumer-C",1:"consumer-C"}  # B left; C takes both in this toy model
print("new ownership:", owners)
resume={0:committed[0]+1,1:committed[1]+1}
print("resume positions after committed offsets:", resume)

print("\nREPLAY")
replay_from={0:0,1:0}
replayed=[(r['partition'],r['offset'],r['key']) for r in records if r['offset'] >= replay_from[r['partition']]]
print("full retained replay:", replayed)

print("\nORDERING SCOPE")
print("o-1 events ordered inside partition 0; there is no useful total order between partition 0 and 1")
Expected evidence

The poison record is quarantined while progress continues. A deliberately slow committed offset shows partition-specific lag. Ownership can move after a rebalance, and replay starts from explicit offsets rather than an implicit “already seen forever” state.

6. Production judgment

Choose partition keys from the ordering/invariant scope, not only load distribution. Measure records/bytes per second, p50/p95/p99 processing latency, committed offset, end offset, time lag, rebalance count/duration, retry/DLQ rate, per-partition skew, consumer CPU/memory, downstream dependency latency, and retained-storage headroom. Retention must exceed realistic outage/rebuild windows or recovery will require a new source snapshot.

Apache Kafka is an optional implementation reference; the official downloads currently list Kafka 4.2.1 as a supported release, with 4.3.0 shown in archived releases despite being newer by version number. This chapter therefore avoids equating “highest version number” with “current recommended deployment.” The lab remains vendor neutral.

The next lesson uses one retained change stream to build several independently operated projections.

Check your understanding

  1. Where does an offset have meaning?
  2. Why can two consumers in one group process different orders concurrently?
  3. What is consumer lag?
  4. Why can committing an offset before the side effect be unsafe?
  5. What should happen to a permanently malformed event?
Review the answers

1. Within a particular stream partition/log scope; it is not a universal global sequence number.

2. They can own different partitions while each partition retains its own order.

3. The gap between the stream end and the consumer’s durable committed progress, measured in records/bytes/time depending on the system.

4. A crash can make the group resume after a record whose effect never became durable, causing loss.

5. Quarantine/dead-letter it with identity and diagnostics under a bounded retry policy rather than block the partition indefinitely.

References

Foundational claims use primary specifications/research or current official documentation where practical. Product references are optional implementation anchors; the mandatory labs are vendor-neutral.

Keep knowledge open

Help the academy stay free and grow.

If these tutorials save you time, a small donation supports new lessons, technical review, diagrams, examples, and long-term maintenance.

ETHEthereum / ERC-20 only
0x716c4Ab160C4B66F31a28AE2448BfF68fc3a2ef0

Send only Ethereum or ERC-20 compatible assets to this address.