Kafka por sessão: preservar ordem não garante executar cada efeito uma única vez
O pipeline em Go combina serialização por sessão e watermark contíguo, mas replay, idempotência e timeouts continuam sendo decisões explícitas.
Um artigo da InfoQ descreve um pipeline Kafka em Go que preserva a sequência de mensagens por sessão em uma plataforma de IA conversacional.[13] O autor relata mais de 40 milhões de mensagens em produção sem violações de ordem observadas, distinguindo essa observação de uma verificação exaustiva de cada sequência.[13] O desenho é interessante justamente por explicitar o custo da recuperação, não por prometer exactly-once.
O que muda para quem desenvolve
Kafka garante ordem dentro da partição, mas isso não resolve sozinho o paralelismo entre sessões que compartilham partições.[13] No desenho apresentado, workers distribuem mensagens por hashing consistente do identificador de sessão, e uma goroutine criada sob demanda processa cada sessão em sequência.[13] Retries transitórios permanecem nessa goroutine, enquanto mensagens posteriores da mesma sessão aguardam.[13]
O commit usa um watermark contíguo por partição: concluir um offset alto não autoriza ultrapassar um offset anterior ainda em execução.[13] Depois de uma falha, mensagens já executadas podem ser repetidas; os efeitos externos dependem da idempotência dos serviços chamados.[13]
Minha leitura é que há três contratos a escrever separadamente: ordem lógica, recuperação do consumo e repetição de efeitos. Uma conversa pode manter a ordem e ainda gerar uma cobrança duplicada se a aplicação confundir avanço do offset com execução única.
O artigo admite forçar a conclusão de lacunas ou offsets travados após timeout, com registro na DLQ, para devolver progresso à partição.[13] Essa escolha pode sacrificar a completude da mensagem.[13] Deve aparecer como exceção operacional explícita, não como sucesso normal.
Como aplicar
Proponha um teste de falhas com sessões concorrentes e efeitos externos simulados:
- Registre uma sequência esperada por sessão e um identificador estável por mensagem.
- Faça uma mensagem falhar temporariamente e confira que as posteriores da mesma sessão aguardam.
- Permita que outra sessão avance e observe se o watermark continua respeitando a lacuna.
- Interrompa o consumidor depois do efeito externo e antes do commit.
- Reinicie e confira tanto o replay quanto a proteção contra repetir o efeito.
Acrescente rebalance, shutdown e saturação de buffers. A proposta é observar o contrato sob interrupção, não apenas medir throughput em condições ideais. Para timeouts forçados, defina quem examina a DLQ e como a operação é compensada ou retomada.
O artigo descreve pausa de partição e retenção dos registros já recebidos para evitar ultrapassagem durante backpressure.[13] Essa pausa pode atrasar sessões independentes na mesma partição.[13] Meça esse custo antes de escolher o nível de isolamento.
Cuidados e limites
O ensaio sintético informou 14.027 mensagens por segundo, com p99 de 185 ms; o alvo de latência da camada de mensagens excluía a inferência variável do LLM.[13] O teste contra o pipeline implantado registrou p99 de 2.843,1 ms, atribuído pelo autor ao caminho WebSocket e à quota do API Gateway.[13]
Os testes de alta carga não validaram independentemente a ordem por sessão.[13] Portanto, não use throughput como prova de ordenação. Antes de adotar o desenho, decida se o custo de Kafka e dessa coordenação é proporcional ao requisito real.
Fonte
[13] Fonte: InfoQ. Data/hora no feed: 07/10/2026 às 06:00:00.