Blog

Outbox transacional: gravar no banco e publicar o evento sem perder nenhum dos dois

Um marketplace confirmava o pedido no banco e, na linha seguinte, publicava pedido.confirmado no broker para o faturamento, o estoque e o e-mail ao cliente. Numa manutenção do broker de 4 minutos, 37 pedidos foram pagos, gravados como confirmados e nunca faturados. Ninguém viu: nada lançou exceção que importasse, a API devolveu 200 e o painel de erros ficou limpo. O financeiro descobriu três dias depois, ao conciliar. O defeito não estava na publicação nem na gravação, estava em fazer as duas como se fossem uma só. Este artigo mostra por que gravar no banco e publicar um evento são duas escritas que nunca serão atômicas, como o padrão outbox move o evento para dentro da mesma transação do dado, como um relay publica com entrega ao menos uma vez e sem furar a ordem por agregado, como o consumidor absorve a duplicata que essa garantia traz, como operar a tabela e como provar com falhas injetadas que nenhum evento se perde.

2026-10-08 / Arquitetura / 17 min

01

Por que gravar e publicar são duas escritas, e uma sempre pode falhar sozinha

O banco e o broker são sistemas diferentes, com processos, redes e falhas diferentes. Não existe COMMIT que cubra os dois. Qualquer sequência que você escolha deixa uma janela em que um lado foi feito e o outro não: um deploy no meio da requisição, uma queda de rede, um timeout do broker, um erro de serialização, um OOM kill. Essa janela dura milissegundos, e por isso o problema passa em teste e aparece quando o volume e a quantidade de deploys sobem. Com 1.200 pedidos por dia e um incidente de infraestrutura por mês, é questão de tempo.

Inverter a ordem não resolve, só troca o sintoma. Publicar antes de gravar cria um evento sobre algo que não existe: o faturamento emite nota de um pedido que o banco recusou. Gravar antes de publicar cria um dado sem evento: o pedido está confirmado e ninguém foi avisado. Publicar dentro da transação, antes do COMMIT, é o pior dos mundos porque parece seguro: se o COMMIT falhar depois, o evento já saiu.

EstratégiaFalha entre os dois passosResultadoDetectável?
Gravar e depois publicarProcesso cai ou broker recusa após o COMMITPedido confirmado sem eventoSó na conciliação, dias depois
Publicar e depois gravarBanco recusa ou processo cai após o publishEvento de um pedido que não existeQuando o consumidor não acha o pedido
Publicar dentro da transaçãoCOMMIT falha após o publishEvento de uma gravação desfeitaRaramente
Retry em memória após o COMMITProcesso reinicia com a fila de retry na memóriaPedido confirmado sem eventoSó na conciliação
Outbox transacionalQualquer umEvento publicado no mínimo uma vez, ou transação desfeita por inteiroSim, a tabela mostra o atraso
Escrita dupla (o que quebra)                       Outbox (o que segura)

app -- UPDATE pedidos ----------> banco             app -- BEGIN ---------------------> banco
app -- publish(pedido.confirmado) -> broker           |-- UPDATE pedidos
        ^                                             |-- INSERT outbox (evento)
        |                                             '-- COMMIT  (os dois, ou nenhum)
   queda aqui = pedido gravado, evento perdido
   ordem inversa = evento publicado, pedido nao     relay -- SELECT ... FOR UPDATE SKIP LOCKED
                                                       |-- publish(evento) ----------> broker
                                                       '-- UPDATE outbox SET publicado_em
                                                    consumidor -- INSERT inbox (message_id) ON CONFLICT
                                                       '-- trata o evento so se for a primeira vez

02

A tabela outbox: o evento nasce na mesma transação do dado

A ideia central é trocar a escrita no broker por uma escrita no próprio banco do negócio. Em vez de publicar, a aplicação insere uma linha em uma tabela outbox dentro da mesma transação que altera o pedido. Como é a mesma transação, a atomicidade que o banco já oferece passa a cobrir o evento: ou o pedido muda e o evento existe, ou nada aconteceu. A publicação real fica para outro processo, que lê a tabela. O esquema abaixo inclui o índice parcial, que mantém a leitura dos pendentes barata mesmo com milhões de linhas já publicadas, e a tabela inbox do lado do consumidor, usada na seção 4.

