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.
Define stream partitions, offsets, consumer groups, retention, replay, lag, rebalancing, and backpressure in a vendor-neutral model.
Explain why ordering is usually per partition/key and why increasing partitions trades ordering scope for parallelism.
Diagnose slow consumers, poison events, group ownership changes, and replay behavior using observable offsets.
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.
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
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.
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")
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
- Where does an offset have meaning?
- Why can two consumers in one group process different orders concurrently?
- What is consumer lag?
- Why can committing an offset before the side effect be unsafe?
- 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.
- Apache Kafka downloads — Current official release/download status used only as a version snapshot.
- Apache Kafka 4.2.1 release announcement — Current supported bugfix release announcement.
- Apache Kafka design documentation — Official architecture concepts for partitions, consumer groups, replication, and log retention.
- Apache Kafka consumer configuration — Official configuration surface relevant to offsets, polling, group membership, and backpressure tuning.