Se você já precisou explicar para o time de negócio por que uma transferência foi debitada duas vezes, sabe que o problema quase nunca está na regra de negócio. Ele mora no espaço entre os serviços: a mensagem que chegou duplicada, o evento que foi publicado antes do commit, o consumidor que caiu no meio do processamento. É nesse espaço que a mensageria deixa de ser um detalhe de infraestrutura e vira uma decisão de engenharia sobre consistência.

Este é o primeiro artigo de uma série sobre mensageria em Java e Spring. Cada artigo parte de uma ferramenta e tenta responder duas perguntas: como ela funciona por baixo dos panos e como usá-la sem criar uma falsa sensação de confiabilidade.

Para isso, vou usar um cenário que exige algumas garantias de verdade: um fluxo de transferências financeiras com análise antifraude, liquidação, notificações e múltiplos serviços, cada um com seu próprio banco. O exemplo usa Java 25, Spring Boot 4.1.1 e Apache Kafka 4.3.1.

A ideia, porém, não é transformar o artigo em um tutorial de “copiar e colar” de configuração. Antes do código, quero construir o modelo mental do Kafka. Assim fica mais fácil entender por que aparecem termos como topic, partition, offset, acks, consumer group, ISR e idempotência em praticamente todo projeto sério com Kafka.

O que é Kafka e onde ele entra no problema?

Kafka é uma plataforma distribuída para publicar, armazenar e consumir eventos. Na prática, ele funciona como um log distribuído: os produtores acrescentam registros ao final de um fluxo, e os consumidores acompanham esse fluxo usando uma posição chamada offset.

Uma forma simples de visualizar isso é imaginar um extrato de movimentações. Em vez de alguém simplesmente dizer “o saldo agora é R$ 1.000”, você registra os fatos que aconteceram: “R$ 200 foram debitados”, “R$ 500 foram creditados”, e assim por diante. Outros sistemas podem ler esses fatos e construir sua própria visão do que aconteceu.

É justamente aí que Kafka se torna interessante. Um mesmo evento pode ser consumido por sistemas diferentes sem que o produtor precise conhecer cada um deles. O antifraude pode ler uma transferência, o ledger pode ler a decisão aprovada e o serviço de notificações pode ler o resultado final.

Kafka foi desenhado para trabalhar com alto volume, distribuição, retenção e reprocessamento. Ele não substitui automaticamente um banco relacional, nem resolve sozinho as transações entre sistemas diferentes. No exemplo deste artigo, o banco continua sendo a fonte da verdade para o saldo e o Kafka apenas transporta os fatos entre os serviços.

Um modelo mental simples

Se você estiver começando, guarde esta cadeia na cabeça:

Aplicação
   │
   │ publica
   ▼
Producer ──▶ Broker ──▶ Topic ──▶ Partition ──▶ Consumer
                    │                 │
                    │                 └── Offset
                    └── Réplicas / ISR

Cada termo resolve um problema diferente:

ConceitoO que éPense nisso como
ProducerAplicação que publica registrosQuem escreve um fato
BrokerServidor Kafka que armazena e serve os registrosUm nó do cluster
TopicNome lógico do fluxo de eventosUma categoria de eventos
PartitionDivisão de um topic em sequências ordenadasUma fila paralela
OffsetPosição de um registro dentro da partiçãoNúmero da linha no log
ConsumerAplicação que lê registrosQuem reage ao fato
Consumer groupConjunto de consumidores que divide as partiçõesUm time processando a mesma fila
ReplicaCópia de uma partição em outro brokerBackup operacional
ISRRéplicas que estão em sincronia com o líderCópias consideradas aptas

A linguagem ajuda, mas existe uma regra ainda mais importante: Kafka garante ordem por partição, não ordem global do topic. Se um tópico tem 12 partições, são 12 sequências que podem ser processadas em paralelo.

Tópicos, partições, chaves e offsets

Um tópico é dividido em partições. Cada partição é uma sequência ordenada de registros, identificados por um offset crescente. O registro é normalmente enviado com uma chave (key), e essa chave ajuda o produtor a decidir em qual partição o registro será armazenado.

No nosso exemplo, escolhemos a conta de origem como chave:

Tópico: pagamentos.transferencia.analisada.v1   (chave = conta de origem)