-- Fila de saida no mesmo banco do negocio. seq define a ordem; id viaja como messageId.
CREATE TABLE outbox (
  seq                  bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
  id                   uuid        NOT NULL DEFAULT gen_random_uuid() UNIQUE,
  aggregate_id         text        NOT NULL,
  tipo                 text        NOT NULL,
  payload              jsonb       NOT NULL,
  criado_em            timestamptz NOT NULL DEFAULT now(),
  publicado_em         timestamptz,
  tentativas           int         NOT NULL DEFAULT 0,
  ultimo_erro          text,
  proxima_tentativa_em timestamptz NOT NULL DEFAULT now()
);

-- Indice parcial: so as linhas pendentes. A tabela cresce, o indice nao.
CREATE INDEX outbox_pendentes ON outbox (seq) WHERE publicado_em IS NULL;

-- Lado do consumidor: um registro por mensagem ja tratada.
CREATE TABLE inbox (
  message_id   uuid PRIMARY KEY,
  processado_em timestamptz NOT NULL DEFAULT now()
);

O lado de escrita fica simples. Note que o INSERT na outbox só acontece se o UPDATE do pedido realmente alterou uma linha: um pedido já confirmado, reenviado por um cliente impaciente, não gera um segundo evento.

// Confirma o pedido e registra o evento na MESMA transacao: os dois ou nenhum.
export async function confirmarPedido(pool, { pedidoId, clienteId, totalCentavos }) {
  const client = await pool.connect();
  try {
    await client.query('BEGIN');
    const { rowCount } = await client.query(
      `UPDATE pedidos SET status = 'confirmado', confirmado_em = now()
        WHERE id = $1 AND status = 'pendente'`,
      [pedidoId],
    );
    if (rowCount === 0) {
      await client.query('ROLLBACK'); // pedido inexistente ou ja confirmado: nenhum evento
      return false;
    }
    await client.query(
      `INSERT INTO outbox (aggregate_id, tipo, payload) VALUES ($1, 'pedido.confirmado', $2)`,
      [pedidoId, JSON.stringify({ pedidoId, clienteId, totalCentavos })],
    );
    await client.query('COMMIT');
    return true;
  } catch (erro) {
    await client.query('ROLLBACK');
    throw erro;
  } finally {
    client.release();
  }
}

O payload deve carregar o que o consumidor precisa para agir sem voltar ao banco do produtor: identificador, valores e o momento do fato. Se o consumidor precisar consultar o serviço do pedido para entender o evento, você recriou o acoplamento síncrono que o evento deveria evitar.

03

O relay: publicar da tabela, sem perder e sem furar a ordem

O relay é um laço que lê eventos pendentes, publica no broker e marca como publicados. A consulta de reserva carrega duas decisões. FOR UPDATE SKIP LOCKED deixa várias instâncias do relay trabalharem ao mesmo tempo sem pegar a mesma linha. E o NOT EXISTS só devolve o evento mais antigo pendente de cada agregado: se o pedido 77 tem um evento de confirmação e outro de cancelamento, o segundo só é liberado depois que o primeiro foi publicado, mesmo com relays concorrentes. Um lote nunca contém dois eventos do mesmo agregado, e a ordem por pedido se preserva sem um relay único.

// Pega o evento mais antigo de cada agregado que esta pendente e vencido.
// NOT EXISTS garante a ordem por agregado; SKIP LOCKED deixa varios relays em paralelo.
const SQL_RESERVAR = `
  SELECT o.seq, o.id, o.aggregate_id, o.tipo, o.payload, o.tentativas
    FROM outbox o
   WHERE o.publicado_em IS NULL
     AND o.proxima_tentativa_em <= now()
     AND NOT EXISTS (
       SELECT 1 FROM outbox anterior
        WHERE anterior.aggregate_id = o.aggregate_id
          AND anterior.publicado_em IS NULL
          AND anterior.seq < o.seq)
   ORDER BY o.seq
   LIMIT $1
   FOR UPDATE OF o SKIP LOCKED`;

const dormir = (ms) => new Promise((resolve) => setTimeout(resolve, ms));

