Designing one event-driven flow is a modelling problem. Running forty services that communicate through events is an architecture problem: who owns each topic, how contracts change without breaking consumers, how a service publishes reliably, how a consumer survives duplicates and poison messages, how one service answers queries about another service's data, and how anyone traces a request that hops through five queues. The basics, events versus commands, the styles of event-driven design and ordering, are in Event-Driven Architecture, in depth, and the trade-offs of splitting into services at all are in Microservices Architecture, in depth.
This article is the estate-level view. It describes a reference architecture, then works through the rules that keep it healthy as it grows: topic ownership, envelopes and schema compatibility, publishing through an outbox, consuming idempotently, retry and dead-letter topology, local read models, tracing and capacity. A worked example adds a new service to a running estate without touching any producer, and the article ends with failure modes and a checklist.
The reference architecture
Figure 1 shows the shape most mature estates converge on. Each service owns a database and writes outgoing events to an outbox table in the same transaction as its state change. A relay, either a poller or change data capture, moves outbox rows to the broker. Topics are owned by exactly one producing service. Every event schema is registered and checked for compatibility before a producer can deploy. Consumers keep an inbox table of processed message IDs, and some build local read models of data owned elsewhere. An event catalog lists every topic, its owner, schema and consumers.
None of these parts is exotic. What makes the estate work is that every service follows the same rules, so a new engineer can predict how any service publishes and consumes without reading its code.
Topic ownership, naming and keys
One owner per topic. Only the service that owns an entity publishes events about it. If two services publish to orders.order.v1, nobody can change its schema safely and ordering per order is lost. Other services that want to influence an order send it a command.
Names carry ownership and version. A convention such as <domain>.<entity>.v<major> makes the owner obvious and leaves room for a breaking change as a new topic. Put the event type in the envelope, not in the topic name, so that one entity's events share a partition order.
Partition key is the entity ID. Brokers such as Kafka guarantee order only within a partition, so all events for one order must share a key. Choosing partition counts and the cost of changing them later are covered in Kafka Partition Architecture.
| Decision | Recommended default | Why |
|---|---|---|
| Event granularity | Business facts (OrderPlaced), not row changes | Consumers depend on meaning, not your table layout |
| Payload | Enough state to act without a call back | Callbacks recreate synchronous coupling and load |
| Retention | Days for notifications; compacted or long for entity state topics | New consumers need history to build read models |
| Private vs public events | Internal topics stay inside the owning team | Only public topics get compatibility guarantees |
Envelopes and schema compatibility
A shared envelope lets every consumer handle routing, deduplication and tracing the same way. The CloudEvents specification defines widely used attribute names, and adopting them saves a debate:
{
"specversion": "1.0",
"id": "01J9Z6K3W7Q8F2T4M5N6P7R8S9", # unique per event; the dedup key
"source": "orders-service",
"type": "OrderPlaced",
"subject": "order/8812", # also the partition key
"time": "2026-10-03T19:22:04Z",
"dataschema": "orders.order.v1/OrderPlaced/7",
"data": { "orderId": "8812", "customerId": "c-301", "totalMinor": 4599,
"currency": "EUR", "lines": [{"sku": "B-17", "qty": 1}] }
}Schemas change constantly, so a registry enforces compatibility at build time. The common modes are backward (new consumers can read old events), forward (old consumers can read new events) and full (both). For public topics, choose full or at least forward compatibility, because you cannot redeploy every consumer at the moment a producer ships. In practice that means: add fields only as optional with defaults, never rename or retype a field, never reuse a removed field's name, and put a breaking change on a new major topic that runs in parallel until consumers move.
Publishing reliably with an outbox
The dual-write problem is the most common data-loss bug in event-driven systems: a service commits to its database and then publishes, and a crash between the two leaves the world inconsistent. The fix is the outbox, explained fully in Transactional outbox architecture. The minimum schema and the relay loop look like this:
CREATE TABLE outbox (
id uuid PRIMARY KEY,
topic text NOT NULL,
key text NOT NULL,
payload jsonb NOT NULL,
headers jsonb NOT NULL, -- traceparent, correlation id
created_at timestamptz NOT NULL DEFAULT now(),
published_at timestamptz
);
CREATE INDEX outbox_unpublished ON outbox (created_at) WHERE published_at IS NULL;
-- relay, run by one leader per database
WITH batch AS (
SELECT id FROM outbox WHERE published_at IS NULL
ORDER BY created_at LIMIT 500 FOR UPDATE SKIP LOCKED)
SELECT o.* FROM outbox o JOIN batch USING (id) ORDER BY o.created_at;
-- publish each row with key = o.key, wait for broker acks,
-- then UPDATE outbox SET published_at = now() WHERE id = ANY(:ids)A crash after the broker acknowledges but before the update republishes the batch, so the relay gives at-least-once delivery and consumers must deduplicate. Change data capture tools that read the database log remove the polling load and preserve commit order, at the cost of another moving part to operate. Delete published rows on a schedule so the table stays small.
Consuming reliably: inbox, retries and dead letters
Consumers face three realities: duplicates, out-of-order arrival across keys, and messages they cannot process. The inbox pattern handles duplicates by recording the event ID in the same transaction as the effect:
def handle(msg):
evt = parse(msg)
with db.transaction() as tx:
inserted = tx.execute(
"INSERT INTO inbox(event_id, received_at) VALUES (%s, now()) "
"ON CONFLICT DO NOTHING", [evt.id]).rowcount
if inserted == 0:
return # duplicate: effect already applied
apply_effect(tx, evt) # business change + any outbox rows
consumer.commit(msg) # offset after the DB commitFailures are split by type. A transient failure, such as a timeout calling a dependency, deserves retries with backoff. A permanent failure, such as a payload that fails validation, will fail forever; retrying it in place blocks every later message on that partition. The usual topology is a main topic, one or more retry topics with increasing delays (for example 1 minute, then 10 minutes), and a dead-letter topic. The handler classifies the error and republishes to the next stage with headers recording the attempt count and the original error.
Be careful with ordering when you divert messages: if event 5 for order 8812 goes to a retry topic, event 6 for the same order may be processed first. For entity-state consumers, either pause the key (park later events for that key until the failed one is resolved) or make the handler order-tolerant by carrying a version number and ignoring events older than the stored version.
Local read models
In a synchronous estate, Shipping asks Customers for an address on every shipment. In an event-driven estate, Shipping keeps a local read model of the customer fields it needs, updated from customers.customer.v1 events. This removes a runtime dependency and a network hop, and it costs staleness and storage. Three rules keep read models honest: copy only the fields you use, store the source event's version or timestamp with each row so you can reason about staleness, and make the projection rebuildable by replaying the topic from the beginning, which requires the topic's retention or compaction to cover the full entity history. Combining this with a separate write model is the CQRS approach covered in CQRS + Event Sourcing Architecture.
Observability for asynchronous flows
Asynchronous systems fail quietly: nothing returns an error, a number just stops moving. Four signals catch most problems.
- Consumer lag per group and partition, in messages and in seconds. Alert on lag in seconds, because a rising message count during a sale may be fine.
- End-to-end latency from the event's
timeto the consumer's commit, as a histogram. - Dead-letter arrivals: any arrival is a bug or a contract break until proven otherwise.
- Outbox age: the oldest unpublished row. A stuck relay looks exactly like a quiet day without it.
Propagate W3C traceparent in message headers and start a span linked to it in the consumer, so one trace shows the HTTP request, the outbox write, the publish and every downstream handler. Keep the event catalog generated from the schema registry and consumer group metadata rather than written by hand, or it will be wrong within a quarter.
Worked example: adding a Loyalty service
An estate publishes orders.order.v1 with 24 partitions at a peak of 2,000 events per second. The business wants a Loyalty service that awards points when an order is delivered and reverses them on refund.
- Loyalty subscribes with its own consumer group. No producer changes and no producer deploys; Orders does not learn that Loyalty exists. This is the main payoff of the architecture.
- Capacity: if one handler instance processes 150 events per second, 2,000 per second needs at least 14 instances, and 24 partitions allow up to 24 active consumers, so there is headroom. Had the topic been created with 8 partitions, Loyalty would have been capped at 8 instances, or 1,200 per second, and would lag at peak.
- Backfill: Loyalty needs historical orders to award points for recent deliveries. It starts from the earliest retained offset, which is 7 days on this topic, and a one-off batch job loads older delivered orders from a data-warehouse export. Seven days of events at an average of about 660 per second is roughly 400 million. At 14 instances replaying at 2,100 per second that takes about 53 hours; scaling to all 24 partitions (3,600 per second) cuts it to about 31 hours. Either way the backfill runs days before launch, not on launch day.
- Idempotency: the inbox table makes the replay safe to restart, and points are keyed by order ID so a redelivered OrderDelivered cannot double-award.
- Ordering: OrderRefunded must not be applied before OrderDelivered. Because both share the order ID key they arrive in order; if one is diverted to retry, the handler parks later events for that order.
Failure modes and when not to use it
| Failure | What you see | Prevention |
|---|---|---|
| Dual write without outbox | Orders exist with no event, or events with no order | Outbox or CDC, never publish-after-commit |
| Breaking schema change | One consumer's DLQ fills after a producer deploy | Registry compatibility check in the producer's CI |
| Poison message on main topic | One partition's lag climbs while others are flat | Classify errors; retry and DLQ topics |
| Hot key | One partition lags; one consumer at 100 percent CPU | Better key choice; split the heavy entity's events |
| Event storm on replay | Downstream services flooded by a reprocessing job | Rate-limit replays; publish replays to a separate topic |
| Hidden synchronous coupling | Consumer calls producer's API per event | Fatter events or a local read model |
| Unowned topic | Nobody can approve a schema change | Catalog entry with owner required to create a topic |
When not to use this architecture: a small team with a handful of services, request-response workflows where the user waits for the answer, and domains that need strong consistency across entities. In those cases a modular monolith or synchronous calls are cheaper to build, test and operate. For a concrete single-domain walk-through, see Designing an Event-Driven Order System.
What to do next
- Publish a one-page estate standard: topic naming, envelope fields, partition key rule, compatibility mode and retention per topic class.
- Add a schema-registry compatibility check to every producer's CI pipeline.
- Find every publish-after-commit in your codebase and move it to an outbox or CDC.
- Add an inbox table to every consumer that has side effects, and commit offsets only after the database commit.
- Create retry and dead-letter topics per consumer group, with error classification in the handler.
- Dashboard consumer lag in seconds, end-to-end latency, DLQ arrivals and outbox age; alert on all four.
- Propagate traceparent and a correlation ID in message headers end to end.
- Generate an event catalog from the registry and consumer groups, and require an owner before a topic can exist.