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:
| Conceito | O que é | Pense nisso como |
|---|---|---|
| Producer | Aplicação que publica registros | Quem escreve um fato |
| Broker | Servidor Kafka que armazena e serve os registros | Um nó do cluster |
| Topic | Nome lógico do fluxo de eventos | Uma categoria de eventos |
| Partition | Divisão de um topic em sequências ordenadas | Uma fila paralela |
| Offset | Posição de um registro dentro da partição | Número da linha no log |
| Consumer | Aplicação que lê registros | Quem reage ao fato |
| Consumer group | Conjunto de consumidores que divide as partições | Um time processando a mesma fila |
| Replica | Cópia de uma partição em outro broker | Backup operacional |
| ISR | Réplicas que estão em sincronia com o líder | Có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:
acks | O que acontece | Consequência |
|---|---|---|
0 | O producer não espera confirmação do broker | Menor latência, mas sem confirmação de recebimento |
1 | O líder confirma depois de gravar localmente | O líder pode falhar antes da replicação e perder o registro |
all | O líder espera as réplicas em sincronia confirmarem | Maior 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?
- At-most-once: o offset é confirmado antes do processamento. Se a aplicação cair, a mensagem pode ser perdida.
- At-least-once: o offset só é confirmado depois do processamento. Se a aplicação cair antes do commit, a mesma mensagem pode voltar. Duplicata faz parte do modelo.
- Exactly-once, dentro do Kafka: com producer idempotente, transações e consumidores
read_committed, um fluxo Kafka → Kafka pode evitar efeitos duplicados e tornar a operação transacional dentro do próprio Kafka.
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.
| Componente | Versão | Observação |
|---|---|---|
| Java | 25 (LTS) | Records, text blocks, sealed interfaces e pattern matching em switch |
| Spring Boot | 4.1.1 | Baseado no Spring Framework 7 |
| Spring for Apache Kafka | 4.1.1 | Cliente Kafka 4.2.1; Jackson 3; sem dependência do Spring Retry |
| Apache Kafka (broker) | 4.3.1 | Somente KRaft |
| PostgreSQL | 18 | Um banco por serviço |
| Jackson | 3 | JsonMapper, exceções não checadas |
Algumas mudanças dessas versões pesam no dia a dia e aparecem ao longo do artigo:
- No Spring Boot 4, o starter do Kafka precisa ser declarado explicitamente:
spring-boot-starter-kafka. - Com o Jackson 3, os serializadores do Spring Kafka ganharam novos nomes:
JacksonJsonSerializer,JacksonJsonDeserializereJsonKafkaHeaderMappersubstituemJsonSerializer,JsonDeserializereDefaultKafkaHeaderMapper. As classes antigas (Jackson 2) seguem funcionando, mas estão deprecadas. - O Spring Kafka 4 removeu a dependência do Spring Retry e passou a usar o suporte a retry do Spring Framework 7. É uma quebra de compatibilidade para quem tinha configuração baseada em
RetryTemplateantigo. - O Kafka 4.0 introduziu o novo protocolo de rebalanceamento de grupos (KIP-848), e o 4.3 já registra em log uma recomendação para abandonar o protocolo classic nos consumidores.
- O Kafka 4.2 tornou as share groups (KIP-932, as “Kafka Queues”) prontas para produção. Voltamos a elas mais adiante.
- O Spring Kafka 4.1.0 corrigiu três CVEs (CVE-2026-41726, CVE-2026-41727 e CVE-2026-41731), envolvendo cache ilimitado por header de seletor, headers de retry forjados e correspondência ampla de pacotes confiáveis. Em um sistema financeiro, manter essa dependência atualizada não é opcional.
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:
- Um dono para cada dado. O
ledger-serviceé o único que altera saldos. Kafka carrega fatos sobre o que aconteceu, não é uma API para alterar o estado de outro serviço. - Nomes de tópico estáveis e versionados:
<domínio>.<entidade>.<evento>.v<versão>. A versão no nome permite evoluir contratos sem quebrar consumidores. - Chave = conta de origem, para que todas as transferências que debitam a mesma conta sejam lidas em ordem.
- Cada serviço tem seu banco. Não existe transação distribuída entre eles, e é por isso que o padrão Transactional Outbox é o coração do exemplo.
Os tópicos e suas características:
| Tópico | Produtor | Consumidores (grupo) | Partições |
|---|---|---|---|
pagamentos.transferencia.solicitada.v1 | transferencias | antifraude | 12 |
pagamentos.transferencia.analisada.v1 | antifraude | ledger | 12 |
pagamentos.transferencia.liquidada.v1 | ledger | notificacoes, transferencias-projecao | 12 |
pagamentos.transferencia.rejeitada.v1 | antifraude e ledger | notificacoes, transferencias-projecao | 12 |
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ção | Para que serve | Neste artigo |
|---|---|---|
acks | Quanto de confirmação o producer exige | all |
enable.idempotence | Evita duplicatas causadas por retries do próprio producer | true |
delivery.timeout.ms | Tempo máximo para o envio completar, incluindo retries do cliente | 120000 |
linger.ms | Dá alguns milissegundos para formar lotes maiores | 5 |
compression.type | Comprime os lotes enviados | zstd |
key.serializer | Converte a chave Java para bytes | StringSerializer |
value.serializer | Converte o valor Java para bytes | StringSerializer |
transaction-id-prefix | Ativa produtores transacionais do Kafka | Usado 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ção | Para que serve | Neste artigo |
|---|---|---|
group-id | Identifica o consumer group | ledger, antifraude, notificacoes |
enable-auto-commit | Define se o cliente confirma offsets automaticamente | false |
auto-offset-reset | Define de onde começar quando não há offset válido | earliest |
isolation.level | Controla a leitura de registros transacionais | read_committed |
max-poll-records | Quantidade máxima de registros devolvidos por poll | 50 no ledger |
max.poll.interval.ms | Tempo máximo esperado entre polls | Deve comportar o tempo de processamento |
concurrency | Quantidade de consumidores/threads gerados pelo container Spring | 3 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:
- Valores monetários em centavos (
long), nuncadouble. Ponto flutuante não representa dinheiro com exatidão. Se o domínio exigir mais casas decimais, useBigDecimalcom escala fixa ou uma unidade menor. - Cada serviço mantém sua própria cópia dos records de que precisa, em vez de compartilhar uma biblioteca de classes. Compartilhar classes acopla o ciclo de deploy dos serviços. Em produção, o contrato costuma viver em um esquema formal (JSON Schema, Avro ou Protobuf com um Schema Registry), com regras de compatibilidade verificadas no pipeline.
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:
- Publicar antes do commit: se o commit falhar, existe um evento sobre uma transferência que nunca existiu.
- Publicar depois do commit: se o processo cair entre um passo e outro, a transferência existe, mas ninguém foi avisado.
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:
- Fecha a principal janela de perda entre banco e broker. A linha só é marcada como publicada depois que o Kafka confirmou o envio. Se o processo cair no meio, o lote volta a ficar pendente e pode ser publicado novamente.
- Não garante que nenhum evento seja publicado duas vezes. Se o processo cair depois do envio e antes do
UPDATE, o mesmo evento será publicado de novo. Esse é o at-least-once em ação, e é por isso que os consumidores precisam ser idempotentes. - Ordem por chave só é estrita com uma instância do relay (ou com o relay particionado por chave). Com
SKIP LOCKEDe várias instâncias, dois lotes concorrentes podem publicar eventos da mesma conta fora de ordem. Se a ordem rígida importa mais do que o paralelismo de publicação, rode um relay único, ou considere captura de mudanças (CDC) com Debezium lendo o log de transações do banco.
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:
acks=allfaz o produtor esperar a confirmação do conjunto de réplicas em sincronia da partição. Combinado commin.insync.replicas=2em um tópico com fator de replicação 3, evita que a escrita seja aceita quando o ISR caiu abaixo do mínimo configurado.enable.idempotence=truefaz o producer evitar duplicatas causadas por retransmissões do próprio cliente. Em versões modernas do cliente, a idempotência já é habilitada por padrão quando não há configurações conflitantes; deixá-la explícita aqui documenta a intenção.delivery.timeout.mslimita o tempo total que uma mensagem pode ficar tentando ser entregue (incluindo retries). Passado esse tempo, osendfalha, e o relay não marca o evento como publicado.linger.msecompression-typesão ajustes de vazão: agrupam registros em lotes e comprimem, à custa de alguns milissegundos de latência.
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:
- Produtor idempotente: elimina duplicatas causadas por retentativas do cliente.
- 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.
- 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:
- Só vale dentro do Kafka. Se o listener também gravasse em um banco de dados, essa escrita não participaria da transação Kafka. Para essa fronteira você continua precisando de outbox e idempotência.
- Retry não bloqueante (tópicos de retry) não combina com transações no container. Quando o listener lança exceção, o commit do container é feito e o registro é encaminhado ao tópico de retry, conforme a documentação do Spring Kafka.
- O tratamento de falhas muda. Com transações, o fluxo passa pelo
AfterRollbackProcessor, e o comportamento padrão também descarta o registro após várias tentativas. Configure um recuperador com Dead Letter Topic ali também, seguindo a seção After-rollback Processor da documentação.
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:
- Insert primeiro, com
ON CONFLICT DO NOTHING. Em PostgreSQL, uma segunda transação que tente inserir a mesma chave espera a primeira terminar. Se a primeira confirmar, a segunda não insere nada e reconhece a duplicata; se a primeira desfizer, a segunda insere normalmente. Isso cobre entregas duplicadas sequenciais e concorrentes. - Contas travadas em ordem (
ORDER BY id FOR UPDATE). Se a transferência A→B e a transferência B→A rodassem ao mesmo tempo travando as contas em ordens opostas, haveria deadlock. Uma ordem global elimina esse cenário. - Rejeição não é exceção. Saldo insuficiente é um resultado de negócio legítimo: a transação é confirmada, registra-se a rejeição e o evento correspondente vai para o outbox. Lançar uma exceção aqui só faria o Kafka reentregar a mesma mensagem e chegar ao mesmo resultado.
- O evento de saída e os lançamentos vivem na mesma transação. Ou tudo acontece, ou nada acontece.
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:
enable-auto-commit: falsecomack-mode: record: o Spring Kafka confirma o offset depois que o listener retorna com sucesso, registro a registro. Nunca deixe o cliente confirmar offsets sozinho em intervalos de tempo, porque isso pode confirmar mensagens que ainda não foram processadas.isolation-level: read_committed: o ledger lêanalisada, escrito de forma transacional pelo antifraude. Sem isso, ele poderia ver registros de transações abortadas.auto-offset-reset: earliest: quando um grupo novo aparece (ou perde seus offsets), ele começa do início do log em vez de ignorar tudo o que já foi publicado. Comlatest, um grupo novo em um serviço financeiro poderia pular transferências.group.protocol: consumer: usa o novo protocolo de rebalanceamento (KIP-848), com atribuição de partições conduzida pelo broker, sem barreira global de sincronização. Ao usar esse protocolo, vários ajustes de sessão e de atribuição passam a ser controlados no lado do broker, então não os configure no cliente.concurrency: 3cria três threads consumidoras dentro da instância. O paralelismo útil total, somando todas as instâncias, nunca passa do número de partições (12 aqui).max-poll-recordse o tempo que o processamento leva por registro precisam caber dentro domax.poll.interval.ms; do contrário, o consumidor é considerado morto e a partição é reatribuída.
Falhas: retry bloqueante e Dead Letter Topic
Existem três categorias de falha, e cada uma pede uma resposta diferente:
| Tipo | Exemplo | Resposta |
|---|---|---|
| Transitória | Banco indisponível por alguns segundos | Retentar com espera crescente |
| Permanente (mensagem venenosa) | JSON inválido, campo obrigatório ausente | Ir direto para o DLT, sem retentar |
| Resultado de negócio | Saldo insuficiente | Nã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:
- Espera exponencial com teto: 0,5 s, 1 s, 2 s, 4 s, 8 s (limitada a 15 s), e só 5 retentativas. Retry infinito em um serviço financeiro esconde problemas e trava partições.
- Exceções não retentáveis: um
JacksonExceptionsignifica que a mensagem é ilegível, e retentar nunca vai consertar isso. Ela vai direto para o DLT. - Mesma partição no DLT: o
DeadLetterPublishingRecovererusa a mesma partição do registro original, então o tópico-dltprecisa ter pelo menos o mesmo número de partições do original.
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:
- Lag por grupo e partição: a distância entre o último offset publicado e o último confirmado. Lag crescente significa consumidor lento ou parado.
- Mensagens no DLT: qualquer valor acima de zero merece alerta.
- Backlog do outbox: quantidade de eventos pendentes e idade do mais antigo. Se o relay parar, nada é publicado, e ninguém percebe pelo lado do Kafka.
- Taxa de rejeições por motivo no ledger: uma alta súbita de
SALDO_INSUFICIENTEouCONTA_INEXISTENTEcostuma indicar um problema a montante.
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:
- Criptografia em trânsito: TLS entre clientes e brokers (e entre brokers), com autenticação mútua ou SASL.
- Autorização mínima: cada serviço usa uma identidade própria com ACLs específicas. O antifraude lê
solicitadae escreve emanalisadaerejeitada; ele não tem motivo algum para lerliquidada. - Sem dados pessoais nos eventos. Usar identificadores de conta e de transferência, em vez de nomes, documentos ou e-mails, reduz muito a superfície de exposição. Um log imutável e replicado é um lugar ruim para dados que talvez precisem ser apagados. Envolva o time jurídico e de compliance para definir retenção e requisitos regulatórios aplicáveis, como a LGPD. Este artigo não é aconselhamento jurídico.
- Criptografia em repouso nos discos dos brokers e controle de acesso ao DLT, já que ele contém mensagens que falharam, muitas vezes com o payload original completo.
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:
- Confiar no tratador de erros padrão, que descarta a mensagem depois das tentativas, sem DLT.
- Deixar o auto-commit de offsets ligado em consumidores que fazem efeitos colaterais.
- Publicar no Kafka e gravar no banco em duas operações separadas (dual write) sem outbox.
- Escolher uma chave de idempotência técnica (offset, UUID aleatório) em vez de uma identidade de negócio.
- Usar
auto.offset.reset=latestem grupos novos de consumidores que não podem perder nada. - Aumentar o número de partições sem considerar que o mapeamento chave → partição muda.
- Retry infinito ou sem espera, que trava a partição e mascara o problema.
- Ignorar o DLT: sem alerta e sem um procedimento de reprocessamento, ele vira um cemitério de mensagens.
- Evoluir um contrato de evento de forma incompatível sem versionar o tópico (por padrão o Spring Boot configura o mapeador para ignorar campos desconhecidos, o que ajuda em mudanças compatíveis, mas convém fixar esse comportamento com um teste de contrato).
- Supor que transações Kafka resolvem consistência com um banco de dados externo.
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:
- Apache Kafka 4.3 — Quickstart
- Apache Kafka 4.3 — Docker
- Apache Kafka 4.3 — configurações do producer
- Apache Kafka 4.3 — configurações de broker e tópicos
- Apache Kafka 4.3 — operações básicas
- Spring Boot — Apache Kafka Support
- Spring for Apache Kafka 4.1 — documentação
- Spring for Apache Kafka 4.1 — retries não bloqueantes