export async function publicarLote(pool, broker, { tamanho = 50 } = {}) {
  const client = await pool.connect();
  try {
    await client.query('BEGIN');
    const { rows } = await client.query(SQL_RESERVAR, [tamanho]);
    for (const evento of rows) {
      try {
        // messageId = outbox.id: o consumidor usa para descartar duplicatas
        await broker.publish({
          topic: evento.tipo,
          key: evento.aggregate_id,
          messageId: evento.id,
          value: evento.payload,
        });
        await client.query('UPDATE outbox SET publicado_em = now() WHERE seq = $1', [evento.seq]);
      } catch (erro) {
        const esperaSegundos = Math.min(2 ** evento.tentativas, 300); // 1s, 2s, 4s ... teto de 5 min
        await client.query(
          `UPDATE outbox
              SET tentativas = tentativas + 1,
                  ultimo_erro = $2,
                  proxima_tentativa_em = now() + make_interval(secs => $3)
            WHERE seq = $1`,
          [evento.seq, String(erro.message).slice(0, 500), esperaSegundos],
        );
      }
    }
    await client.query('COMMIT');
    return rows.length;
  } catch (erro) {
    await client.query('ROLLBACK');
    throw erro;
  } finally {
    client.release();
  }
}

export function iniciarRelay(pool, broker, { intervaloMs = 500 } = {}) {
  let ativo = true;
  (async () => {
    while (ativo) {
      try {
        const publicados = await publicarLote(pool, broker);
        if (publicados === 0) await dormir(intervaloMs);
      } catch (erro) {
        console.error('relay da outbox falhou', erro);
        await dormir(2000);
      }
    }
  })();
  return () => {
    ativo = false;
  };
}

A garantia é ao menos uma vez, não exatamente uma. Se o processo cair depois do publish e antes do COMMIT da marcação, a transação do relay é desfeita e o evento será publicado de novo. Isso é o comportamento desejado: perder um evento é irrecuperável, duplicar um evento é tratável. Por isso o messageId vai junto, para o consumidor reconhecer a repetição. Alguns pontos de operação decidem se o relay se comporta bem:

  • Publique com confirmação do broker (acks no Kafka, publisher confirms no RabbitMQ). Um publish que retorna sem confirmação pode ter se perdido, e marcar como publicado seria perder o evento.
  • Use a chave de partição igual ao aggregate_id, para que os eventos do mesmo pedido cheguem na mesma partição, na ordem em que foram publicados.
  • Aplique backoff por evento, como no código, em vez de repetir em laço apertado. Um evento com payload recusado não deve impedir os outros agregados de seguir.
  • Mantenha o lote pequeno, de dezenas de linhas. A transação do relay segura locks enquanto fala com o broker, e um lote grande com um broker lento prende linhas por muito tempo.

04

O consumidor tem que aguentar a duplicata

Como a entrega é ao menos uma vez, o consumidor precisa produzir o mesmo resultado quando recebe a mesma mensagem duas vezes. O faturamento não pode emitir duas notas para o mesmo pedido.confirmado. O mecanismo mais direto é a tabela inbox: registrar o messageId e executar o efeito na mesma transação, com ON CONFLICT DO NOTHING para detectar a repetição. Se o efeito falhar, o ROLLBACK desfaz também a marca, e a reentrega do broker tenta de novo, sem criar a falsa impressão de que a mensagem foi tratada.

// Consumidor idempotente: marcar a mensagem e tratar o efeito na MESMA transacao.
// Se o efeito falhar, a marca some junto no ROLLBACK e a reentrega tenta de novo.
export async function consumir(pool, mensagem, tratar) {
  const client = await pool.connect();
  try {
    await client.query('BEGIN');
    const { rowCount } = await client.query(
      'INSERT INTO inbox (message_id) VALUES ($1) ON CONFLICT DO NOTHING',
      [mensagem.messageId],
    );
    if (rowCount === 0) {
      await client.query('ROLLBACK');
      return 'duplicada';
    }
    await tratar(client, mensagem.value); // ex.: INSERT INTO faturas ... usando o mesmo client
    await client.query('COMMIT');
    return 'processada';
  } catch (erro) {
    await client.query('ROLLBACK');
    throw erro;
  } finally {
    client.release();
  }
}

