Outbox Pattern sem Redis: quando uma tabela no PostgreSQL já basta
Se a sua aplicação já depende de PostgreSQL, boa parte das filas de jobs pode sair com uma tabela, uma transação curta e
FOR UPDATE SKIP LOCKED— sem Redis, sem RabbitMQ, sem um daemon a mais para operar. Este artigo vai do problema do Transactional Outbox até o worker completo: schema, lease, fencing token, retry com backoff, retenção e testes determinísticos. E diz com honestidade onde esse desenho quebra.
1. O problema real: duas escritas que deveriam ser uma só
Considere uma operação que cria um pedido e precisa disparar o envio de um e-mail:
BEGIN INSERT INTO orders (...) ENQUEUE email_job COMMIT
Se o INSERT e o enqueue acontecem em sistemas diferentes, surgem duas falhas clássicas:
- O pedido é confirmado no banco, mas o processo cai antes de publicar o job. O e-mail nunca é enviado.
- O job é publicado, mas a transação do pedido sofre rollback. O worker recebe uma tarefa que aponta para algo que não existe ou que nunca deveria ter sido confirmado.
Publicar depois do COMMIT não elimina o problema: o processo pode cair entre essas duas operações. Publicar dentro da transação também não resolve, porque o broker não participa da mesma transação distribuída.
O Transactional Outbox grava o evento na mesma transação que altera os dados de negócio. Um relay separado lê a tabela e publica o trabalho depois. Assim, se a transação confirma, o registro do outbox também existe; se ela desfaz, os dois desaparecem juntos. É a solução documentada no padrão de Transactional Outbox.
⚠️ O padrão não promete exactly-once.
O relay pode publicar a mensagem e cair antes de marcar o registro como enviado. Na próxima tentativa, a mensagem é publicada de novo. Por isso o consumidor precisa ser idempotente — e o schema precisa te dar como fazer isso (veja a colunadeduplication_keyna seção 4).
2. Outbox e job queue não são a mesma coisa
Uma fila de jobs interna distribui trabalho entre workers que pertencem à própria aplicação. Uma tabela jobs com reivindicação concorrente costuma bastar.
Um outbox é uma garantia de consistência entre uma mudança no banco e uma mensagem que precisa sair dali. O registro é consumido por um relay que publica em RabbitMQ, Kafka, webhook ou outro serviço.
Os dois padrões podem compartilhar a mesma tabela e o mesmo polling, mas o contrato é diferente:
| Dimensão | Job queue interna | Transactional Outbox |
|---|---|---|
| Objetivo | Um worker processa a tarefa | Uma mudança confirmada gera uma mensagem |
| Consumidor | Worker da própria aplicação | Relay que publica fora do banco |
| Estado crítico | status, attempts, lease_* | published_at, relay_* |
| Precisa de lease? | Sim, o trabalho é longo | Não necessariamente — o relay é curto |
| Semântica | at-least-once | at-least-once |
Os dois schemas estão nas seções 4 e 5.
3. FOR UPDATE SKIP LOCKED e a reivindicação atômica
O PostgreSQL mantém locks de linha até o fim da transação. FOR UPDATE impede que outra transação modifique ou reivindique a mesma linha enquanto o lock existir. Com SKIP LOCKED, em vez de esperar por uma linha já travada, a consulta pula essa linha e procura outra. A própria documentação ressalta que isso produz uma visão inconsistente e, portanto, não é uma opção geral de leitura; o uso apropriado inclui tabelas que funcionam como filas, com múltiplos consumidores.
Uma reivindicação típica é atômica:
WITH claimed AS ( SELECT id FROM jobs WHERE status = 'pending' AND run_at <= now() ORDER BY priority DESC, run_at, id LIMIT 20 FOR UPDATE SKIP LOCKED ) UPDATE jobs AS j SET status = 'running', attempts = j.attempts + 1, lease_until = now() + interval '30 seconds', lease_token = gen_random_uuid(), worker_id = $1, updated_at = now() FROM claimed WHERE j.id = claimed.id RETURNING j.*;
O UPDATE ... RETURNING é o que importa: seleção e mudança de estado acontecem na mesma transação. Dois workers não recebem a mesma linha por causa da combinação do lock com o SKIP LOCKED. O guia da Neon mostra o mesmo mecanismo aplicado a uma fila de jobs e mede o custo de cada alternativa de reivindicação.
Use uma ordenação estável. Sem ORDER BY, o PostgreSQL não promete ordem; em uma fila isso cria comportamento difícil de explicar. E alinhe um índice parcial ao predicado — sem ele, o SKIP LOCKED varre a tabela inteira a cada poll:
CREATE INDEX jobs_ready_idx ON jobs (priority DESC, run_at, id) WHERE status = 'pending';
💡
SKIP LOCKEDnão é lock distribuído mágico.
Ele coordena consumidores que acessam o mesmo PostgreSQL. Se o worker fizer todo o trabalho externo dentro da transação, o lock fica aberto durante uma operação lenta, prejudicando throughput e vacuum. O desenho normal é: reivindicar rápido, confirmar, executar fora da transação, finalizar com atualização condicional.
4. Schema de produção: a fila de jobs
-- gen_random_uuid() é nativo desde o PostgreSQL 13. -- A extensão pgcrypto só é necessária em versões anteriores. CREATE TYPE job_status AS ENUM ('pending', 'running', 'succeeded', 'failed'); CREATE TABLE jobs ( id uuid PRIMARY KEY DEFAULT gen_random_uuid(), kind text NOT NULL, payload jsonb NOT NULL, status job_status NOT NULL DEFAULT 'pending', priority integer NOT NULL DEFAULT 0, run_at timestamptz NOT NULL DEFAULT now(), attempts integer NOT NULL DEFAULT 0, max_attempts integer NOT NULL DEFAULT 8, lease_until timestamptz, lease_token uuid, worker_id text, deduplication_key text, last_error text, created_at timestamptz NOT NULL DEFAULT now(), updated_at timestamptz NOT NULL DEFAULT now() ); -- Deduplicação: o consumidor pergunta "já processei isso?" sem colisão entre kinds. CREATE UNIQUE INDEX jobs_dedup_idx ON jobs (kind, deduplication_key) WHERE deduplication_key IS NOT NULL; CREATE INDEX jobs_ready_idx ON jobs (priority DESC, run_at, id) WHERE status = 'pending'; CREATE INDEX jobs_expired_lease_idx ON jobs (lease_until) WHERE status = 'running'; CREATE INDEX jobs_retention_idx ON jobs (status, updated_at);
A transação de negócio e a criação do job ficam juntas:
BEGIN; INSERT INTO books (id, title, owner_id) VALUES ($1, $2, $3); INSERT INTO jobs (kind, payload, deduplication_key) VALUES ( 'index-book', jsonb_build_object('book_id', $1), 'index-book:' || $1::text ) ON CONFLICT (kind, deduplication_key) WHERE deduplication_key IS NOT NULL DO NOTHING; COMMIT;
Se a inserção do livro falhar, o job também não existe. Se o COMMIT ocorrer, o worker encontra o job depois, ainda que o processo web morra no instante seguinte. E se o usuário clicar duas vezes, a segunda tentativa colide no índice e o DO NOTHING evita o job duplicado.
5. Schema do outbox, quando o destino é externo
Para o relay que publica em um broker ou webhook, o estado crítico é published_at, não status:
CREATE TABLE outbox ( id uuid PRIMARY KEY DEFAULT gen_random_uuid(), aggregate_type text NOT NULL, aggregate_id text NOT NULL, event_type text NOT NULL, payload jsonb NOT NULL, attempts integer NOT NULL DEFAULT 0, relay_until timestamptz, relay_token uuid, published_at timestamptz, created_at timestamptz NOT NULL DEFAULT now() ); CREATE INDEX outbox_unpublished_idx ON outbox (created_at) WHERE published_at IS NULL;
O relay reivindica fora da transação de publicação, e só marca depois de publicar:
WITH relay AS ( SELECT id FROM outbox WHERE published_at IS NULL ORDER BY created_at, id LIMIT 100 FOR UPDATE SKIP LOCKED ) UPDATE outbox AS o SET relay_until = now() + interval '30 seconds', relay_token = gen_random_uuid(), attempts = o.attempts + 1 FROM relay WHERE o.id = relay.id RETURNING o.id, o.event_type, o.payload;
-- chamado só depois de o publish ter retornado com sucesso UPDATE outbox SET published_at = now(), relay_until = NULL, relay_token = NULL WHERE id = $1 AND relay_token = $2;
⚠️
published_at: não marque junto com a reivindicação.
Se você escrevepublished_atno mesmoUPDATEque seleciona, e o relay cai antes de publicar, o evento é perdido para sempre. A ordem obrigatória é: reivindicar, publicar, marcar.
6. Lease, heartbeat e fencing token
Um lock de linha só vive até o fim da transação. Como o worker precisa processar o job por vários segundos ou minutos, ele libera o lock depois da reivindicação e usa um lease: lease_until indica até quando a posse é válida.
O worker renova o lease com uma escrita condicionada:
UPDATE jobs SET lease_until = now() + interval '30 seconds', updated_at = now() WHERE id = $1 AND status = 'running' AND worker_id = $2 AND lease_token = $3 AND lease_until > now();
O lease_token é um UUID novo a cada reivindicação. Ele funciona como fencing token: impede que um worker antigo sobrescreva o estado depois que outro worker recuperou o job.
Note que o heartbeat não rotaciona o token. Se ele rotacionasse, o worker se expulsaria de si mesmo no meio do trabalho. O token só muda na reivindicação.
A finalização também é sempre condicional:
UPDATE jobs SET status = 'succeeded', lease_until = NULL, lease_token = NULL, worker_id = NULL, updated_at = now() WHERE id = $1 AND status = 'running' AND lease_token = $2;
Se o número de linhas afetadas for zero, o worker perdeu o lease e não deve sobrescrever o estado atual.
⚠️ Fencing não desfaz o que já foi feito lá fora.
Esse fencing protege as transições na tabela, mas não apaga efeitos externos já realizados por um worker zumbi. Se o worker antigo enviou um e-mail antes de perder a conexão, o banco não pode desfazer esse e-mail. Para efeitos externos, use idempotency keys, operações condicionais no destino ou um segundo outbox. Heartbeat não transforma uma operação externa em exactly-once.
Um reclaimer devolve leases vencidos à fila — e precisa aplicar o mesmo backoff do retry, senão um job que derruba o processo volta em loop apertado:
UPDATE jobs SET status = CASE WHEN attempts >= max_attempts THEN 'failed' ELSE 'pending' END, run_at = CASE WHEN attempts >= max_attempts THEN run_at ELSE now() + make_interval(secs => LEAST(300, power(2, attempts)::int)) END, lease_until = NULL, worker_id = NULL, lease_token = NULL, updated_at = now() WHERE status = 'running' AND lease_until < now();
O dimensionamento do lease precisa considerar relógios, duração máxima esperada, pausas de GC, suspensão de máquina e atraso de rede. Lease curto demais provoca duplicatas; lease longo demais retarda a recuperação.
7. Retry com backoff exponencial
Falhas transitórias não devem ser reprocessadas em loop apertado:
atraso = min(base * 2^(attempts - 1), máximo) + jitter
O jitter aleatório evita que muitos jobs que falharam juntos voltem simultaneamente.
UPDATE jobs SET status = CASE WHEN attempts >= max_attempts THEN 'failed' ELSE 'pending' END, run_at = CASE WHEN attempts >= max_attempts THEN run_at ELSE now() + ($2::interval) END, last_error = $3, lease_until = NULL, worker_id = NULL, lease_token = NULL, updated_at = now() WHERE id = $1 AND status = 'running' AND lease_token = $4;
Jobs em failed precisam ser observáveis e reprocessáveis manualmente. E o sistema tem que distinguir erro permanente — payload inválido, por exemplo — de erro transitório, como a indisponibilidade de uma API. Nem tudo deve ser retentado.
8. Retenção: a tabela cresce, e ninguém planeja isso
Uma tabela de jobs só funciona em produção se alguém apagar o que já terminou. Sem isso, o índice jobs_ready_idx permanece indexando um WHERE status = 'pending' que nunca esvazia, e o vacuum passa a ser o gargalo.
DELETE FROM jobs WHERE status IN ('succeeded', 'failed') AND updated_at < now() - interval '7 days';
Rode em lotes pequenos, fora do pico, acompanhe o crescimento e ajuste a janela. Se o volume for alto o suficiente para o DELETE em massa gerar bloat, particione por range em created_at e derrube partições inteiras em vez de apagar linha a linha — o DROP PARTITION é praticamente instantâneo e libera o espaço de uma vez.
9. Implementação Go: o núcleo do worker
O exemplo usa database/sql e lib/pq só para deixar o mecanismo explícito. Em produção, prefira pgx.
package jobs import ( "context" "database/sql" "encoding/json" "fmt" "time" ) type Job struct { ID string Kind string Payload json.RawMessage Attempts int LeaseToken string } type Clock interface{ Now() time.Time } type RealClock struct{} func (RealClock) Now() time.Time { return time.Now().UTC() } const claimSQL = ` WITH claimed AS ( SELECT id FROM jobs WHERE status = 'pending' AND run_at <= now() ORDER BY priority DESC, run_at, id LIMIT $2 FOR UPDATE SKIP LOCKED ) UPDATE jobs AS j SET status = 'running', attempts = j.attempts + 1, lease_until = now() + $3::interval, lease_token = gen_random_uuid(), worker_id = $1, updated_at = now() FROM claimed WHERE j.id = claimed.id RETURNING j.id, j.kind, j.payload, j.attempts, j.lease_token` func Claim(ctx context.Context, db *sql.DB, worker string, batch int, lease time.Duration) ([]Job, error) { rows, err := db.QueryContext(ctx, claimSQL, worker, batch, lease.String()) if err != nil { return nil, err } defer rows.Close() var result []Job for rows.Next() { var j Job if err := rows.Scan(&j.ID, &j.Kind, &j.Payload, &j.Attempts, &j.LeaseToken); err != nil { return nil, err } result = append(result, j) } return result, rows.Err() } func Complete(ctx context.Context, db *sql.DB, jobID, token string) error { res, err := db.ExecContext(ctx, ` UPDATE jobs SET status = 'succeeded', lease_until = NULL, lease_token = NULL, worker_id = NULL, updated_at = now() WHERE id = $1 AND status = 'running' AND lease_token = $2`, jobID, token) if err != nil { return err } n, err := res.RowsAffected() if err != nil { return err } if n != 1 { return fmt.Errorf("lease lost for job %s", jobID) } return nil }
Duas observações que o exemplo curto esconde:
- Consuma todas as linhas. O lock de linha só é liberado quando o result set é esgotado ou fechado. Interromper o loop no meio e segurar o
rowsaberto mantém linhas travadas por mais tempo do que o necessário. - O
lease.String()vira'30s', que o PostgreSQL aceita comointerval. Funciona, mas em driver que faz parse de tipo estrito, prefira passar a duração explicitamente.
Um worker completo ainda precisa de heartbeat, cancelamento por context.Context, classificação de erros, backoff, métricas e encerramento gracioso.
10. Testes determinísticos
Testes de concorrência não devem depender de time.Sleep(40 * time.Millisecond). Esse número é resultado de uma máquina, uma carga e uma configuração específicas — não é uma propriedade do padrão. Se um número de benchmark vira asserção, o teste vira flaky por construção.
Injete um relógio controlável:
type FakeClock struct{ t time.Time } func (c *FakeClock) Now() time.Time { return c.t } func (c *FakeClock) Advance(d time.Duration) { c.t = c.t.Add(d) }
Com isso, o cenário de perda de lease é verificado sem esperar:
- o worker A reivindica o job e recebe o token
T1; - o relógio avança além do lease;
- o reclaimer devolve o job à fila;
- o worker B reivindica o mesmo job e recebe
T2; - uma finalização com
T1afeta zero linhas; - uma finalização com
T2conclui o job.
O teste de SKIP LOCKED deve abrir duas transações reais contra PostgreSQL, bloquear explicitamente uma linha na primeira e verificar que a segunda escolhe outra. Mockar SQL não prova a semântica de lock do banco — prova que você sabe qual string passar.
11. Quando o PostgreSQL é a escolha certa
✅ Quando usar?
O app já depende de PostgreSQL, obrigatoriamente
O volume é moderado e previsível
Os jobs são internos ao mesmo sistema
Payload, estado e histórico precisam ser consultáveis em SQL
A implantação precisa funcionar em Electron, Docker, Raspberry Pi ou ambiente offline
Adicionar outro daemon custa mais operação do que a capacidade que ele traria
A prioridade é reduzir componentes e manter uma fonte de verdade
❌ Quando NÃO usar?
O throughput ou a latência que o banco não entrega
Você precisa de pub/sub e fan-out para muitos consumidores
Roteamento por tópico, exchange ou prioridade avançada é requisito
Isolar a carga de jobs do banco transacional é obrigatório
Os consumidores vivem fora do domínio do banco
Retenção e distribuição de mensagens são o produto
Os custos reais: polling, crescimento e limpeza da tabela, bloat, vacuum, tuning, visibilidade de filas, conexões de workers, manutenção de leases e cuidado com transações longas. O banco não é de graça — ele apenas concentra infraestrutura que você já tem.
12. Quando uma tecnologia dedicada se paga
Redis e BullMQ. BullMQ é uma fila Node.js construída sobre Redis, com concorrência horizontal, scripts Lua, pipelining e entrega at-least-once no pior caso (documentação do BullMQ). É natural para equipes Node que querem delayed jobs, repeatable jobs, eventos de fila e dashboards prontos.
RabbitMQ. Oferece acknowledgements de consumidor e publisher confirms, o que dá semântica de pelo menos uma vez, redelivery e consumidores idempotentes, além de filas duráveis e quorum queues replicadas (guia de confiabilidade). É o certo quando roteamento, exchanges e múltiplos consumidores são requisitos de primeira classe.
Celery. Sistema de task queue para distribuir trabalho entre processos ou máquinas. A documentação atual é explícita: os transportes RabbitMQ e Redis são feature complete, e todo o resto — incluindo SQLite para desenvolvimento local — é experimental (documentação do Celery). Ou seja, o transporte SQL que às vezes aparece como "opção" do Celery é exatamente a categoria que a própria documentação marca como experimental. Ele não elimina as decisões sobre throughput, leases, observabilidade e idempotência.
A virada recente: BullMQ com backend PostgreSQL
O enquadramento acima mudou nos últimos anos. O BullMQ passou a oferecer um backend PostgreSQL oficial, com paridade de API completa — filas, workers, flows, schedulers, rate limiting, priorização, delayed jobs, deduplicação, métricas e eventos. O código da aplicação fica idêntico entre backends.
Por baixo, ele mapeia as primitivas do BullMQ para recursos do PostgreSQL: as transições de estado rodam como funções SQL dentro de transações, e o wait-for-job bloqueante é implementado com LISTEN/NOTIFY em vez do BZPOPMIN do Redis. As tabelas vivem num schema próprio (padrão bullmq) e as migrações são explícitas, idempotentes e protegidas por advisory lock.
Isso muda a decisão. Ela deixa de ser binária:
| Opção | Quando escolher |
|---|---|
Tabela sua, com SKIP LOCKED | Volume moderado, domínio simples, zero dependência, você quer controlar cada linha |
| BullMQ com backend PostgreSQL | Equipe Node, quer a API e o ecossistema do BullMQ, sem adicionar Redis |
| BullMQ ou Celery com Redis/RabbitMQ | Volume alto, múltiplos consumidores, roteamento, isolamento de carga |
E o que dizer sobre performance
A documentação do backend PostgreSQL do BullMQ publica uma tabela que vale mais que qualquer benchmark caseiro — e vem com a ressalva de que os números são de um Mac M-series com PostgreSQL local e jobs sem operação, e servem só para ordem de grandeza (docs do backend):
| Operação | PostgreSQL (jobs/s) | Redis, mesma máquina (jobs/s) |
|---|---|---|
add(), um a um | 7.000 | 7.500 |
add() concorrente (Promise.all) | 15.000 | 38.000 |
addBulk(), em lote e concorrente | 45.000 | 52.000 |
| Processamento, 1 worker | 2.300 | 6.000 |
| Processamento, concorrência 8 a 32 | 11.000 | 18.000 |
A leitura honesta: enfileirar é praticamente o mesmo, e o gargalo real está no processamento, onde o PostgreSQL fica cerca de 1,5 a 2 vezes atrás. Enfileirar sequencial e em lote é mais próximo justamente porque a escrita durável é amortizada — por round-trip ou pelo lote.
A comparação correta não é "qual é mais rápido em um benchmark". É: qual contrato de entrega, escala, isolamento e operação o sistema precisa cumprir?
13. Checklist de decisão
Escolha a fila no PostgreSQL quando a resposta for "sim" para a maior parte:
- O banco já é obrigatório?
- O trabalho pode esperar alguns segundos de polling?
- A taxa de jobs cabe com folga na capacidade de escrita e leitura do banco?
- Os workers compartilham esse banco?
- É aceitável processar novamente em caso de crash?
- O consumidor pode ser idempotente — e você tem
deduplication_keypara provar? - Alguém é dono da rotina de limpeza da tabela?
- A equipe quer evitar um serviço adicional em instalações locais?
Adicione Redis, RabbitMQ ou outra tecnologia quando houver necessidade clara de:
- Throughput ou latência que o PostgreSQL não entrega
- Pub/sub e fan-out para muitos consumidores
- Roteamento por tópicos, exchanges ou prioridades avançadas
- Isolamento da carga de jobs em relação ao banco transacional
- Consumidores fora do domínio do banco
- Retenção e distribuição de mensagens como produto próprio
- Recursos operacionais já maduros na equipe
E se você é Node e quer a API do BullMQ sem carregar um Redis: olhe o backend PostgreSQL antes de concluir que a única escolha é escrever a tabela na mão.
TL;DR
- O Transactional Outbox resolve consistência entre a transação de negócio e a intenção de publicar. O
SKIP LOCKEDresolve reivindicação concorrente sem workers se bloqueando. São problemas diferentes. - A garantia real é at-least-once. Fencing token protege a tabela, não o mundo externo. Idempotência é obrigação do consumidor — e precisa de coluna no schema para ser provável.
- Um
lease_tokennovo por reivindicação é o que impede o worker zumbi de sobrescrever o estado de quem o recuperou. Heartbeat renova o lease, mas não rotaciona o token. - Backoff sem jitter provoca tempestade; lease sem heartbeat expõe o job a recuperação precoce; tabela sem rotina de limpeza vira o gargalo.
- PostgreSQL dá conta de volume moderado com folga. A decisão não é "tabela ou broker", e sim qual contrato de entrega e operação você precisa cumprir.
Leitura adicional
- Pattern: Transactional Outbox — Chris Richardson
- PostgreSQL Documentation: SELECT — referência de
FOR UPDATE SKIP LOCKED - Queue system using SKIP LOCKED in Postgres — Neon
- RabbitMQ: Reliability Guide
- BullMQ: What is BullMQ
- BullMQ: PostgreSQL backend
- Celery: Introduction — Celery 5.6