Programação e desenvolvimento de software

Como criar um pipeline Kafka que preserve a ordem das sessões usando Go

Joshua Oluikpe explica um projeto prático para processar mensagens do Kafka em ordem estrita dentro de cada sessão, mantendo sessões independentes capazes de operar em paralelo. O projeto utiliza hashing consistente, um worker Go por sessão, novas tentativas dentro do mesmo fluxo e gerenciamento de pontos de commit contíguos para evitar a perda de mensagens após uma falha.

2026-10-07
6 min de leitura
1 visualizações
certi.news Editorial Team
Como criar um pipeline Kafka que preserve a ordem das sessões usando Go

Não basta depender da ordem do Kafka dentro da partição quando milhares de sessões compartilham a mesma partição. Em sistemas de conversação e tarefas sequenciais, uma mensagem de correção como «defina o destino como Paris» pode chegar depois de uma mensagem «reserve uma viagem para Londres», mas executá-la antes da primeira corrompe o contexto. A solução descrita por Joshua Oluikpe na InfoQ transfere parte da responsabilidade pela ordenação para a camada de aplicação escrita em Go.

Ordenação dentro da sessão e paralelismo entre sessões

O projeto se baseia em dois níveis de workers. O nível de distribuição recebe os registros do Kafka e os encaminha de acordo com o identificador da sessão usando hashing consistente, de modo que as mensagens da mesma sessão cheguem ao mesmo worker, enquanto sessões diferentes possam ser distribuídas entre vários workers. Em seguida, cada sessão ativa recebe uma goroutine dedicada, que processa uma mensagem por vez e na ordem de chegada.

Assim, uma sessão lenta não se torna um obstáculo para outras sessões, como ocorre em um pool de workers plano quando a nova tentativa de uma única sessão ocupa o worker por completo. As goroutines são criadas conforme necessário e removidas após um período de inatividade, fazendo com que o consumo de recursos esteja relacionado ao número de sessões ativas, e não ao número total de sessões possíveis.

Novas tentativas sem ultrapassar mensagens

As novas tentativas ocorrem dentro da própria goroutine da sessão. Quando ocorre um erro temporário, como o esgotamento do tempo limite de um serviço dependente, o worker aguarda usando recuo exponencial com uma dose de aleatoriedade e então reprocessa a mesma mensagem antes de avançar para as mensagens seguintes. Como as mensagens seguintes permanecem no canal atrás da mensagem com falha, a quarta mensagem não pode ultrapassar a terceira.

Já os erros irrecuperáveis, como dados inválidos, são enviados diretamente para a lista de mensagens mortas (DLQ). Essa separação entre erros temporários e finais elimina a necessidade de uma máquina de estados separada para gerenciar novas tentativas, mas também significa que as mensagens de determinada sessão podem se acumular durante o período de recuo.

Commit seguro e recuperação após falhas

Quando várias sessões processam suas mensagens em paralelo, não é permitido fazer commit apenas do maior offset concluído. Se o offset 104 for concluído antes do 102 e o commit for feito até o 104, uma falha do consumidor poderá fazer com que o 102 seja definitivamente ultrapassado. Por isso, o sistema acompanha, para cada partição, os offsets em processamento e os que foram concluídos, e só move o ponto de commit até o maior offset concluído de forma contígua.

Como consequência, a reinicialização pode começar em um offset mais antigo e reprocessar algumas mensagens. O projeto utiliza um identificador de evento estável para ajudar a eliminar duplicações, mas não afirma alcançar uma execução «exatamente uma vez»; os efeitos externos continuam dependendo da capacidade dos serviços dependentes de lidar com duplicações de forma idempotente. Em casos raros, uma lacuna persistente pode ser ultrapassada após o fim de um período determinado para recuperar o progresso da partição, com o registro dessa situação como uma falha operacional e o encaminhamento da mensagem para a DLQ, uma compensação explícita entre continuidade e completude do processamento.

O que é necessário para a produção real?

As garantias de ordenação não são concluídas apenas com a lógica básica. Durante a redistribuição das partições, o consumidor interrompe o recebimento de novos registros, aguarda a conclusão do trabalho em andamento e então faz commit do último watermark contíguo antes de transferir a propriedade. Quando os canais dos workers ficam cheios, ele interrompe temporariamente a busca de mensagens em vez de descartá-las, armazenando os registros previamente obtidos em buffers específicos de cada partição para que mensagens mais recentes não avancem sobre eles.

O sistema também precisa detectar individualmente offsets paralisados, conectar o rastreamento através dos limites assíncronos do Kafka, manter os sinais de atraso necessários para a expansão automática durante a interrupção da busca e escoar os trabalhos e gravar a DLQ antes do encerramento. O produto também aplica hashing consistente no lado do envio para garantir a ordem das mensagens durante a produção assíncrona.

Resultados e limitações

Em um teste sintético que envolveu 50.000 mensagens e 10 sessões, com a injeção de erros temporários em 10% das mensagens, o sistema registrou 14.027 mensagens por segundo sem violações de ordenação, com latência mediana de 11,7 milissegundos. Em um teste de carga com 1.000 sessões, a taxa efetiva de envio chegou a 48.805 mensagens por segundo, de uma meta de 50.000, sem erros de envio, mas a latência de cauda ultrapassou o orçamento de 200 milissegundos devido à cota de taxa de mensagens da API Gateway, e não ao mecanismo de ordenação das sessões.

O autor afirma que o sistema processou mais de 40 milhões de mensagens de produção sem violações de ordenação observáveis, encaminhando apenas 460 mensagens para a DLQ, cerca de 0,002%. No entanto, esse número se baseia no monitoramento operacional, e não em uma verificação abrangente da sequência de cada mensagem; além disso, o teste de carga mais elevado não verificou de forma independente a ordenação de cada sessão. Portanto, o valor prático do projeto está em transformar a ordenação de uma suposição no nível do Kafka em uma propriedade imposta pela aplicação, mantendo a necessidade de testes explícitos de sequência e de garantias de idempotência nos serviços dependentes.

Fonte da notícia
InfoQ - Architecture Articles
Abrir fonte original ↗
c
Autor

certi.news Editorial Team

Na mesma categoria

Você também pode gostar

Ver todas as notícias