Quando o efeito não acontece no mesmo banco, como uma chamada a uma API externa, a transação não cobre os dois. Nesse caso, passe a chave de idempotência da própria mensagem para a API de destino, usando o messageId. É o mesmo princípio da cobrança que não pode ser feita duas vezes: a deduplicação mora onde o efeito acontece.

05

Operar a outbox: backlog, limpeza e a alternativa de CDC

A outbox troca um risco de perda silenciosa por um risco visível: o atraso. Isso é uma melhora, desde que alguém olhe. O alerta que importa é a idade do evento mais antigo pendente, não apenas a quantidade: 500 eventos pendentes em um pico que se esvazia em segundos é normal, um único evento de 10 minutos é um relay parado. A tabela também precisa de limpeza, em lotes pequenos, para não segurar lock nem inflar o armazenamento, e a inbox deve guardar as mensagens por mais tempo do que o broker pode reentregar.

-- Painel e alerta: o que importa e a idade do evento mais antigo, nao so a contagem.
SELECT count(*)                                                    AS pendentes,
       coalesce(extract(epoch FROM now() - min(criado_em)), 0)::int AS idade_mais_antiga_s,
       count(*) FILTER (WHERE tentativas >= 5)                     AS com_falha_repetida
  FROM outbox
 WHERE publicado_em IS NULL;

-- Limpeza em lotes pequenos (rodar em loop ate afetar 0 linhas), para nao segurar lock longo.
DELETE FROM outbox
 WHERE seq IN (
   SELECT seq FROM outbox
    WHERE publicado_em < now() - interval '7 days'
    ORDER BY seq
    LIMIT 5000);

-- A inbox guarda mensagens por mais tempo do que o broker pode reentregar (ex.: 30 dias).
DELETE FROM inbox WHERE processado_em < now() - interval '30 days';

O relay por consulta (polling) é a escolha simples e suficiente para a maioria dos sistemas, com latência de centenas de milissegundos. Quando o volume é alto ou a latência precisa ser menor, a leitura do log de transações do banco (CDC, como o Debezium) elimina o polling, ao custo de operar mais uma peça. A tabela compara as opções.

AbordagemGarantiaCusto operacionalQuando usar
Outbox com relay por consultaAo menos uma vez, ordem por agregadoBaixo: um processo e uma tabelaPadrão. Até alguns milhares de eventos por segundo
Outbox com CDC (log do banco)Ao menos uma vez, ordem do logMédio: conector, slot de replicação, monitoramentoVolume alto ou latência de dezenas de milissegundos
Publicar após o COMMIT com retry em memóriaPode perder em queda do processoBaixoSó quando perder o evento é aceitável
Transação distribuída (2PC/XA)Atômica entre banco e brokerAlto: poucos brokers suportam, trava recursosRaramente justificável
Event sourcingO log de eventos é a fonte da verdadeAlto: muda o modelo inteiroQuando o histórico já é o produto

06

Provar que não perde: falhas injetadas em vez de confiança

Um teste que só confere o caminho feliz não prova nada sobre a outbox. As garantias aparecem quando algo falha, então o teste precisa provocar a falha. Os três testes abaixo cobrem as propriedades que importam: um pedido que não muda não gera evento, um broker fora do ar não perde o evento e a recuperação publica uma única vez, e uma entrega duplicada é descartada pelo consumidor. Eles usam um banco de teste real, porque o comportamento de transação e de SKIP LOCKED não existe em um mock.

import test from 'node:test';
import assert from 'node:assert/strict';
import { pool, limparTabelas, criarPedido } from './setup.js'; // banco de teste real, nao mock
import { confirmarPedido } from '../src/confirmar-pedido.js';
import { publicarLote } from '../src/relay.js';
import { consumir } from '../src/consumidor.js';

const brokerQueFalha = () => ({ publicadas: [], falhar: true,
  async publish(msg) {
    if (this.falhar) throw new Error('broker indisponivel');
    this.publicadas.push(msg);
  } });

test('pedido ja confirmado nao deixa evento orfao', async () => {
  await limparTabelas();
  const pedidoId = await criarPedido({ status: 'confirmado' });
  assert.equal(await confirmarPedido(pool, { pedidoId, clienteId: 'c1', totalCentavos: 1000 }), false);
  const { rows } = await pool.query('SELECT 1 FROM outbox');
  assert.equal(rows.length, 0);
});

test('broker fora do ar nao perde o evento e a recuperacao publica uma vez', async () => {
  await limparTabelas();
  const pedidoId = await criarPedido({ status: 'pendente' });
  await confirmarPedido(pool, { pedidoId, clienteId: 'c1', totalCentavos: 1000 });

  const broker = brokerQueFalha();
  await publicarLote(pool, broker);
  const { rows: [pendente] } = await pool.query('SELECT publicado_em, tentativas FROM outbox');
  assert.equal(pendente.publicado_em, null);
  assert.equal(pendente.tentativas, 1);

  broker.falhar = false;
  await pool.query('UPDATE outbox SET proxima_tentativa_em = now()'); // pula o backoff no teste
  await publicarLote(pool, broker);
  await publicarLote(pool, broker); // segunda rodada nao pode republicar
  assert.equal(broker.publicadas.length, 1);
});

test('entrega duplicada e descartada pelo consumidor', async () => {
  await limparTabelas();
  const mensagem = { messageId: '0b1f6a52-6d8c-4a39-9c1e-3f1d2a7c9e10', value: { pedidoId: 'p1' } };
  let efeitos = 0;
  const tratar = async () => { efeitos += 1; };
  assert.equal(await consumir(pool, mensagem, tratar), 'processada');
  assert.equal(await consumir(pool, mensagem, tratar), 'duplicada');
  assert.equal(efeitos, 1);
});

Em produção, a conciliação continua útil como rede de segurança: uma consulta diária que cruza pedidos confirmados nas últimas 24 horas com eventos publicados detecta qualquer caminho de código que alguém ainda tenha deixado de fora da outbox, como um script de correção que atualiza o pedido por fora da função confirmarPedido.

FAQ

Perguntas frequentes

Por que não usar uma transação distribuída (2PC) entre o banco e o broker?

Porque a maioria dos brokers não participa de XA, e onde participa o custo é alto: a transação segura recursos nos dois lados enquanto espera o coordenador, e uma falha do coordenador deixa transações em dúvida bloqueando linhas. A outbox obtém o resultado prático, nenhum evento perdido, usando só a transação local que o banco já oferece, e aceita a duplicata ocasional, que o consumidor idempotente absorve.

A outbox não vira um gargalo por escrever uma linha a mais em toda transação?

O custo é um INSERT adicional na mesma transação, normalmente abaixo de um milissegundo, sem round trip extra porque usa a mesma conexão. O que pode virar problema é a tabela crescer sem limpeza e a leitura dos pendentes varrer linhas antigas. O índice parcial sobre publicado_em IS NULL e a remoção em lotes dos eventos já publicados mantêm o custo estável com o tempo.

Dá para garantir entrega exatamente uma vez com a outbox?

Entre a outbox e o broker não: a garantia é ao menos uma vez, porque o relay pode publicar e cair antes de marcar. O que se obtém é o efeito exatamente uma vez, combinando a entrega ao menos uma vez com um consumidor idempotente que descarta o messageId repetido. Para o negócio o resultado é o mesmo, uma nota emitida por pedido, e é o único jeito honesto de prometer isso em sistemas distribuídos.

O evento precisa existir se, e somente se, o dado existir

Gravar no banco e publicar no broker nunca serão atômicos, e qualquer ordem entre os dois passos deixa uma janela de perda ou de evento fantasma que só aparece em produção, na conciliação, dias depois. A outbox resolve movendo o evento para dentro da transação do dado, publicando por um relay com entrega ao menos uma vez e ordem por agregado, e deixando o consumidor idempotente absorver a duplicata. O que sobra é um risco mensurável, o atraso do evento mais antigo, em vez de uma perda silenciosa. Com testes que injetam falha no broker e na entrega, e uma conciliação diária como rede de segurança, a pergunta deixa de ser se algum evento se perdeu e passa a ser quanto tempo ele levou para sair.