Partição 0: [CC-104 #0] [CC-104 #1] [CC-104 #2] ...
Partição 1: [CC-220 #0] [CC-220 #1] [CC-220 #2] ...
Partição 2: [CC-555 #0] [CC-042 #1] ...

Isso significa que as transferências da CC-104 ficam na mesma partição e mantêm sua ordem relativa. Ao mesmo tempo, CC-104 e CC-220 podem ser processadas em paralelo porque estão em partições diferentes.

Por que a quantidade de partições importa?

A quantidade de partições define, na prática, quanto paralelismo um consumer group consegue explorar. Um grupo com três consumidores pode processar até três partições em paralelo; um grupo com 20 consumidores não terá 20 unidades de trabalho úteis se o tópico tiver apenas 12 partições.

É por isso que “vou começar com uma partição e aumento depois” merece cuidado. O Kafka permite aumentar a quantidade de partições, mas o mapeamento da chave pode mudar. Uma chave que antes caía na partição 2 pode passar a cair em outra partição depois do aumento. Registros antigos continuam onde estavam, então não é seguro tratar o crescimento de partições como se fosse uma operação neutra para a ordenação por chave.

E o offset?

O offset não identifica globalmente um evento no Kafka. Ele é uma posição dentro de uma partição.

Se a partição 4 está no offset 100, isso significa que o consumidor já avançou até aquela posição naquela partição. O conjunto de offsets do consumer group é que permite ao Kafka retomar o consumo depois de uma reinicialização.

Isso também explica uma característica importante: Kafka não apaga um registro só porque um consumidor leu. O registro segue no log enquanto a política de retenção permitir. Um segundo grupo pode começar a consumir do mesmo tópico sem interferir no primeiro.

Consumer groups: paralelismo sem duplicar trabalho entre consumidores

Dentro de um consumer group, uma partição é atribuída a apenas um consumidor por vez. Isso faz o grupo dividir o trabalho.

Topic: transferencias

Partition 0 ──▶ Consumer A
Partition 1 ──▶ Consumer B
Partition 2 ──▶ Consumer C
Partition 3 ──▶ Consumer A

Se eu criar outro grupo, ele recebe os mesmos eventos de forma independente:

Grupo "ledger"           Grupo "notificacoes"
Partition 0 ─▶ A         Partition 0 ─▶ X
Partition 1 ─▶ B         Partition 1 ─▶ Y
Partition 2 ─▶ C         Partition 2 ─▶ Z

Esse é o fan-out. O produtor publicou uma vez, mas vários grupos conseguem consumir aquele fato para finalidades diferentes.

Uma confusão comum para quem está começando é achar que dois consumidores com grupos diferentes “competem” pela mensagem. Não competem. Eles possuem posições independentes no mesmo log.

O que significa acks?

acks não é um comando do Kafka. É uma configuração do producer que define quanto de confirmação ele exige do cluster antes de considerar a publicação concluída.

Os valores principais são:

acksO que aconteceConsequência
0O producer não espera confirmação do brokerMenor latência, mas sem confirmação de recebimento
1O líder confirma depois de gravar localmenteO líder pode falhar antes da replicação e perder o registro
allO líder espera as réplicas em sincronia confirmaremMaior durabilidade; é a opção usada no exemplo

Quando acks=all, a confirmação não significa “o dado está em todos os brokers do cluster”. Significa que o líder espera a confirmação do conjunto de réplicas em sincronia (ISR) daquela partição.

É aqui que entra min.insync.replicas.

Suponha um topic com fator de replicação 3:

Partition 0

Broker 1  ── líder
Broker 2  ── réplica
Broker 3  ── réplica

ISR = [Broker 1, Broker 2, Broker 3]

Se acks=all e min.insync.replicas=2, uma escrita só pode ser aceita enquanto pelo menos duas réplicas estiverem em sincronia. Com as três réplicas no ISR, o acks=all continua esperando a confirmação do conjunto inteiro em sincronia; min.insync.replicas=2 funciona como um limite mínimo de segurança para que a escrita seja aceita.

Essa diferença é pequena no texto, mas enorme na prática. acks=all e min.insync.replicas trabalham juntos: um define o nível de confirmação pedido pelo producer, e o outro impede a escrita quando o cluster perdeu réplicas demais para aquela política.

Replicação, ISR e retenção

Kafka distribui uma partição entre brokers e mantém réplicas. Uma réplica que está acompanhando o líder faz parte do ISR. Quando um broker fica para trás, pode deixar de ser considerado in-sync.

Outro conceito importante é retenção. Kafka não pensa em uma mensagem como “um item que desaparece depois de ser consumido”. O registro continua armazenado até que a política do tópico ou do broker determine sua remoção. Uma configuração como retention.ms=604800000 representa uma retenção de até sete dias por tempo.

Isso dá a Kafka uma característica muito útil: um consumidor novo pode ler eventos antigos com auto.offset.reset=earliest, enquanto um consumidor que já possui offsets continua de onde parou.

Não confunda retenção com backup ou arquivo permanente. Se o tópico tem sete dias de retenção, depois desse período os registros podem ser removidos. Se o evento precisa ser preservado por motivos contábeis ou regulatórios, o desenho deve tratar isso explicitamente.

At-most-once, at-least-once e exactly-once

Outra forma de entender Kafka é perguntar: o que acontece se o consumidor cair no momento errado?

Existe uma ressalva que vale repetir: exactly-once do Kafka não transforma uma escrita no PostgreSQL em parte da transação Kafka. No momento em que um banco externo entra no fluxo, você precisa resolver essa fronteira de consistência com outras técnicas, como o Transactional Outbox e a idempotência.

As versões usadas neste artigo

O artigo foi revisado em 02 de outubro de 2026. Como Kafka, Spring Boot e Spring Kafka evoluem rapidamente, fixe as versões nos seus projetos e confira as notas de release antes de copiar um exemplo para produção. Nesta revisão, uso as versões estáveis verificadas nas documentações oficiais: Kafka 4.3.1, Spring Boot 4.1.1 e Spring Kafka 4.1.1.

ComponenteVersãoObservação
Java25 (LTS)Records, text blocks, sealed interfaces e pattern matching em switch
Spring Boot4.1.1Baseado no Spring Framework 7
Spring for Apache Kafka4.1.1Cliente Kafka 4.2.1; Jackson 3; sem dependência do Spring Retry
Apache Kafka (broker)4.3.1Somente KRaft
PostgreSQL18Um banco por serviço
Jackson3JsonMapper, exceções não checadas

Algumas mudanças dessas versões pesam no dia a dia e aparecem ao longo do artigo:

Uma distinção que costuma gerar confusão: a versão do broker e a versão da biblioteca cliente não precisam ser iguais. Neste artigo, o broker é Kafka 4.3.1, enquanto o Spring Kafka 4.1.1 traz o cliente Kafka 4.2.1. O Spring Boot gerencia essa compatibilidade por meio do seu BOM, por isso não faz sentido adicionar manualmente uma versão diferente de kafka-clients sem uma razão concreta.

A arquitetura do exemplo

O fluxo é uma coreografia de eventos. Nenhum serviço chama outro diretamente; cada um reage a fatos publicados em tópicos:

Cliente ──POST /transferencias──▶ transferencias-service
                                        │  (Postgres: transferencia + outbox)
                                        ▼  relay do outbox
                          [pagamentos.transferencia.solicitada.v1]
                                        │
                                        ▼
                               antifraude-service        Kafka → Kafka, transacional
                                  │             │
                    (aprovada)    │             │  (bloqueada)
                                  ▼             │
             [pagamentos.transferencia.analisada.v1]
                                  │             │
                                  ▼             │
                            ledger-service      │        Postgres: contas, lançamentos, outbox
                             │          │       │
                  (liquidada)│          │(rejeitada)
                             ▼          ▼       ▼
        [pagamentos.transferencia.liquidada.v1]  [pagamentos.transferencia.rejeitada.v1]
                             │                          │
             ┌───────────────┴──────────────┬───────────┘
             ▼                              ▼
   notificacoes-service            transferencias-service
   (retry topics + DLT)            (projeção do status)

Alguns princípios de design orientam esse desenho:

Os tópicos e suas características:

TópicoProdutorConsumidores (grupo)Partições
pagamentos.transferencia.solicitada.v1transferenciasantifraude12
pagamentos.transferencia.analisada.v1antifraudeledger12
pagamentos.transferencia.liquidada.v1ledgernotificacoes, transferencias-projecao12
pagamentos.transferencia.rejeitada.v1antifraude e ledgernotificacoes, transferencias-projecao12

O tópico rejeitada tem dois produtores, o que simplifica o exemplo mas mistura responsabilidades. Em produção vale considerar tópicos separados por produtor ou, no mínimo, um único dono documentado.

Instalando Kafka localmente

Para acompanhar este artigo, você não precisa montar um cluster com vários servidores. Um broker local já é suficiente para entender os conceitos e executar o projeto. O próprio projeto Apache Kafka publica a imagem oficial apache/kafka e fornece uma forma simples de subir um broker em modo KRaft com Docker.

O exemplo usa apache/kafka:4.3.1. A partir do Kafka 4.0, ZooKeeper deixou de fazer parte do modo suportado, Kafka 4.3 opera em KRaft. Para desenvolvimento local, isso simplifica bastante a infraestrutura: temos o próprio Kafka cuidando dos metadados do cluster por meio do quorum baseado em Raft.

Opção 1: subir Kafka com Docker

A forma mais rápida é:

docker pull apache/kafka:4.3.1

docker run --name kafka -p 9092:9092 apache/kafka:4.3.1

A documentação oficial também apresenta esse modelo de execução para a imagem JVM do Kafka 4.3.1.

Opção 2: Docker Compose

No projeto do artigo, prefiro Compose porque o Kafka e o PostgreSQL sobem juntos:

# docker-compose.yml (desenvolvimento)
services:
  kafka:
    image: apache/kafka:4.3.1
    container_name: kafka
    ports:
      - "9092:9092"

  postgres:
    image: postgres:18-alpine
    environment:
      POSTGRES_USER: app
      POSTGRES_PASSWORD: app
      POSTGRES_DB: transferencias
    ports:
      - "5432:5432"

Suba o ambiente com:

docker compose up -d

Confira se os containers estão em execução:

docker compose ps

Para parar o ambiente:

docker compose down

Como estamos usando um único broker, os tópicos locais precisam de replication-factor=1. No código, a propriedade app.kafka.replicas representa isso e vale 1 localmente e 3 no cenário de produção descrito mais adiante.

Opção 3: rodar Kafka sem Docker

Também é possível baixar a distribuição do Kafka e iniciar o servidor diretamente com os scripts oficiais. Para Kafka 4.3.1, o quickstart usa Java 17 ou superior:

tar -xzf kafka_2.13-4.3.1.tgz
cd kafka_2.13-4.3.1

KAFKA_CLUSTER_ID="$(bin/kafka-storage.sh random-uuid)"
bin/kafka-storage.sh format --standalone -t "$KAFKA_CLUSTER_ID" -c config/server.properties
bin/kafka-server-start.sh config/server.properties

Para este artigo, porém, Docker deixa o ambiente mais reproduzível e reduz a quantidade de passos que não têm relação direta com o desenvolvimento da aplicação.

Primeiros comandos do Kafka

Antes de conectar Java ao Kafka, vale fazer um teste manual. Isso ajuda a separar problemas de infraestrutura de problemas no código.

Quando usamos a imagem oficial via Compose, os utilitários ficam dentro do container em /opt/kafka/bin/.

Criar um tópico

docker compose exec kafka /opt/kafka/bin/kafka-topics.sh \
  --bootstrap-server kafka:9092 \
  --create \
  --topic pagamentos.transferencia.solicitada.v1 \
  --partitions 3 \
  --replication-factor 1

--partitions 3 cria três partições. --replication-factor 1 é necessário aqui porque só existe um broker.

Listar tópicos

docker compose exec kafka /opt/kafka/bin/kafka-topics.sh \
  --bootstrap-server kafka:9092 \
  --list

Inspecionar um tópico

docker compose exec kafka /opt/kafka/bin/kafka-topics.sh \
  --bootstrap-server kafka:9092 \
  --describe \
  --topic pagamentos.transferencia.solicitada.v1

O --describe mostra informações como quantidade de partições, fator de replicação, líder, réplicas e ISR. Esse é um dos primeiros comandos que costumo usar quando quero saber como um tópico realmente está distribuído.

Publicar manualmente

docker compose exec -it kafka /opt/kafka/bin/kafka-console-producer.sh \
  --bootstrap-server kafka:9092 \
  --topic pagamentos.transferencia.solicitada.v1

Depois, digite mensagens e pressione Enter. Cada linha será publicada como um registro.

Para experimentar com chave e valor, use:

docker compose exec -it kafka /opt/kafka/bin/kafka-console-producer.sh \
  --bootstrap-server kafka:9092 \
  --topic pagamentos.transferencia.solicitada.v1 \
  --property parse.key=true \
  --property key.separator=:

Agora você pode digitar algo como:

CC-104:{"transferenciaId":"001","valorCentavos":30000}
CC-104:{"transferenciaId":"002","valorCentavos":10000}
CC-220:{"transferenciaId":"003","valorCentavos":5000}

A chave CC-104 passa a fazer parte da decisão de particionamento. É exatamente essa ideia que usaremos no producer Java.

Consumir desde o início

docker compose exec -it kafka /opt/kafka/bin/kafka-console-consumer.sh \
  --bootstrap-server kafka:9092 \
  --topic pagamentos.transferencia.solicitada.v1 \
  --from-beginning

--from-beginning é especialmente útil em laboratório porque permite enxergar também os registros que já existiam no tópico.

Ver grupos de consumidores

docker compose exec kafka /opt/kafka/bin/kafka-consumer-groups.sh \
  --bootstrap-server kafka:9092 \
  --list

Para detalhar um grupo:

docker compose exec kafka /opt/kafka/bin/kafka-consumer-groups.sh \
  --bootstrap-server kafka:9092 \
  --describe \
  --group ledger

Esse comando mostra, entre outras informações, CURRENT-OFFSET, LOG-END-OFFSET e LAG. O lag é a distância entre aquilo que já foi produzido e aquilo que o grupo já confirmou como consumido.

Alterar configuração de um tópico

Configurações de tópico podem ser alteradas com kafka-configs.sh. Por exemplo:

docker compose exec kafka /opt/kafka/bin/kafka-configs.sh \
  --bootstrap-server kafka:9092 \
  --entity-type topics \
  --entity-name pagamentos.transferencia.solicitada.v1 \
  --alter \
  --add-config retention.ms=604800000

No ambiente do artigo, boa parte disso será declarada via código Spring (NewTopic), mas conhecer a CLI é importante para investigação e operação.

Configurações que aparecem no código

Agora que os conceitos estão definidos, as propriedades do application.yml deixam de parecer uma lista aleatória de chaves.

Configurações do producer

ConfiguraçãoPara que serveNeste artigo
acksQuanto de confirmação o producer exigeall
enable.idempotenceEvita duplicatas causadas por retries do próprio producertrue
delivery.timeout.msTempo máximo para o envio completar, incluindo retries do cliente120000
linger.msDá alguns milissegundos para formar lotes maiores5
compression.typeComprime os lotes enviadoszstd
key.serializerConverte a chave Java para bytesStringSerializer
value.serializerConverte o valor Java para bytesStringSerializer
transaction-id-prefixAtiva produtores transacionais do KafkaUsado no antifraude

Uma distinção importante: enable.idempotence=true não significa que seu negócio é idempotente. Ele protege contra certas duplicações causadas pelo próprio producer. Não impede que o mesmo evento seja publicado novamente por outra instância, por um relay reiniciado ou por uma aplicação que decidiu reenviar a operação. Para isso, o consumidor continua precisando de uma chave de idempotência adequada.

Configurações do consumer

ConfiguraçãoPara que serveNeste artigo
group-idIdentifica o consumer groupledger, antifraude, notificacoes
enable-auto-commitDefine se o cliente confirma offsets automaticamentefalse
auto-offset-resetDefine de onde começar quando não há offset válidoearliest
isolation.levelControla a leitura de registros transacionaisread_committed
max-poll-recordsQuantidade máxima de registros devolvidos por poll50 no ledger
max.poll.interval.msTempo máximo esperado entre pollsDeve comportar o tempo de processamento
concurrencyQuantidade de consumidores/threads gerados pelo container Spring3 no ledger

Mais adiante vamos ver essas propriedades aplicadas aos serviços. O objetivo aqui é criar o mapa mental: producer publica, consumer lê, group divide o trabalho, offset marca o progresso e as configurações definem as garantias e o comportamento operacional.

Ambiente Java e dependências

As dependências principais de cada serviço:

<parent>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-parent</artifactId>
    <version>4.1.1</version>
    <relativePath/>
</parent>

<properties>
    <java.version>25</java.version>
</properties>

<dependencies>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-webmvc</artifactId>
    </dependency>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-jdbc</artifactId>
    </dependency>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-kafka</artifactId>
    </dependency>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-flyway</artifactId>
    </dependency>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-actuator</artifactId>
    </dependency>
    <dependency>
        <groupId>org.postgresql</groupId>
        <artifactId>postgresql</artifactId>
        <scope>runtime</scope>
    </dependency>
</dependencies>

O starter spring-boot-starter-kafka é o ponto de integração entre Spring Boot e Spring for Apache Kafka. As propriedades ficam em spring.kafka.*, e o Boot também suporta propriedades adicionais do cliente Kafka quando elas não aparecem diretamente na classe de configuração do Boot.

Uma nota sobre o Flyway: para PostgreSQL, além do starter, pode ser necessário o módulo flyway-database-postgresql. As dependências de teste aparecem na seção de testes.

Nos trechos de código a seguir, omiti os imports óbvios de java.* e do próprio Spring quando não ajudam a entender o exemplo. Os imports que mudaram com o Spring Boot 4, o Jackson 3 e o Spring Kafka 4 estão sempre presentes.

Os contratos dos eventos

Eventos são contratos entre serviços. Com Java 25, records são a forma mais natural de representá-los: imutáveis, concisos e bem suportados pelo Jackson 3.

package com.exemplo.pagamentos.transferencias.eventos;

import java.time.Instant;
import java.util.UUID;

public record TransferenciaSolicitada(
        UUID eventoId,
        UUID transferenciaId,
        String contaOrigem,
        String contaDestino,
        long valorEmCentavos,
        String moeda,
        Instant ocorridoEm) {}
public record TransferenciaAnalisada(
        UUID eventoId,
        UUID transferenciaId,
        String contaOrigem,
        String contaDestino,
        long valorEmCentavos,
        String moeda,
        int scoreRisco,
        Instant ocorridoEm) {}

public record TransferenciaLiquidada(UUID eventoId, UUID transferenciaId, Instant ocorridoEm) {}

public record TransferenciaRejeitada(
        UUID eventoId, UUID transferenciaId, String motivo, Instant ocorridoEm) {}

Duas escolhas merecem justificativa:

Serviço 1: transferencias-service e o Transactional Outbox

O problema do dual write

A operação natural ao receber uma transferência seria gravar no banco e depois publicar o evento no Kafka. São dois sistemas, sem transação atômica entre eles, e qualquer ordem falha em algum cenário:

O Transactional Outbox resolve isso transformando “publicar” em uma escrita no mesmo banco, dentro da mesma transação da regra de negócio. A gravação da transferência e a do evento na tabela outbox_evento são atômicas. Um processo separado, o relay, lê o outbox e publica no Kafka.

Kafka tem transações, mas elas cobrem apenas leitura e escrita dentro do Kafka. O Spring Kafka permite sincronizar uma transação de banco com uma transação Kafka, mas isso é best effort: o commit dos dois lados não é atômico, e uma falha no segundo commit deixa os sistemas divergentes. Para dinheiro, essa margem de erro não serve.

O esquema

-- V1__schema.sql (transferencias-service)
CREATE TABLE transferencia (
    id                UUID PRIMARY KEY,
    conta_origem      VARCHAR(32)  NOT NULL,
    conta_destino     VARCHAR(32)  NOT NULL,
    valor_centavos    BIGINT       NOT NULL CHECK (valor_centavos > 0),
    moeda             CHAR(3)      NOT NULL,
    status            VARCHAR(20)  NOT NULL,
    motivo_rejeicao   VARCHAR(255),
    criada_em         TIMESTAMPTZ  NOT NULL DEFAULT now(),
    atualizada_em     TIMESTAMPTZ  NOT NULL DEFAULT now()
);

CREATE TABLE outbox_evento (
    id            UUID PRIMARY KEY,
    topico        VARCHAR(255) NOT NULL,
    chave         VARCHAR(64)  NOT NULL,
    tipo          VARCHAR(100) NOT NULL,
    payload       TEXT         NOT NULL,
    criado_em     TIMESTAMPTZ  NOT NULL DEFAULT now(),
    publicado_em  TIMESTAMPTZ
);

CREATE INDEX idx_outbox_pendentes ON outbox_evento (criado_em) WHERE publicado_em IS NULL;

O índice parcial mantém a consulta do relay rápida mesmo com milhões de linhas históricas, porque só indexa o que ainda não foi publicado.

A porta e o adaptador de publicação

Se você leu o artigo sobre Hexagonal Architecture, a estrutura vai parecer familiar: a regra de negócio não deve saber que o Kafka existe. Ela depende de uma porta (PublicadorDeEventos), e o outbox é o adaptador que a implementa.

package com.exemplo.pagamentos.transferencias.aplicacao;

import java.util.UUID;

public interface PublicadorDeEventos {

    void publicar(String topico, String chave, String tipo, UUID eventoId, Object evento);
}
package com.exemplo.pagamentos.transferencias.infra.outbox;

import com.exemplo.pagamentos.transferencias.aplicacao.PublicadorDeEventos;
import org.springframework.jdbc.core.simple.JdbcClient;
import org.springframework.stereotype.Component;
import org.springframework.transaction.annotation.Propagation;
import org.springframework.transaction.annotation.Transactional;
import tools.jackson.databind.json.JsonMapper;

import java.util.UUID;

@Component
class OutboxPublicador implements PublicadorDeEventos {

    private final JdbcClient jdbc;
    private final JsonMapper jsonMapper;

    OutboxPublicador(JdbcClient jdbc, JsonMapper jsonMapper) {
        this.jdbc = jdbc;
        this.jsonMapper = jsonMapper;
    }

    @Override
    @Transactional(propagation = Propagation.MANDATORY)
    public void publicar(String topico, String chave, String tipo, UUID eventoId, Object evento) {
        jdbc.sql("""
                INSERT INTO outbox_evento (id, topico, chave, tipo, payload)
                VALUES (:id, :topico, :chave, :tipo, :payload)
                """)
                .param("id", eventoId)
                .param("topico", topico)
                .param("chave", chave)
                .param("tipo", tipo)
                .param("payload", jsonMapper.writeValueAsString(evento))
                .update();
    }
}

Dois detalhes fazem diferença aqui. Primeiro, Propagation.MANDATORY faz o método falhar se for chamado fora de uma transação, o que impede o erro clássico de gravar o evento sem estar atrelado à regra de negócio. Segundo, o Jackson 3 lança JacksonException, que é uma exceção não checada, então o código não precisa de try/catch para serializar.

O caso de uso e a idempotência na borda

A API recebe um header Idempotency-Key com um UUID gerado pelo cliente. Usamos esse valor como identificador da transferência: se o cliente reenviar a mesma requisição (por timeout, por exemplo), a inserção conflita e nenhum novo evento é gerado.

package com.exemplo.pagamentos.transferencias.aplicacao;

import com.exemplo.pagamentos.transferencias.eventos.TransferenciaSolicitada;
import org.springframework.jdbc.core.simple.JdbcClient;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;

import java.time.Instant;
import java.util.Optional;
import java.util.UUID;

@Service
public class TransferenciaService {

    private final JdbcClient jdbc;
    private final PublicadorDeEventos publicador;

    public TransferenciaService(JdbcClient jdbc, PublicadorDeEventos publicador) {
        this.jdbc = jdbc;
        this.publicador = publicador;
    }

    @Transactional
    public Transferencia solicitar(UUID id, String origem, String destino, long valor, String moeda) {
        if (origem.equals(destino)) {
            throw new TransferenciaInvalidaException("Conta de origem e destino devem ser diferentes");
        }

        int inseridas = jdbc.sql("""
                INSERT INTO transferencia (id, conta_origem, conta_destino, valor_centavos, moeda, status)
                VALUES (:id, :origem, :destino, :valor, :moeda, 'PENDENTE')
                ON CONFLICT (id) DO NOTHING
                """)
                .param("id", id)
                .param("origem", origem)
                .param("destino", destino)
                .param("valor", valor)
                .param("moeda", moeda)
                .update();

        if (inseridas == 1) {
            var evento = new TransferenciaSolicitada(
                    UUID.randomUUID(), id, origem, destino, valor, moeda, Instant.now());
            publicador.publicar(Topicos.SOLICITADA, origem, "TransferenciaSolicitada",
                    evento.eventoId(), evento);
        }

        return buscar(id).orElseThrow();
    }

    @Transactional(readOnly = true)
    public Optional<Transferencia> buscar(UUID id) {
        return jdbc.sql("""
                SELECT id, conta_origem, conta_destino, valor_centavos, moeda, status
                FROM transferencia WHERE id = :id
                """)
                .param("id", id)
                .query(Transferencia.class)
                .optional();
    }
}
public record Transferencia(
        UUID id, String contaOrigem, String contaDestino,
        long valorCentavos, String moeda, String status) {}

public final class Topicos {

    public static final String SOLICITADA = "pagamentos.transferencia.solicitada.v1";
    public static final String ANALISADA = "pagamentos.transferencia.analisada.v1";
    public static final String LIQUIDADA = "pagamentos.transferencia.liquidada.v1";
    public static final String REJEITADA = "pagamentos.transferencia.rejeitada.v1";

    private Topicos() {}
}

public class TransferenciaInvalidaException extends RuntimeException {

    public TransferenciaInvalidaException(String mensagem) {
        super(mensagem);
    }
}

Note que a chave do evento é a conta de origem, o que garante que as transferências de uma mesma conta cheguem ordenadas ao consumidor.

O controller devolve 202 Accepted, porque o processamento é assíncrono. O cliente consulta o status depois:

package com.exemplo.pagamentos.transferencias.web;

@RestController
@RequestMapping("/transferencias")
public class TransferenciaController {

    private final TransferenciaService service;

    public TransferenciaController(TransferenciaService service) {
        this.service = service;
    }

    @PostMapping
    public ResponseEntity<TransferenciaResponse> solicitar(
            @RequestHeader("Idempotency-Key") UUID chaveIdempotencia,
            @RequestBody SolicitarTransferenciaRequest request) {

        Transferencia t = service.solicitar(chaveIdempotencia, request.contaOrigem(),
                request.contaDestino(), request.valorEmCentavos(), request.moeda());

        return ResponseEntity.accepted().body(TransferenciaResponse.de(t));
    }

    @GetMapping("/{id}")
    public ResponseEntity<TransferenciaResponse> consultar(@PathVariable UUID id) {
        return service.buscar(id)
                .map(t -> ResponseEntity.ok(TransferenciaResponse.de(t)))
                .orElseGet(() -> ResponseEntity.notFound().build());
    }

    public record SolicitarTransferenciaRequest(
            String contaOrigem, String contaDestino, long valorEmCentavos, String moeda) {}

    public record TransferenciaResponse(UUID id, String status) {

        static TransferenciaResponse de(Transferencia t) {
            return new TransferenciaResponse(t.id(), t.status());
        }
    }
}

Em produção, valide o corpo (Bean Validation) e mapeie TransferenciaInvalidaException para 400 em um @RestControllerAdvice. Omiti isso para manter o foco na mensageria.

O relay: do banco para o Kafka

O relay é um processo agendado que lê o outbox e publica no Kafka. O ponto sensível é permitir várias instâncias do serviço sem que duas leiam as mesmas linhas: FOR UPDATE SKIP LOCKED faz cada instância pegar um lote diferente.

package com.exemplo.pagamentos.transferencias.infra.outbox;

import org.apache.kafka.clients.producer.ProducerRecord;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.jdbc.core.simple.JdbcClient;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import org.springframework.transaction.PlatformTransactionManager;
import org.springframework.transaction.support.TransactionTemplate;

import java.nio.charset.StandardCharsets;
import java.util.List;
import java.util.UUID;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;

@Component
public class OutboxRelay {

    private static final Logger log = LoggerFactory.getLogger(OutboxRelay.class);
    private static final int TAMANHO_LOTE = 100;
    private static final long TIMEOUT_ENVIO_MS = 10_000;

    private final JdbcClient jdbc;
    private final KafkaTemplate<String, String> kafka;
    private final TransactionTemplate transacao;

    public OutboxRelay(JdbcClient jdbc, KafkaTemplate<String, String> kafka,
                       PlatformTransactionManager transactionManager) {
        this.jdbc = jdbc;
        this.kafka = kafka;
        this.transacao = new TransactionTemplate(transactionManager);
    }

    @Scheduled(fixedDelay = 200)
    public void publicarPendentes() {
        transacao.executeWithoutResult(status -> {
            List<EventoOutbox> lote = jdbc.sql("""
                    SELECT id, topico, chave, tipo, payload
                    FROM outbox_evento
                    WHERE publicado_em IS NULL
                    ORDER BY criado_em
                    LIMIT :limite
                    FOR UPDATE SKIP LOCKED
                    """)
                    .param("limite", TAMANHO_LOTE)
                    .query(EventoOutbox.class)
                    .list();

            if (lote.isEmpty()) {
                return;
            }

            List<CompletableFuture<?>> envios = lote.stream().map(this::enviar).toList();
            aguardar(envios);

            jdbc.sql("UPDATE outbox_evento SET publicado_em = now() WHERE id IN (:ids)")
                    .param("ids", lote.stream().map(EventoOutbox::id).toList())
                    .update();

            log.debug("Outbox: {} evento(s) publicado(s)", lote.size());
        });
    }

    private CompletableFuture<?> enviar(EventoOutbox evento) {
        var registro = new ProducerRecord<>(evento.topico(), evento.chave(), evento.payload());
        registro.headers().add("evento-id", evento.id().toString().getBytes(StandardCharsets.UTF_8));
        registro.headers().add("evento-tipo", evento.tipo().getBytes(StandardCharsets.UTF_8));
        return kafka.send(registro);
    }

    private void aguardar(List<CompletableFuture<?>> envios) {
        try {
            CompletableFuture.allOf(envios.toArray(CompletableFuture[]::new))
                    .get(TIMEOUT_ENVIO_MS, TimeUnit.MILLISECONDS);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            throw new IllegalStateException("Publicação do outbox interrompida", e);
        } catch (ExecutionException | TimeoutException e) {
            throw new IllegalStateException("Falha ao publicar lote do outbox", e);
        }
    }

    record EventoOutbox(UUID id, String topico, String chave, String tipo, String payload) {}
}

Lembre de anotar a aplicação com @EnableScheduling. Vale entender o que esse código garante e o que ele não garante:

Com o tempo, a tabela cresce. Uma rotina periódica que apaga eventos já publicados e antigos (DELETE FROM outbox_evento WHERE publicado_em < now() - interval '7 days') mantém o outbox pequeno.

Configuração do produtor

# application.yml (transferencias-service)
spring:
  application:
    name: transferencias-service
  kafka:
    bootstrap-servers: ${KAFKA_BOOTSTRAP_SERVERS:localhost:9092}
    producer:
      acks: all
      compression-type: zstd
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.apache.kafka.common.serialization.StringSerializer
      properties:
        enable.idempotence: true
        delivery.timeout.ms: 120000
        linger.ms: 5
    template:
      observation-enabled: true
    listener:
      observation-enabled: true

app:
  kafka:
    replicas: 1   # 3 em produção

Cada opção tem uma razão de ser:

Os valores de string dos serializadores são intencionais. O payload já é um JSON gerado na hora de gravar o outbox, então o produtor só precisa enviar texto.

Declarando os tópicos

O Spring Boot cria os tópicos declarados como NewTopic na inicialização, via KafkaAdmin. É ótimo em desenvolvimento e em testes; em produção, gerencie tópicos como infraestrutura versionada (Terraform, Strimzi ou equivalente) e desative a criação automática de tópicos nos brokers.

package com.exemplo.pagamentos.transferencias.infra.kafka;

import com.exemplo.pagamentos.transferencias.aplicacao.Topicos;
import org.apache.kafka.clients.admin.NewTopic;
import org.apache.kafka.common.config.TopicConfig;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.config.TopicBuilder;

import java.time.Duration;

@Configuration
class TopicosConfig {

    @Bean
    NewTopic solicitada(@Value("${app.kafka.replicas:3}") int replicas) {
        return TopicBuilder.name(Topicos.SOLICITADA)
                .partitions(12)
                .replicas(replicas)
                .config(TopicConfig.MIN_IN_SYNC_REPLICAS_CONFIG, String.valueOf(Math.min(2, replicas)))
                .config(TopicConfig.RETENTION_MS_CONFIG, String.valueOf(Duration.ofDays(7).toMillis()))
                .build();
    }
}

O número de partições merece cuidado, porque aumentá-lo depois muda o mapeamento chave → partição. Registros novos de uma mesma conta podem passar a cair em outra partição, quebrando a ordem em relação aos antigos. Dimensione com folga desde o início, pensando no paralelismo máximo de consumo que você vai precisar.

A projeção de status

O transferencias-service também é consumidor: ele escuta os resultados finais para atualizar o status da transferência. Como é uma atualização condicional, ela é naturalmente idempotente: só transiciona a partir de PENDENTE.

package com.exemplo.pagamentos.transferencias.infra.kafka;

import com.exemplo.pagamentos.transferencias.eventos.TransferenciaLiquidada;
import com.exemplo.pagamentos.transferencias.eventos.TransferenciaRejeitada;
import org.springframework.jdbc.core.simple.JdbcClient;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;
import tools.jackson.databind.json.JsonMapper;

@Component
class ResultadoTransferenciaListener {

    private final JdbcClient jdbc;
    private final JsonMapper jsonMapper;

    ResultadoTransferenciaListener(JdbcClient jdbc, JsonMapper jsonMapper) {
        this.jdbc = jdbc;
        this.jsonMapper = jsonMapper;
    }

    @KafkaListener(topics = Topicos.LIQUIDADA, groupId = "transferencias-projecao")
    void aoLiquidar(String json) {
        var evento = jsonMapper.readValue(json, TransferenciaLiquidada.class);
        jdbc.sql("""
                UPDATE transferencia SET status = 'LIQUIDADA', atualizada_em = now()
                WHERE id = :id AND status = 'PENDENTE'
                """)
                .param("id", evento.transferenciaId())
                .update();
    }

    @KafkaListener(topics = Topicos.REJEITADA, groupId = "transferencias-projecao")
    void aoRejeitar(String json) {
        var evento = jsonMapper.readValue(json, TransferenciaRejeitada.class);
        jdbc.sql("""
                UPDATE transferencia
                SET status = 'REJEITADA', motivo_rejeicao = :motivo, atualizada_em = now()
                WHERE id = :id AND status = 'PENDENTE'
                """)
                .param("id", evento.transferenciaId())
                .param("motivo", evento.motivo())
                .update();
    }
}

Nada disso está protegido por tratamento de erro customizado, e é aí que mora uma armadilha importante do Spring Kafka. Se você não definir um tratador de erros, o comportamento padrão do DefaultErrorHandler é tentar entregar o registro algumas vezes (nove retentativas, sem espera) e, depois, registrar o erro em log e seguir em frente, descartando a mensagem. Para uma projeção de status isso pode ser tolerável; para o ledger, jamais. Vamos ver como tratar isso corretamente no próximo serviço.

Serviço 2: antifraude-service e o exactly-once dentro do Kafka

O antifraude lê solicitada, calcula um risco e publica em analisada (se aprovada) ou em rejeitada (se bloqueada). É um estágio que só lê e escreve no Kafka, e por isso é o cenário em que as transações do Kafka fazem sentido.

O que as transações do Kafka garantem

O exactly-once do Kafka combina três mecanismos:

  1. Produtor idempotente: elimina duplicatas causadas por retentativas do cliente.
  2. Transações: um conjunto de escritas em vários tópicos e partições, junto com o commit do offset de consumo, é confirmado ou abortado atomicamente.
  3. Isolamento read_committed: consumidores só enxergam registros de transações confirmadas.

No Spring Boot, basta definir um prefixo de transação para o produtor. O Boot configura sozinho um KafkaTransactionManager e o conecta ao container do listener:

# application.yml (antifraude-service)
spring:
  application:
    name: antifraude-service
  kafka:
    bootstrap-servers: ${KAFKA_BOOTSTRAP_SERVERS:localhost:9092}
    producer:
      acks: all
      transaction-id-prefix: tx-antifraude-${HOSTNAME:local}-
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.apache.kafka.common.serialization.StringSerializer
    consumer:
      group-id: antifraude
      enable-auto-commit: false
      auto-offset-reset: earliest
      isolation-level: read_committed
      properties:
        group.protocol: consumer

O prefixo precisa ser único por instância da aplicação. Por isso o exemplo usa o nome do host. Com o EOSMode.V2, que é o único modo suportado hoje, o transactional.id não precisa mais ser reaproveitado entre reinícios para proteger contra “zumbis”, mas continua sendo obrigatório que instâncias simultâneas não o compartilhem.

O listener

package com.exemplo.pagamentos.antifraude;

import com.exemplo.pagamentos.antifraude.eventos.TransferenciaAnalisada;
import com.exemplo.pagamentos.antifraude.eventos.TransferenciaRejeitada;
import com.exemplo.pagamentos.antifraude.eventos.TransferenciaSolicitada;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Component;
import tools.jackson.databind.json.JsonMapper;

import java.nio.charset.StandardCharsets;
import java.time.Instant;
import java.util.UUID;

@Component
public class AnaliseAntifraudeListener {

    private static final int LIMITE_BLOQUEIO = 80;

    private final KafkaTemplate<String, String> kafka;
    private final JsonMapper jsonMapper;
    private final MotorDeRisco motor;

    public AnaliseAntifraudeListener(KafkaTemplate<String, String> kafka,
                                     JsonMapper jsonMapper, MotorDeRisco motor) {
        this.kafka = kafka;
        this.jsonMapper = jsonMapper;
        this.motor = motor;
    }

    @KafkaListener(topics = Topicos.SOLICITADA, groupId = "antifraude")
    public void analisar(ConsumerRecord<String, String> registro) {
        TransferenciaSolicitada t = jsonMapper.readValue(registro.value(), TransferenciaSolicitada.class);
        int score = motor.calcularRisco(t);

        if (score >= LIMITE_BLOQUEIO) {
            var rejeitada = new TransferenciaRejeitada(
                    idDerivado("rejeitada-antifraude", t.transferenciaId()),
                    t.transferenciaId(), "BLOQUEADA_ANTIFRAUDE", Instant.now());
            kafka.send(Topicos.REJEITADA, t.contaOrigem(), jsonMapper.writeValueAsString(rejeitada));
            return;
        }

        var analisada = new TransferenciaAnalisada(
                idDerivado("analisada", t.transferenciaId()),
                t.transferenciaId(), t.contaOrigem(), t.contaDestino(),
                t.valorEmCentavos(), t.moeda(), score, Instant.now());
        kafka.send(Topicos.ANALISADA, t.contaOrigem(), jsonMapper.writeValueAsString(analisada));
    }

    private static UUID idDerivado(String tipo, UUID transferenciaId) {
        return UUID.nameUUIDFromBytes((tipo + ":" + transferenciaId).getBytes(StandardCharsets.UTF_8));
    }
}
@Component
public class MotorDeRisco {

    private static final long LIMITE_VALOR_ALTO_CENTAVOS = 10_000_000L; // R$ 100.000,00

    // Regra de exemplo; um motor real combinaria histórico, dispositivo, geografia etc.
    public int calcularRisco(TransferenciaSolicitada t) {
        return t.valorEmCentavos() >= LIMITE_VALOR_ALTO_CENTAVOS ? 95 : 10;
    }
}

Quando o listener roda dentro da transação iniciada pelo container, cada kafka.send(...) participa dela automaticamente. Ao final do método, o Spring Kafka envia os offsets consumidos à transação e a confirma; se o método lançar uma exceção, a transação é abortada, as escritas ficam invisíveis para consumidores read_committed, e o registro de entrada é relido.

O idDerivado merece atenção. O identificador do evento é determinístico, calculado a partir da transferência, e não um UUID.randomUUID(). Se o antifraude receber a mesma solicitação duas vezes (por causa de uma duplicata do relay do outbox, por exemplo), os dois eventos de saída terão o mesmo identificador, e os consumidores seguintes conseguem reconhecer a duplicata. Isso vai importar no ledger.

O que essa transação não cobre

É tentador pensar que ligar transações resolve todos os problemas de consistência. Não resolve, e as limitações são importantes:

Se o fluxo do seu serviço não é Kafka para Kafka, provavelmente você não precisa de transações Kafka. Precisa de idempotência, que é o tema do serviço seguinte.

Serviço 3: ledger-service e o consumidor idempotente

O ledger é o serviço mais crítico: ele altera saldos. Aqui a regra é clara: um mesmo evento processado duas vezes não pode produzir dois débitos. Como a entrega no Kafka é at-least-once, a garantia precisa vir do processamento.

O esquema

-- V1__schema.sql (ledger-service)
CREATE TABLE conta (
    id              VARCHAR(32) PRIMARY KEY,
    saldo_centavos  BIGINT   NOT NULL CHECK (saldo_centavos >= 0),
    moeda           CHAR(3)  NOT NULL
);

CREATE TABLE lancamento (
    id                BIGSERIAL PRIMARY KEY,
    transferencia_id  UUID         NOT NULL,
    conta_id          VARCHAR(32)  NOT NULL REFERENCES conta (id),
    tipo              CHAR(1)      NOT NULL CHECK (tipo IN ('D', 'C')),
    valor_centavos    BIGINT       NOT NULL CHECK (valor_centavos > 0),
    criado_em         TIMESTAMPTZ  NOT NULL DEFAULT now()
);

CREATE TABLE transferencia_processada (
    transferencia_id  UUID PRIMARY KEY,
    resultado         VARCHAR(20)  NOT NULL,
    processada_em     TIMESTAMPTZ  NOT NULL DEFAULT now()
);

CREATE TABLE outbox_evento (
    id            UUID PRIMARY KEY,
    topico        VARCHAR(255) NOT NULL,
    chave         VARCHAR(64)  NOT NULL,
    tipo          VARCHAR(100) NOT NULL,
    payload       TEXT         NOT NULL,
    criado_em     TIMESTAMPTZ  NOT NULL DEFAULT now(),
    publicado_em  TIMESTAMPTZ
);

CREATE INDEX idx_outbox_pendentes ON outbox_evento (criado_em) WHERE publicado_em IS NULL;

Cada transferência gera duas linhas em lancamento, um débito (D) e um crédito (C): é o modelo de partidas dobradas, que permite auditar qualquer saldo somando os lançamentos. A restrição CHECK (saldo_centavos >= 0) é uma rede de proteção no banco caso alguma regra de aplicação falhe. O outbox_evento é idêntico ao do serviço anterior, junto com o PublicadorDeEventos e o OutboxRelay. Em uma organização real, esse trecho de infraestrutura deveria virar uma biblioteca interna reutilizável.

A chave de idempotência é uma identidade de negócio

A tabela transferencia_processada guarda o identificador da transferência como chave primária. Essa escolha é deliberada: a chave de idempotência precisa representar o fato de negócio (“esta transferência já foi liquidada”), não um detalhe técnico como o offset da mensagem ou um UUID aleatório do evento. Se a chave fosse um identificador aleatório gerado a cada publicação, duplicatas geradas mais acima no fluxo passariam despercebidas.

A regra de negócio

O resultado de uma liquidação é um conjunto fechado de possibilidades. Um sealed interface com records expressa isso bem, e o compilador garante que todos os casos sejam tratados:

package com.exemplo.pagamentos.ledger.aplicacao;

import java.util.UUID;

public sealed interface ResultadoLiquidacao {

    record Liquidada(UUID transferenciaId) implements ResultadoLiquidacao {}

    record Rejeitada(UUID transferenciaId, String motivo) implements ResultadoLiquidacao {}

    record Duplicada(UUID transferenciaId) implements ResultadoLiquidacao {}
}
package com.exemplo.pagamentos.ledger.aplicacao;

import com.exemplo.pagamentos.ledger.eventos.TransferenciaAnalisada;
import com.exemplo.pagamentos.ledger.eventos.TransferenciaLiquidada;
import com.exemplo.pagamentos.ledger.eventos.TransferenciaRejeitada;
import org.springframework.jdbc.core.simple.JdbcClient;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;

import java.time.Instant;
import java.util.List;
import java.util.UUID;

@Service
public class LiquidacaoService {

    private final JdbcClient jdbc;
    private final PublicadorDeEventos publicador;

    public LiquidacaoService(JdbcClient jdbc, PublicadorDeEventos publicador) {
        this.jdbc = jdbc;
        this.publicador = publicador;
    }

    @Transactional
    public ResultadoLiquidacao liquidar(TransferenciaAnalisada t) {
        // 1. Idempotência: só um processamento por transferência consegue inserir esta linha.
        //    Uma entrega concorrente da mesma transferência espera o commit deste e depois não insere nada.
        int novas = jdbc.sql("""
                INSERT INTO transferencia_processada (transferencia_id, resultado)
                VALUES (:id, 'EM_PROCESSAMENTO')
                ON CONFLICT (transferencia_id) DO NOTHING
                """)
                .param("id", t.transferenciaId())
                .update();

        if (novas == 0) {
            return new ResultadoLiquidacao.Duplicada(t.transferenciaId());
        }

        // 2. Trava as duas contas em ordem determinística (evita deadlock entre transferências cruzadas).
        List<Conta> contas = jdbc.sql("""
                SELECT id, saldo_centavos, moeda FROM conta
                WHERE id IN (:ids)
                ORDER BY id
                FOR UPDATE
                """)
                .param("ids", List.of(t.contaOrigem(), t.contaDestino()))
                .query(Conta.class)
                .list();

        Conta origem = contas.stream().filter(c -> c.id().equals(t.contaOrigem())).findFirst().orElse(null);
        Conta destino = contas.stream().filter(c -> c.id().equals(t.contaDestino())).findFirst().orElse(null);

        // 3. Regras de negócio: rejeição é um resultado normal, não um erro.
        if (origem == null || destino == null) {
            return rejeitar(t, "CONTA_INEXISTENTE");
        }
        if (!origem.moeda().equals(t.moeda()) || !destino.moeda().equals(t.moeda())) {
            return rejeitar(t, "MOEDA_INCOMPATIVEL");
        }
        if (origem.saldoCentavos() < t.valorEmCentavos()) {
            return rejeitar(t, "SALDO_INSUFICIENTE");
        }

        // 4. Partidas dobradas: débito na origem, crédito no destino.
        lancar(t.transferenciaId(), origem.id(), "D", -t.valorEmCentavos(), t.valorEmCentavos());
        lancar(t.transferenciaId(), destino.id(), "C", t.valorEmCentavos(), t.valorEmCentavos());

        // 5. O evento de resultado entra no outbox, na mesma transação dos lançamentos.
        var evento = new TransferenciaLiquidada(UUID.randomUUID(), t.transferenciaId(), Instant.now());
        publicador.publicar(Topicos.LIQUIDADA, t.contaOrigem(), "TransferenciaLiquidada",
                evento.eventoId(), evento);
        marcarResultado(t.transferenciaId(), "LIQUIDADA");

        return new ResultadoLiquidacao.Liquidada(t.transferenciaId());
    }

    private void lancar(UUID transferenciaId, String contaId, String tipo, long delta, long valor) {
        jdbc.sql("UPDATE conta SET saldo_centavos = saldo_centavos + :delta WHERE id = :id")
                .param("delta", delta)
                .param("id", contaId)
                .update();

        jdbc.sql("""
                INSERT INTO lancamento (transferencia_id, conta_id, tipo, valor_centavos)
                VALUES (:transferenciaId, :contaId, :tipo, :valor)
                """)
                .param("transferenciaId", transferenciaId)
                .param("contaId", contaId)
                .param("tipo", tipo)
                .param("valor", valor)
                .update();
    }

    private ResultadoLiquidacao rejeitar(TransferenciaAnalisada t, String motivo) {
        var evento = new TransferenciaRejeitada(UUID.randomUUID(), t.transferenciaId(), motivo, Instant.now());
        publicador.publicar(Topicos.REJEITADA, t.contaOrigem(), "TransferenciaRejeitada",
                evento.eventoId(), evento);
        marcarResultado(t.transferenciaId(), "REJEITADA");
        return new ResultadoLiquidacao.Rejeitada(t.transferenciaId(), motivo);
    }

    private void marcarResultado(UUID transferenciaId, String resultado) {
        jdbc.sql("UPDATE transferencia_processada SET resultado = :resultado WHERE transferencia_id = :id")
                .param("resultado", resultado)
                .param("id", transferenciaId)
                .update();
    }

    record Conta(String id, long saldoCentavos, String moeda) {}
}

Vale examinar as decisões desse método, porque cada uma responde a um risco concreto:

Um ponto de honestidade sobre ordem: o Kafka garante a ordem por chave dentro de uma partição, mas a correção dos saldos não depende dela. Uma transferência toca duas contas, e a conta de destino pode ter registros em outras partições sendo processados em paralelo por outros consumidores. O que mantém os saldos corretos são os locks de linha no banco e a idempotência. A ordem do Kafka ajuda a manter a sequência lógica das operações de uma conta, mas não é o mecanismo de consistência.

O listener

package com.exemplo.pagamentos.ledger.infra.kafka;

import com.exemplo.pagamentos.ledger.aplicacao.LiquidacaoService;
import com.exemplo.pagamentos.ledger.aplicacao.ResultadoLiquidacao;
import com.exemplo.pagamentos.ledger.eventos.TransferenciaAnalisada;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;
import tools.jackson.databind.json.JsonMapper;

@Component
class LiquidacaoListener {

    private static final Logger log = LoggerFactory.getLogger(LiquidacaoListener.class);

    private final LiquidacaoService service;
    private final JsonMapper jsonMapper;

    LiquidacaoListener(LiquidacaoService service, JsonMapper jsonMapper) {
        this.service = service;
        this.jsonMapper = jsonMapper;
    }

    @KafkaListener(topics = Topicos.ANALISADA, groupId = "ledger")
    void aoReceber(ConsumerRecord<String, String> registro) {
        TransferenciaAnalisada evento =
                jsonMapper.readValue(registro.value(), TransferenciaAnalisada.class);

        ResultadoLiquidacao resultado = service.liquidar(evento);

        switch (resultado) {
            case ResultadoLiquidacao.Liquidada l ->
                    log.info("Transferência {} liquidada (partição {}, offset {})",
                            l.transferenciaId(), registro.partition(), registro.offset());
            case ResultadoLiquidacao.Rejeitada r ->
                    log.info("Transferência {} rejeitada: {}", r.transferenciaId(), r.motivo());
            case ResultadoLiquidacao.Duplicada d ->
                    log.info("Transferência {} já processada; entrega duplicada ignorada", d.transferenciaId());
        }
    }
}

O consumidor lê o valor como String e faz o parse manualmente. Essa escolha tem uma vantagem forte para sistemas financeiros: se a mensagem for inválida, o Dead Letter Topic recebe exatamente os bytes originais, sem nenhuma transformação, o que facilita a investigação e o reprocessamento.

Configuração do consumidor

# application.yml (ledger-service)
spring:
  application:
    name: ledger-service
  kafka:
    bootstrap-servers: ${KAFKA_BOOTSTRAP_SERVERS:localhost:9092}
    producer:
      acks: all
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.apache.kafka.common.serialization.StringSerializer
      properties:
        enable.idempotence: true
    consumer:
      group-id: ledger
      enable-auto-commit: false
      auto-offset-reset: earliest
      isolation-level: read_committed
      max-poll-records: 50
      properties:
        group.protocol: consumer
    listener:
      ack-mode: record
      concurrency: 3
      observation-enabled: true
    template:
      observation-enabled: true

app:
  kafka:
    replicas: 1   # 3 em produção

Cada parâmetro tem um motivo:

Falhas: retry bloqueante e Dead Letter Topic

Existem três categorias de falha, e cada uma pede uma resposta diferente:

TipoExemploResposta
TransitóriaBanco indisponível por alguns segundosRetentar com espera crescente
Permanente (mensagem venenosa)JSON inválido, campo obrigatório ausenteIr direto para o DLT, sem retentar
Resultado de negócioSaldo insuficienteNão é falha: é um evento rejeitada

No ledger, a ordem por conta importa, então usamos retry bloqueante: o consumidor retenta o mesmo registro antes de avançar, preservando a ordem. O preço é que a partição fica parada durante as tentativas; como as esperas são curtas e limitadas, é um preço aceitável. Esgotadas as tentativas, o registro vai para um Dead Letter Topic (DLT).

package com.exemplo.pagamentos.ledger.infra.kafka;

import org.apache.kafka.common.TopicPartition;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.listener.DeadLetterPublishingRecoverer;
import org.springframework.kafka.listener.DefaultErrorHandler;
import org.springframework.kafka.support.ExponentialBackOffWithMaxRetries;
import tools.jackson.core.JacksonException;

@Configuration
class TratamentoDeErrosConfig {

    @Bean
    DefaultErrorHandler kafkaErrorHandler(KafkaTemplate<String, String> template) {
        var recuperador = new DeadLetterPublishingRecoverer(template,
                (registro, excecao) -> new TopicPartition(registro.topic() + "-dlt", registro.partition()));

        var espera = new ExponentialBackOffWithMaxRetries(5);
        espera.setInitialInterval(500L);
        espera.setMultiplier(2.0);
        espera.setMaxInterval(15_000L);

        var tratador = new DefaultErrorHandler(recuperador, espera);
        tratador.addNotRetryableExceptions(JacksonException.class, IllegalArgumentException.class);
        return tratador;
    }
}

Como o Boot detecta um bean CommonErrorHandler, ele conecta esse tratador ao container de todos os listeners da aplicação. As decisões desse trecho:

Declare o DLT junto dos demais tópicos:

@Bean
NewTopic analisadaDlt(@Value("${app.kafka.replicas:3}") int replicas) {
    return TopicBuilder.name(Topicos.ANALISADA + "-dlt")
            .partitions(12)
            .replicas(replicas)
            .config(TopicConfig.MIN_IN_SYNC_REPLICAS_CONFIG, String.valueOf(Math.min(2, replicas)))
            .build();
}

Um DLT sem ninguém olhando é só um lugar mais elegante de perder mensagens. O recuperador adiciona headers com a causa da falha (tópico, partição e offset originais, classe e mensagem da exceção), que devem alimentar um alerta. Um listener simples que conta as mensagens e registra o contexto já resolve o essencial:

package com.exemplo.pagamentos.ledger.infra.kafka;

import io.micrometer.core.instrument.MeterRegistry;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.common.header.Header;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.stereotype.Component;

import java.nio.charset.StandardCharsets;

@Component
class MonitorDeDeadLetter {

    private static final Logger log = LoggerFactory.getLogger(MonitorDeDeadLetter.class);

    private final MeterRegistry meterRegistry;

    MonitorDeDeadLetter(MeterRegistry meterRegistry) {
        this.meterRegistry = meterRegistry;
    }

    @KafkaListener(topics = Topicos.ANALISADA + "-dlt", groupId = "ledger-dlt-monitor")
    void registrar(ConsumerRecord<String, String> registro) {
        String causa = cabecalho(registro, KafkaHeaders.DLT_EXCEPTION_MESSAGE);
        String origem = cabecalho(registro, KafkaHeaders.DLT_ORIGINAL_TOPIC);

        meterRegistry.counter("ledger.dlt.mensagens", "topico_origem", origem).increment();
        log.error("Mensagem enviada ao DLT. chave={} origem={} causa={}", registro.key(), origem, causa);
    }

    private static String cabecalho(ConsumerRecord<String, String> registro, String nome) {
        Header header = registro.headers().lastHeader(nome);
        return header == null ? "desconhecido" : new String(header.value(), StandardCharsets.UTF_8);
    }
}

O ciclo completo tem uma última etapa operacional: alguém precisa investigar o DLT, corrigir a causa (um bug, um dado ruim) e republicar as mensagens no tópico original. Como o ledger é idempotente, reprocessar é seguro.

Serviço 4: notificacoes-service e retry não bloqueante

Notificar o cliente é um caso diferente do ledger: a ordem entre notificações de contas distintas não importa, e uma falha momentânea no provedor de e-mail não deveria parar o consumo das demais mensagens. Aqui o retry não bloqueante (tópicos de retry) é a escolha certa: a mensagem que falhou vai para um tópico de retry com atraso, e o consumidor segue em frente.

package com.exemplo.pagamentos.notificacoes;

import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.retrytopic.RetryTopicConfiguration;
import org.springframework.kafka.retrytopic.RetryTopicConfigurationBuilder;
import tools.jackson.core.JacksonException;

import java.util.List;

@Configuration
class RetryNotificacoesConfig {

    @Bean
    RetryTopicConfiguration retryNotificacoes(KafkaTemplate<String, String> template) {
        return RetryTopicConfigurationBuilder.newInstance()
                .exponentialBackoff(1_000, 2, 30_000)
                .maxAttempts(4)
                .includeTopics(List.of(Topicos.LIQUIDADA, Topicos.REJEITADA))
                .notRetryOn(JacksonException.class)
                .create(template);
    }
}
@Component
class NotificacaoListener {

    private final NotificacaoService notificacoes;
    private final JsonMapper jsonMapper;

    NotificacaoListener(NotificacaoService notificacoes, JsonMapper jsonMapper) {
        this.notificacoes = notificacoes;
        this.jsonMapper = jsonMapper;
    }

    @KafkaListener(topics = Topicos.LIQUIDADA, groupId = "notificacoes")
    void aoLiquidar(String json) {
        var evento = jsonMapper.readValue(json, TransferenciaLiquidada.class);
        notificacoes.notificarLiquidacao(evento.transferenciaId());
    }

    @KafkaListener(topics = Topicos.REJEITADA, groupId = "notificacoes")
    void aoRejeitar(String json) {
        var evento = jsonMapper.readValue(json, TransferenciaRejeitada.class);
        notificacoes.notificarRejeicao(evento.transferenciaId(), evento.motivo());
    }
}

O framework cria os tópicos de retry e o DLT automaticamente, com nomes derivados do tópico original. Se nenhum método de tratamento do DLT for configurado, um consumidor padrão apenas registra as mensagens em log; em um cenário real, associe um tratador que persista a falha ou dispare um alerta.

Repare que este é o fan-out em ação: liquidada e rejeitada são lidos por dois grupos diferentes (notificacoes e transferencias-projecao), cada um com seus próprios offsets, sem que o ledger saiba disso.

A idempotência não sumiu, só mudou de forma. Se o mesmo evento chegar duas vezes, o cliente receberia dois e-mails. O NotificacaoService deve registrar (transferenciaId, tipo) como chave única antes de enviar, com a mesma técnica do ledger.

Kafka Queues: quando os grupos de compartilhamento entram

O Kafka 4.2 declarou as share groups (KIP-932) prontas para produção, e o Spring Kafka 4.1 tem suporte a elas, incluindo modos de confirmação (ShareAckMode) e a renovação de lease (ShareAcknowledgment.renew()). Em vez de atribuir cada partição a um consumidor, uma share group permite que vários consumidores leiam da mesma partição, com confirmação individual de cada registro e contagem de tentativas de entrega.

Isso combina com cargas que se comportam como fila: processamento item a item, sem dependência de ordem, em que você quer escalar consumidores além do número de partições. As notificações deste artigo são um candidato natural. O ledger, por outro lado, continua fazendo sentido em consumer groups tradicionais, porque a ordem por chave e a semântica de log fazem parte do que queremos dele. É um recurso novo, então avalie com testes de carga antes de migrar cargas críticas.

Observabilidade

Em mensageria, o que você não mede vira incidente silencioso. Com o Actuator e o Micrometer, o Spring Boot já publica métricas do cliente Kafka. Habilitamos a observação (observation-enabled) nos templates e nos listeners, o que propaga contexto de tracing pelos headers das mensagens e permite seguir uma transferência de ponta a ponta entre os serviços.

Os sinais que eu monitoraria primeiro:

Um indicador simples para o backlog do outbox:

@Component
class MetricasOutbox {

    MetricasOutbox(MeterRegistry registry, JdbcClient jdbc) {
        Gauge.builder("outbox.pendentes", () ->
                        jdbc.sql("SELECT count(*) FROM outbox_evento WHERE publicado_em IS NULL")
                                .query(Long.class).single())
                .description("Eventos aguardando publicação no Kafka")
                .register(registry);
    }
}

Segurança e conformidade

Nada do que fizemos até aqui vale se o log de eventos for um ponto de vazamento. Alguns cuidados básicos para o ambiente de produção:

Testes

Testar comportamento distribuído com mocks engana. Para validar o que realmente importa neste desenho, o teste precisa de um Kafka e de um PostgreSQL de verdade, e o Testcontainers faz isso sem esforço. O Spring Boot 4 usa o Testcontainers 2.x, cujos módulos e pacotes mudaram em relação à série 1.x:

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-test</artifactId>
    <scope>test</scope>
</dependency>
<dependency>
    <groupId>org.testcontainers</groupId>
    <artifactId>testcontainers-junit-jupiter</artifactId>
    <scope>test</scope>
</dependency>
<dependency>
    <groupId>org.testcontainers</groupId>
    <artifactId>testcontainers-kafka</artifactId>
    <scope>test</scope>
</dependency>
<dependency>
    <groupId>org.testcontainers</groupId>
    <artifactId>testcontainers-postgresql</artifactId>
    <scope>test</scope>
</dependency>
<dependency>
    <groupId>org.awaitility</groupId>
    <artifactId>awaitility</artifactId>
    <scope>test</scope>
</dependency>

O teste mais valioso para o ledger é exatamente aquele que a arquitetura promete: a mesma transferência entregue duas vezes debita apenas uma vez.

package com.exemplo.pagamentos.ledger;

import com.exemplo.pagamentos.ledger.eventos.TransferenciaAnalisada;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.jdbc.core.simple.JdbcClient;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.test.context.DynamicPropertyRegistry;
import org.springframework.test.context.DynamicPropertySource;
import org.testcontainers.junit.jupiter.Container;
import org.testcontainers.junit.jupiter.Testcontainers;
import org.testcontainers.kafka.KafkaContainer;
import org.testcontainers.postgresql.PostgreSQLContainer;
import org.testcontainers.utility.DockerImageName;
import tools.jackson.databind.json.JsonMapper;

import java.time.Duration;
import java.time.Instant;
import java.util.UUID;

import static org.assertj.core.api.Assertions.assertThat;
import static org.awaitility.Awaitility.await;

@SpringBootTest
@Testcontainers
class LiquidacaoIdempotenciaIT {

    @Container
    static KafkaContainer kafka = new KafkaContainer(DockerImageName.parse("apache/kafka:4.3.1"));

    @Container
    static PostgreSQLContainer postgres =
            new PostgreSQLContainer(DockerImageName.parse("postgres:18-alpine"));

    @DynamicPropertySource
    static void propriedades(DynamicPropertyRegistry registry) {
        registry.add("spring.kafka.bootstrap-servers", kafka::getBootstrapServers);
        registry.add("spring.datasource.url", postgres::getJdbcUrl);
        registry.add("spring.datasource.username", postgres::getUsername);
        registry.add("spring.datasource.password", postgres::getPassword);
        registry.add("app.kafka.replicas", () -> "1");
    }

    @Autowired KafkaTemplate<String, String> kafkaTemplate;
    @Autowired JdbcClient jdbc;
    @Autowired JsonMapper jsonMapper;

    @Test
    void mesmaTransferenciaEntregueDuasVezesDebitaUmaUnicaVez() {
        jdbc.sql("INSERT INTO conta (id, saldo_centavos, moeda) VALUES ('CC-001', 100000, 'BRL')").update();
        jdbc.sql("INSERT INTO conta (id, saldo_centavos, moeda) VALUES ('CC-002', 0, 'BRL')").update();

        var evento = new TransferenciaAnalisada(
                UUID.randomUUID(), UUID.randomUUID(), "CC-001", "CC-002",
                30_000, "BRL", 10, Instant.now());
        String json = jsonMapper.writeValueAsString(evento);

        // Simula a entrega duplicada típica do at-least-once
        kafkaTemplate.send(Topicos.ANALISADA, "CC-001", json);
        kafkaTemplate.send(Topicos.ANALISADA, "CC-001", json);

        await().atMost(Duration.ofSeconds(30)).untilAsserted(() -> {
            long saldoOrigem = saldo("CC-001");
            long saldoDestino = saldo("CC-002");
            long lancamentos = jdbc.sql("SELECT count(*) FROM lancamento").query(Long.class).single();

            assertThat(saldoOrigem).isEqualTo(70_000);
            assertThat(saldoDestino).isEqualTo(30_000);
            assertThat(lancamentos).isEqualTo(2);
        });
    }

    private long saldo(String conta) {
        return jdbc.sql("SELECT saldo_centavos FROM conta WHERE id = :id")
                .param("id", conta)
                .query(Long.class)
                .single();
    }
}

O await garante que o teste espera o consumo assíncrono, e a propriedade app.kafka.replicas=1 adapta a declaração dos tópicos ao broker único do container. Um bom complemento é um teste com uma mensagem inválida, verificando que ela chega ao DLT sem travar a partição, e outro que simula uma falha transitória do banco.

Checklist das armadilhas mais comuns

Depois de ver o desenho completo, vale reunir os erros que mais aparecem em sistemas reais:

Quando Kafka não é a resposta

Kafka traz um custo operacional real: cluster, monitoramento, dimensionamento de partições, gestão de contratos e uma curva de aprendizado para o time. Ele brilha quando você precisa de um log durável, reprocessável e com alto volume, com vários consumidores independentes. Se o seu problema é distribuir tarefas entre workers com roteamento flexível e baixa complexidade operacional, ou se você já está inteiramente em uma nuvem com um serviço gerenciado de mensageria, outras ferramentas podem se encaixar melhor. É exatamente esse contraste que os próximos artigos da série exploram, ferramenta por ferramenta.

Conclusão

Repare no que sustentou o pipeline inteiro: nenhuma das garantias vem de uma única configuração mágica do Kafka. A durabilidade vem de acks=all com réplicas mínimas em sincronia. A ausência de perda na fronteira entre banco e broker vem do Transactional Outbox. A ausência de duplicidade de efeito vem de consumidores idempotentes com uma chave de negócio. A ordem por conta vem da chave da mensagem, e a correção dos saldos vem dos locks do banco. E a capacidade de operar o sistema vem de retries, DLTs e métricas pensados desde o início.

Kafka é uma ferramenta excelente para propagar fatos entre serviços, mas ele não retira de você a responsabilidade de decidir o que acontece quando uma mensagem chega duas vezes, chega fora de hora ou não chega. Se eu tivesse que resumir este artigo em uma pergunta para levar ao seu próprio projeto, seria esta: se qualquer mensagem do seu sistema for entregue duas vezes amanhã, o que exatamente acontece com o dinheiro? Se a resposta depender de sorte, o desenho ainda não está pronto.

Referências oficiais

As explicações de infraestrutura e configuração deste artigo foram conferidas na documentação oficial das versões usadas: