A message broker is easy to install and hard to use well. Teams adopt Kafka, Pulsar or a managed log, publish a few topics, and a year later find that every consumer reaches back into the producer's database, that one poisoned message stalls a partition for hours, and that nobody can rebuild a read model because the topic only keeps seven days. None of that is a broker bug. It is the result of choosing patterns by accident.

This article is a catalogue of the event streaming patterns that recur in real systems, explained from the log upwards: what each one guarantees, what it costs, the code that makes it correct, and the way it fails. It finishes with a worked order pipeline, a trade-off table and a checklist. Examples use Kafka vocabulary because it is the most common, but every pattern applies to any partitioned, replayable log.

The patterns in one picture

ProducersOrder serviceDB + outbox tableCustomer serviceCDC connectorPayment servicelarge receiptsTopics (partitioned logs)orders.eventskeyed by order_idcustomers.statecompacted, keyedpayments.eventsclaim-check refsorders.retry / .dlqparked failuresObject storereceipt blobsrelayCDCConsumer groupsNotificationsfan-out group 1Enrichment jobstream-table joinOrder read modelidempotent upsertAnalytics sinklake, replayabletableon error
Producers publish through an outbox relay or CDC; each concern reads through its own consumer group; compacted topics act as tables; failures are parked in retry and dead-letter topics; large payloads live in an object store.

What the log guarantees

Start with what the log actually provides. A topic is split into partitions; each partition is an append-only sequence of records with monotonically increasing offsets. The producer chooses the partition, usually by hashing the record key, so all records with the same key land in one partition in the order they were written. A consumer group divides partitions among its members, and each member records the offset it has processed. Two different groups read the same records independently.

Three consequences drive every pattern below. Ordering exists only per partition, so anything that must be ordered must share a key. Delivery is at least once unless you do extra work, because a consumer can process a record and crash before committing its offset. Retention is a policy, not a property: records are deleted by age or size, or compacted down to the latest value per key. Patterns are mostly ways of living with those three facts.

Notification or state transfer

The first decision is what an event carries. An event notification says that something happened and carries an identifier: order 81723 was paid. Consumers that need detail call the owning service back. It keeps payloads small and the owner in control of its data, but every consumer now couples to the owner's API at runtime, and a burst of events becomes a burst of synchronous reads at the worst moment.

Event-carried state transfer puts the relevant state in the event: order id, customer id, line items, totals, status and version. Consumers build their own local copy and never call back. That buys availability and replayability, at the price of larger events, a schema that is now a public contract, and data that is duplicated in every consumer. The usual answer is a mix: carry the state that consumers routinely need, keep sensitive or bulky fields behind the owner, and never publish a field you would be unwilling to support for years.

Fan-out with consumer groups

Fan-out is the pattern logs are best at. Each independent concern gets its own consumer group: notifications, search indexing, fraud scoring and analytics all read orders.events at their own pace. A slow analytics job does not slow notifications, and a new consumer can start from the earliest retained offset to backfill. Contrast this with a work queue, where a message is consumed once and gone.

Within a group, parallelism is capped by partition count: twelve partitions means at most twelve active consumers, and a thirteenth sits idle. Choose partition count for the peak consumer parallelism you expect over the topic's life, because increasing it later remaps keys to new partitions and breaks per-key ordering for records in flight. Very uneven keys, such as one seller producing a third of all orders, create hot partitions no consumer count can fix; key by order id instead if per-seller ordering is not required.

Publishing safely: outbox and CDC

The hardest bug in event-driven systems is the dual write: a service updates its database and then publishes an event, and a crash between the two leaves them disagreeing forever. The fix is to make one of them the source of truth. In the transactional outbox, the service inserts the event into an outbox table in the same database transaction as the business change, and a relay publishes outbox rows to the log. With change data capture, a connector such as Debezium reads the database's replication log and turns committed row changes into events.

Both give at-least-once publication, so duplicates are possible and consumers must tolerate them. The outbox lets you shape intentional domain events; raw CDC exposes your table layout as a contract unless you transform it. Configure the producer side so the broker does not add its own duplicates or losses:

# producer settings for an outbox relay
acks=all                     # wait for all in-sync replicas
enable.idempotence=true      # broker de-duplicates producer retries per partition
max.in.flight.requests.per.connection=5   # ordering preserved with idempotence on
# topic settings
replication.factor=3
min.insync.replicas=2        # with acks=all, a write needs 2 live copies

Compacted topics as tables, and stream-table joins

A topic with cleanup.policy=compact keeps at least the latest record for each key and deletes older ones in the background; a record with a null value, a tombstone, eventually removes the key. That turns a topic into a replicated table: customers.state keyed by customer id always contains every customer's current row, and a new consumer can rebuild the full table by reading from offset zero.

The pattern this enables is the stream-table join. An enrichment job reads customers.state into a local store, then for each order event looks up the customer and emits an enriched order. In Kafka Streams this is a KStream-KTable join; in Flink it is a lookup or temporal join. Both inputs must be co-partitioned, meaning the same key and the same partition count, or the job must repartition one side first. The join is only as fresh as the table: if the customer update arrives after the order, the order is enriched with the old value. Decide whether that is acceptable or whether you need event-time temporal joins and a wait window.

Idempotent projections and read models

A read model, or materialized view, is a consumer that folds events into a query-friendly store: an order-status table, a search index, a per-customer summary. Because delivery is at least once, the fold must be idempotent. The simplest robust approach is to carry a version or source offset with every event and make the write conditional on it, so a replayed or duplicated event is a no-op. Commit the consumer offset only after the write succeeds.

from confluent_kafka import Consumer
import json, psycopg

consumer = Consumer({
    "bootstrap.servers": "kafka:9092",
    "group.id": "order-read-model",
    "enable.auto.commit": False,
    "auto.offset.reset": "earliest",
    "isolation.level": "read_committed",
})
consumer.subscribe(["orders.events"])
db = psycopg.connect("dbname=readmodel")

UPSERT = '''
INSERT INTO order_view (order_id, status, total_cents, version)
VALUES (%(order_id)s, %(status)s, %(total_cents)s, %(version)s)
ON CONFLICT (order_id) DO UPDATE
SET status = EXCLUDED.status, total_cents = EXCLUDED.total_cents,
    version = EXCLUDED.version
WHERE order_view.version < EXCLUDED.version
'''

while True:
    msg = consumer.poll(1.0)
    if msg is None:
        continue
    if msg.error():
        raise RuntimeError(msg.error())
    event = json.loads(msg.value())
    with db.transaction():
        db.execute(UPSERT, event)          # stale or duplicate versions do nothing
    consumer.commit(message=msg, asynchronous=False)

The WHERE version < guard does double duty: duplicates do nothing, and an out-of-order older event cannot overwrite a newer state. Rebuilding the view becomes a routine operation: create a new table, start a new consumer group from the earliest offset, and switch readers when it catches up. That only works if the topic retains enough history, which is why topics that feed read models are often compacted rather than time-limited.

Retry topics and dead letters

A consumer that cannot process a record has three bad defaults: retry forever and block the partition, skip it and lose data silently, or crash and restart into the same record. The retry and dead-letter pattern makes the choice explicit. Classify failures first. Transient errors, such as a timeout calling a dependency, are retried in place with backoff for a bounded time. Records that keep failing are copied, with error metadata in headers, to a retry topic consumed after a delay, and after a fixed number of attempts to a dead-letter topic that a human or a repair job owns. Permanent errors, such as a payload that fails schema validation, go straight to the dead-letter topic.

Moving a record to another topic gives up per-key ordering: later events for the same order can overtake the parked one. If ordering matters, park the key, not just the record, by recording blocked keys and diverting their later events too until the first is resolved. Always alert on dead-letter growth; a dead-letter topic nobody reads is the skip-silently option with extra storage.

Claim check for large payloads

Logs are tuned for many small records. Brokers cap record size, commonly around 1 MB by default, and large records hurt replication and consumer memory. The claim check pattern stores the bulky payload, such as a PDF receipt or an image, in an object store and publishes an event that carries a reference and a content hash. Consumers fetch the blob only if they need it. Keep blobs at least as long as the topic's retention, and make references immutable, so a replay reads the same bytes.

Schemas are contracts

Once several teams consume a topic, its schema is an API. Register schemas in a schema registry, serialize with Avro, Protobuf or JSON Schema, and enforce a compatibility mode on every change. Backward compatibility, the common default, means new consumers can read old data: you may add optional fields with defaults and remove fields, but not rename or change types. Put an event type and a schema version in every record so consumers can branch safely, and treat a breaking change as a new topic with a migration period in which producers write both.

Worked example: sizing an order pipeline

Take an order service that peaks at 3,000 order events per second, with an average record of 1.5 KB after compression and a requirement to rebuild any read model within an hour. Peak ingress is 3,000 x 1.5 KB, about 4.5 MB/s, trivial for a single broker, so throughput is not what sizes the topic. Consumer parallelism is. The notification consumer spends about 20 ms per event calling an email provider, so one consumer handles 50 events per second and peak needs 60 consumers. That sets a floor of 60 partitions; choosing 64 leaves headroom.

Retention follows from replay needs. Keeping 30 days of history at a daily average of 800 events per second is 800 x 86,400 x 30 x 1.5 KB, about 3.1 TB before replication, 9.3 TB with three replicas, which argues for tiered storage or for compacting a separate orders.state topic that holds only the latest version of each order. Rebuilding the order read model from that compacted topic, perhaps 40 million keys at 1.5 KB, means reading 60 GB; at 100 MB/s across the consumer group that is ten minutes, comfortably inside the one-hour target.

Failure modes

  • Consumer lag spiral. Lag grows, retention deletes unread records, and the consumer silently skips data when its offset falls off the log. Alert on lag in time, not records, and set retention well above the worst recovery time.
  • Rebalance storms. A consumer that takes longer than the poll interval limit to process a batch is kicked out of the group, triggering a rebalance that pauses everyone. Bound batch processing time and use cooperative rebalancing where available.
  • Poison pill. One malformed record crashes every consumer that reaches it. Validate at the edge and route parse failures to the dead-letter topic.
  • Lost ordering after scaling. Adding partitions remaps keys. Plan partition counts up front, or migrate to a new topic with a cut-over.
  • Leaked internals. Raw CDC events expose column names, and the first table refactor breaks five teams. Publish domain events, not tables.
  • Unbounded duplicates. A non-idempotent consumer double-charges or double-emails after a crash. Every side effect needs a dedup key.

Trade-offs

PatternBuys youCosts you
Event notificationSmall events, owner keeps dataRuntime coupling, callback storms
State transferAutonomy, replay, availabilityLarger contracts, duplicated data
Outbox or CDCNo dual-write divergenceRelay to operate, at-least-once
Compacted tableRebuildable current stateNo history, tombstone discipline
Stream-table joinLocal enrichment, no lookupsCo-partitioning, staleness window
Retry and DLQPartitions keep flowingOrdering lost unless keys are parked
Claim checkSmall records, big payloadsSecond store, lifetime coupling

Related reading: the difference between a log and an event store is in Event Sourcing vs Event Streaming; transactions end to end are in Exactly-Once Semantics in Streaming; compaction internals are in Kafka log compaction architecture; contracts are covered in schema registry patterns; the relay itself is in the outbox pattern; and parked messages in dead-letter queue design.

What to do next

  1. List every topic you own with its key, partition count, retention, cleanup policy and consumer groups; mark any topic where ordering matters but the key does not guarantee it.
  2. Find every dual write, a database commit followed by a publish, and replace it with an outbox or CDC.
  3. For each consumer, name its dedup key and make the side effect conditional on it; test by replaying a day of records into a staging copy.
  4. Decide per topic whether events carry state or only identifiers, and write the reason in the schema's documentation.
  5. Add retry and dead-letter topics with a bounded attempt count, error headers and an alert on dead-letter growth.
  6. Alert on consumer lag in seconds and confirm retention exceeds your worst recovery time by a wide margin.
  7. Rebuild one read model from scratch in staging and time it; if it misses your target, add a compacted state topic.
  8. Turn on schema compatibility checks in CI so a breaking change fails a build rather than a consumer.
Key takeaway: Event streaming patterns are ways of living with three facts about a log: ordering only within a partition, at-least-once delivery, and retention as policy. Publish through an outbox or CDC, carry state when consumers need autonomy, give every concern its own consumer group, make every projection idempotent with a version guard, use compacted topics as rebuildable tables, park failures in retry and dead-letter topics, and treat schemas as public contracts.