WritingEngineering

The outbox relay is the only process allowed to publish to RabbitMQ

Run a transactional outbox with RabbitMQ so a committed row reaches consumers once, even when the relay dies between the broker's confirm and its own commit.

A capture of GBP 42.50 commits at 09:14:02. The payments row flips to captured, and the ledger service needs to hear about it. The handler commits, then calls channel.publish('payments', 'payment.captured', ...). The RabbitMQ connection dropped a second earlier, the publish throws, the handler logs it, and the HTTP response has already gone back as 200. The ledger never posts the entry, and nobody finds out until month-end reconciliation opens an exception for a payment that was captured and never booked.

The transactional outbox fixes this, but only if you apply it without exceptions. My thesis is that the relay that reads the outbox table must be the only process in the service that ever calls publish. Publishing from the handler after commit "for lower latency" and keeping the outbox as a fallback recreates the dual write with a second publisher, and its duplicates arrive at the worst moment: when the broker has just recovered and both publishers are busy.

Two publishers are the dual write wearing a different coat

Writing the outbox row in the same transaction as the payments row means the database holds one truth: either both rows exist or neither does. The moment a handler publishes directly after commit, that truth has two readers with two clocks. The handler publishes at commit plus a few milliseconds. The relay publishes when it next polls, perhaps 200 milliseconds later, perhaps 40 seconds later if it fell behind during a broker outage. Each sends the same payment.captured for the same payment, and the consumer gets two copies with no shared sequence to reconcile them.

The duplicate rate is not constant, which is why the design looks fine in staging. In a quiet hour the handler's publish lands first and the relay finds the row already sent, so almost nothing duplicates. During a broker reconnect the handler's publishes fail, the relay catches up with a backlog of several thousand rows, and the handler's publishes start succeeding again at the same time. That is the hour when the ledger receives two of everything. I have built services on NestJS with MySQL and RabbitMQ, and the rule I apply is that the relay is the only module with a publish call in it, enforced by a lint rule on imports of the channel.

The relay loop: claim, publish, confirm, mark

The relay is a loop with four steps, and the order of the last two is the whole design. It claims a batch of pending rows with FOR UPDATE SKIP LOCKED, so a second worker takes the next batch rather than blocking, publishes each row on a confirm channel with the outbox row id as the AMQP messageId, waits for the confirm, and only then marks the row sent.

  1. Claim

    Open a transaction and select pending rows ordered by id with FOR UPDATE SKIP LOCKED. Unclaimed rows stay available to other workers.
  2. Publish

    Send each row's payload to the exchange on a confirm channel, with the outbox id as messageId and the aggregate sequence in a header.
  3. Confirm

    Wait for the broker's basic.ack for that message. A nack or a closed channel means the row stays pending and the batch rolls back.
  4. Mark

    Update the row to sent, commit the transaction, and release the locks. Rows are deleted later by a retention job, never here.
TypeScript
// Simplified: one relay tick against MySQL and an amqplib ConfirmChannel
async function relayTick(db: Db, ch: ConfirmChannel): Promise<number> {
  return db.transaction(async (tx) => {
    const rows = await tx.query<OutboxRow>(
      `SELECT id, routing_key, payload, aggregate_id, seq
         FROM outbox
        WHERE status = 'pending'
        ORDER BY id
        LIMIT 100
        FOR UPDATE SKIP LOCKED`,
    );
    for (const row of rows) {
      await new Promise<void>((resolve, reject) =>
        ch.publish('payments', row.routing_key, Buffer.from(row.payload), {
          persistent: true,
          messageId: row.id,
          headers: { aggregateId: row.aggregate_id, seq: row.seq },
        }, (err) => (err ? reject(err) : resolve())),
      );
      await tx.query(
        `UPDATE outbox SET status = 'sent', sent_at = NOW(6) WHERE id = ?`,
        [row.id],
      );
    }
    return rows.length;
  });
}

The confirm is what makes the mark safe. Without confirm mode, publish returns as soon as the bytes are in the socket buffer, and a broker that dies a moment later has never seen the message while your row already says sent. With confirm mode, the broker acknowledges a persistent message to a durable queue only once it has written it to disk, so everything the relay does after that point can be repeated without losing data.

flowchart TD
  H[Capture handler] --> T[One transaction]
  T --> P[payments row captured]
  T --> O[outbox row pending]
  O --> R[Relay]:::accent
  R --> B[RabbitMQ exchange]
  B -- confirm --> R
  R --> S[outbox row sent]
  B --> Q[ledger queue]
  Q --> C[Ledger consumer]
One transaction writes both rows; only the relay talks to the broker

Walking through the crash that duplicates a message

