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-stock message 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.

Six separate decisions, all visible on the flow of one order 5 app-contracts : exchange name, routing keys, DTOs, JSON Schemas, fixtures the only shared code; outbox, queues, handlers and serializers stay in each service order 1 owns the orders 2 publishes facts: order.order.placed.v1 6 outbox table written with the order 3 app.events topic exchange billing.events read by billing owns the invoices shipping.events read by shipping owns the shipments stock.events read by stock owns the stock and the reservations 6 in each consumer: processed_event and failure transport after the commit 4 Mercure hub topics as IRIs SSE shop and admin fronts EventSource , then a fresh read of the API 1 2 3 4 5 6 Functional : one owner per piece of data Message : past facts, versioned Transport : one exchange, one queue per service Push : SSE to the fronts Code : only the contracts are shared Transactional : outbox and idempotency

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:

  • order owns the orders, their lines and the customers (Order, OrderLine, Customer);
  • stock owns the available quantities and the reservations (Reservation);
  • shipping owns the shipments (Shipment);
  • billing owns 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.placed.v1, order.order.cancelled.v1 no event /orders/*, private
stock stock.reservation.confirmed.v1 order.order.placed.v1, order.order.cancelled.v1 /stock/*, public
shipping shipping.shipment.dispatched.v1 order.order.*.v1, stock.reservation.confirmed.v1 /shipments/*, private
billing no event order.order.*.v1, shipping.shipment.dispatched.v1 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

  1. Change together, in app-contracts, the JSON Schema, the DTO and, for a new event, its routing constant.
  2. For an optional addition, bump the revision (1.0 to 1.1) and add a fixture, without touching the older ones.
  3. For any other change, create the v2 event and publish it next to the v1, as two outbox rows in the same transaction.
  4. Let CI validate each frozen fixture against its schema and its DTO: a mismatch blocks the merge.
  5. Bind each consumer to the v2, handler included, then delete the v1 binding key in RabbitMQ: Messenger never removes one. Meanwhile, each fact arrives twice, under two event_ids: the handler discards the duplicate by business key.
  6. Retire the v1 once no queue receives it and no v1 message can be replayed.

Adding a consumer

Take billing, which arrives after the other services.

  1. Declare the billing.events queue and its binding keys, without #, in the configuration of billing. The producer does not change.
  2. Write a handler for each event these binding keys let through, wildcards included.
  3. Set up the receiving side: envelope serializer, processed_event middleware after doctrine_transaction, retries and failure transport.
  4. Deploy, then run messenger:setup-transports: the queue receives the flow from that moment on.
  5. If the history matters, replay past events from each producer's outbox (order, shipping), with no guaranteed order between them. processed_event absorbs the overlap with the normal flow.
  6. 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.

  1. Identify the producer, the routing key, the period and the target queue.
  2. Check that the target queue is bound to this key: the replay goes through the default exchange, which ignores bindings.
  3. Check that the period stays within the outbox horizon and within the consumer's processed_event horizon.
  4. Count the rows without publishing anything (--dry-run option), then republish.
  5. 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

  1. List the messages with messenger:failed:show, and read each exception with the -vv option.
  2. Fix the cause: missing data, a handler bug, an unavailable dependency.
  3. Process again with messenger:failed:retry. If another copy was applied in the meantime, processed_event discards it.
  4. Discard with messenger:failed:remove a message that must not be processed.
  5. For a consumer written in another language, read the queue fed by the DLX and its x-death header.

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_event table in each consumer.
  • A single quarantine level per queue.
  • Aligned horizons: outbox retention, processed_event purge, stay in quarantine.
  • Written procedures for contracts, consumers, replays and the quarantine.
  • For each axis, the signal that should make you review it.

Sources