The problem
Our online shop is split into services.
order records the orders,
stock reserves the goods,
shipping prepares the parcels and
billing invoices. They exchange events over
RabbitMQ, and Vue fronts display the state of the orders.
Architecture decisions rarely arrive together in such a system. They are made one at a time, incident after incident:
- two services change the status of an order, and nothing says which one is authoritative;
-
a
reserve-stockmessage is broadcast to everyone, and two services execute it; - a shared queue splits the messages instead of broadcasting them;
- every front polls every API at regular intervals;
- a common package piles up Doctrine entities from several services;
- an event is lost when a service restarts, another one arrives twice.
Each symptom gets its fix, sound when taken alone. Put end to end, the fixes clash: a copy with no owner, a format with no version, a retry rule per use case. What is missing is not one more technical choice, but a grid that ties the choices together.
The options
Decide as you go, service by service
- What it solves: each service moves at its own pace, without waiting for a common frame.
- What it costs: the conventions drift from one service to the next. Plugging in a new consumer then means reading the code of every producer.
A common library that hides the messaging
- What it solves: a shared library publishes, consumes, retries and logs on behalf of the services. The conventions apply by default.
- What it costs: every service depends on this library and on its version. Changing it forces the services to be redeployed together, which is what the split was meant to avoid.
Six axes decided once, each with its signal
- What it solves: each question has a default answer, written before the first service. The services stay free in their code, not in their contracts.
- What it costs: a document to keep up to date, and rules to check in CI and in review.
Our recommendation
Decide six axes for the whole system, in this order, each with a default choice and the signal that calls it into question.
The order matters: each axis builds on the previous one. We do not name an event before knowing which service owns the data. We do not choose the queues before knowing which messages flow.
1. Functional: one owner per piece of data
Each piece of data has a single owning service. Only that service writes it, and only that service publishes its changes:
-
orderowns the orders, their lines and the customers (Order,OrderLine,Customer); -
stockowns the available quantities and the reservations (Reservation); -
shippingowns the shipments (Shipment); billingowns the invoices.
A service that needs another one's data asks for it over
HTTP, or keeps a copy fed by the events.
shipping thus keeps a copy of each order's
status, without ever changing it. If ordering matters, a
version number prevents a late event from overwriting a more
recent state.
2. Message: event, command, query
Three kinds of messages flow, and each has its channel:
- an event is a past fact, broadcast through the event exchange to whoever wants to hear it;
- a command asks one specific service for an action, point to point;
- a query expects an immediate answer: it stays synchronous, over HTTP.
Events follow the
<context>.<entity>.<past-tense
fact>.v<N>
convention: order.order.placed.v1,
stock.reservation.confirmed.v1. The context is
the service that owns the fact, which ties this axis to the
previous one. The reserve-stock message from
the start was a command in disguise: broadcast, it could be
executed by two services, or by none.
Each event travels in a common JSON envelope: name, date,
producer, schema revision, correlation identifier and
payload. Its event_id, a UUID v7
set by the producer, lets the consumers discard duplicates.
3. Transport: one exchange, one queue per service
A single topic exchange, app.events, carries
every event. Each consuming service owns a queue named after
it (stock.events, shipping.events,
billing.events) and chooses its binding keys.
The producer publishes once, and RabbitMQ delivers a copy to
each interested queue.
The binding keys stay precise: never #, and the
* wildcard only for a family of events handled
in full. A queue is only split on a measured signal, and the
split queues keep the service's prefix. The details are in
our article on the RabbitMQ topology.
4. Push: to the fronts, over SSE
The fronts do not read RabbitMQ. A Mercure hub pushes the changes to them over Server-Sent Events (SSE), after the commit of the service that publishes them. Each browser tab opens a single connection, whatever the number of services.
Each service publishes under its own topic prefixes, with a token restricted to these prefixes. Personal data goes out as private updates. A single service sets the subscriber cookie, on the parent domain it shares with the hub and the fronts.
The guarantees are not those of the broker. A lost update only delays the display: the front reads the API again on every connection, because Mercure notifies and the API is authoritative.
5. Code: share the contracts, nothing else
A monorepo makes sharing code easy, and that is a trap. We
only share a lightweight package,
app-contracts, with no runtime dependency. It
holds the exchange name, the routing keys, one
readonly DTO per event, the JSON Schemas and
frozen fixtures.
Everything else lives in each service: Doctrine entities, outbox, queues, binding keys, handlers, Messenger serializer. A shared entity would couple two databases, a shared handler two deployments.
Messages evolve by adding optional fields, and any other
change creates a v2, published next to the
v1. The producer validates each message against
its schema, and the consumer ignores the fields it does not
know.
Our article on message contracts in a monorepo
details these rules.
6. Transactional: outbox, idempotency, boundaries
A transaction only covers one database, that of a single service. It includes neither the broker, nor the Mercure hub, nor another service. Two mechanisms bridge the gap, one on each side of the broker.
On the producer side, the outbox: the event is written to a table, in the same transaction as the business change. A relay publishes each row, waits for the broker's confirmation, marks it as published, then commits. At worst a duplicate, never a loss: our article on the Outbox pattern details this relay.
On the consumer side, idempotency: a middleware placed after
doctrine_transaction inserts the
event_id into a
processed_event table, in the effect's
transaction. A duplicate inserts nothing, and no handler is
called. An effect outside the database, such as an email,
needs its own idempotency key or a message sent after the
commit.
That leaves the message that fails on every attempt. Messenger retries it with a growing delay, then moves it to a failure transport, where it stays readable and replayable. RabbitMQ's DLX remains for consumers written in another language.
The services map
The table crosses the first axes on our running example. The
topic prefixes are relative to
https://example.com.
| Service | Emits | Consumes | To the fronts |
|---|---|---|---|
order |
order.order.,
order.order.
|
no event | /orders/*, private |
stock |
stock.reservation.
|
order.order.,
order.order.
|
/stock/*, public |
shipping |
shipping.shipment.
|
order.order.*.v1,
stock.reservation.
|
/shipments/*, private |
billing |
no event |
order.order.*.v1,
shipping.shipment.
|
no update |
A row tells what a service declares: its queue, its binding
keys, its topics. The Consumes column tells who will be
affected by a contract change. billing consumes
without emitting anything: its invoices are read through its
API.
The trade-offs we accept
More parts to operate. An outbox and its
relay per producer, a processed_event table per
consumer, a Mercure hub, a contracts package. Each needs
monitoring and, for the tables, a purge.
Eventual consistency. Between two services, a piece of data converges with a delay: the relay's, then the consumer's. A screen that must show the exact state reads the owner's API again.
Conventions to enforce. The naming and the
optional additions only hold if they are checked. We hand
them to CI, through the shared constants and the frozen
fixtures. Code review catches the rest, such as a
v2 published without its v1.
Horizons to align. The outbox retention,
the processed_event purge and the stay in the
failure transport depend on one another. Purging
processed_event too early lets a replay process
already applied events again.
The signals that should make us review an axis:
- replaying becomes routine rather than an exceptional move: a log such as Kafka becomes relevant;
- the relay, measured, falls too far behind or loads the database: change data capture (CDC) replaces the polling. If the volume is the cause, Kafka becomes relevant;
- the screens assemble several services, or the permissions can no longer be expressed as topic prefixes: a BFF comes back;
- a consumer is written in another language: the schemas are used to generate its code, and a DLX handles its quarantine;
- heavy messages delay critical ones in the same queue: that queue gets split.
Implementation
What remains are routine moves, which we write as procedures, versioned with the code.
Changing a contract
-
Change together, in
app-contracts, the JSON Schema, the DTO and, for a new event, its routing constant. -
For an optional addition, bump the revision (
1.0to1.1) and add a fixture, without touching the older ones. -
For any other change, create the
v2event and publish it next to thev1, as two outbox rows in the same transaction. - Let CI validate each frozen fixture against its schema and its DTO: a mismatch blocks the merge.
-
Bind each consumer to the
v2, handler included, then delete thev1binding key in RabbitMQ: Messenger never removes one. Meanwhile, each fact arrives twice, under twoevent_ids: the handler discards the duplicate by business key. -
Retire the
v1once no queue receives it and nov1message can be replayed.
Adding a consumer
Take billing, which arrives after the other
services.
-
Declare the
billing.eventsqueue and its binding keys, without#, in the configuration ofbilling. The producer does not change. - Write a handler for each event these binding keys let through, wildcards included.
-
Set up the receiving side: envelope serializer,
processed_eventmiddleware afterdoctrine_transaction, retries and failure transport. -
Deploy, then run
messenger:setup-transports: the queue receives the flow from that moment on. -
If the history matters, replay past events from each
producer's outbox (
order,shipping), with no guaranteed order between them.processed_eventabsorbs the overlap with the normal flow. - Update the subscription catalogue, or generate it again from the services' configuration.
Replaying events
A replay republishes rows from a producer's outbox to the queue of a single consumer. Our article "Kafka or RabbitMQ?" explains why this replay is enough while it stays occasional.
- Identify the producer, the routing key, the period and the target queue.
- Check that the target queue is bound to this key: the replay goes through the default exchange, which ignores bindings.
-
Check that the period stays within the outbox horizon and
within the consumer's
processed_eventhorizon. -
Count the rows without publishing anything (
--dry-runoption), then republish. -
After a bug that produced a wrong effect without an error,
fix the data first. Then remove the matching rows from
processed_event, otherwise the consumer discards the replayed event.
Handling the quarantine
-
List the messages with
messenger:failed:show, and read each exception with the-vvoption. - Fix the cause: missing data, a handler bug, an unavailable dependency.
-
Process again with
messenger:failed:retry. If another copy was applied in the meantime,processed_eventdiscards it. -
Discard with
messenger:failed:removea message that must not be processed. -
For a consumer written in another language, read the queue
fed by the DLX and its
x-deathheader.
Checklist
- One owning service for each piece of data; elsewhere, read-only copies.
-
Past-tense events, named
<context>.<entity>.<past-tense fact>.v<N>. - Commands point to point and queries over HTTP, outside the event exchange.
-
A single topic exchange, one queue per consuming service,
no
#binding key. - Towards the fronts, each service's own topic prefixes and a token restricted to them.
- A contracts package with no runtime dependency; entities, queues and handlers outside it.
-
An outbox in each producer, a
processed_eventtable in each consumer. - A single quarantine level per queue.
-
Aligned horizons: outbox retention,
processed_eventpurge, stay in quarantine. - Written procedures for contracts, consumers, replays and the quarantine.
- For each axis, the signal that should make you review it.
Sources
- martinfowler.com, What do you mean by “Event-Driven”?: event notification and event-carried state transfer.
- Enterprise Integration Patterns, Event Message, Command Message and Publish-Subscribe Channel.
- microservices.io, Database per service, Transactional outbox and Idempotent Consumer.
- RabbitMQ, Topics tutorial, Exchanges (default exchange) and Dead Letter Exchanges.
- Symfony, Messenger: AMQP transport, Doctrine middleware, retries and failure transport.
- JSON Schema, 2020-12 specification.
- Mercure, protocol specification and documentation.
- MDN, Server-sent events.