Take outbox row 7f3a, the payment.captured for the GBP 42.50 capture. At 09:14:02.310 a relay worker claims it with 99 other rows. At 09:14:02.344 it publishes 7f3a with messageId: "7f3a", and at 09:14:02.351 the broker confirms: the message is on disk in the ledger queue. The worker runs the UPDATE ... SET status = 'sent' inside its open transaction, moves on to row 7f3b, and at 09:14:02.360 the container is killed by the orchestrator for a deploy. The transaction never commits. MySQL rolls back, 7f3a is pending again, and its lock is released.

The replacement worker starts at 09:14:05, claims 7f3a, publishes it again, and the ledger queue now holds two messages with the same messageId. The consumer receives the first, inserts the ledger entry for GBP 42.50 and a row in processed_messages keyed by (consumer, message_id) in the same transaction, commits, and acks. It receives the second, attempts the same insert, hits the unique constraint, acks without doing anything else, and the ledger shows one entry.

flowchart TD
  A[Relay claims row 7f3a] --> P1[Publish message id 7f3a]
  P1 --> K[Broker confirms to disk]
  K --> X[Relay dies before commit]:::accent
  X --> A2[New relay claims 7f3a again]
  A2 --> P2[Publishes 7f3a again]
  P2 --> Q[Queue holds two copies]
  Q --> D1[Consumer applies first]
  Q --> D2[Consumer drops second on id]
The crash window sits between the broker's confirm and the relay's commit

The relay duplicated because a process died in a window a few milliseconds wide, and the consumer absorbed it because every message carried a stable id. A handler publish produces the same duplicate with no crash at all, every time the relay is slower than the handler, which is most of the time.

The consumer's half of the contract

RabbitMQ does not offer exactly-once delivery, so the relay promises at-least-once, the consumer promises idempotence, and the messageId is the contract between them. The consumer opens a transaction, inserts (consumer_name, message_id) into processed_messages, applies its effect in the same transaction, commits, then acks. If the insert fails on the unique constraint, it acks and returns. If the process dies after commit but before the ack, the broker redelivers, the insert fails, and the consumer acks. If it dies before commit, nothing was written, the broker redelivers, and the effect is applied once.

Ordering is the part of the contract the broker cannot keep for you. A single queue delivers in order to a single consumer, but a prefetch of 50 across three consumers, or one redelivery after a nack, puts payment.refunded ahead of payment.captured for the same payment. So the relay sends the aggregate's own sequence number in a header, and the consumer applies events by that sequence, holding one whose predecessor has not arrived. The table below lists the stopping points I design against; each has a different recovery.

Where the process stopsWhat the broker holdsWhat happens next
Relay dies after claim, before publishNothingLock released, row pending, next tick publishes once
Relay dies after confirm, before commitOne copy on diskRow pending again, second copy published, consumer drops it on id
Broker nacks or closes the channelNothing durablePromise rejects, batch rolls back, relay backs off and retries
Consumer dies after commit, before ackRedelivery pendingRedelivered, dedupe insert fails, consumer acks

Where the single-publisher rule costs you

The first cost is latency. A consumer hears about the capture when the relay next polls, not when the handler commits. A poll every 100 milliseconds against an index on status, id is cheap on MySQL, and nothing a payments consumer does is interactive, so I accept it. If a product needs the ledger entry before the HTTP response returns, post it inside the capture transaction; do not publish from the handler.

The second cost is the table itself. An outbox that is never pruned grows by one row per event for ever, and the pending scan slows as the index fills with sent rows. A retention job that deletes sent rows older than seven days, in batches of a few thousand, keeps the hot index small. Do not delete inside the relay: a row deleted before commit was published and can no longer be proved to have been.

The third cost is cross-aggregate ordering, which has no clean fix. Two relay workers with SKIP LOCKED publish two batches in parallel, so payment.captured for payment A may reach the exchange after payment.refunded for payment B even though A committed first. Within one payment the sequence header keeps order; across payments there is none. A consumer that needs a global order, such as a running balance per merchant, derives it from its own ledger, or you run one relay worker and accept the throughput ceiling.

The fourth cost is the consumer that cannot be made idempotent. Sending an email is the usual case: no transaction spans the processed_messages insert and the SMTP call. Insert the dedupe row and commit before sending, so a crash after the send is a dropped redelivery and a crash before it is a lost email. I choose the lost email over the duplicate, and log the message id so the gap can be resent by hand.

Build the relay first, before any consumer exists, and make it the only file in the service that imports the channel. Then give every consumer a processed_messages table and an insert in the same transaction as its effect. Once those two pieces are in place, the broker can fail, the relay can be killed mid-batch, and the handler can go on returning 200 at 09:14:02 without any of it reaching the ledger twice.

Written by Md Nasim Anjum, senior full-stack engineer in Manchester. He builds payment orchestration, KYC and KYB compliance platforms and conversational AI.

Get in touchAll writingRSS

More writing