- Na verificação do Bufstream 0.1.0~0.1.3, um sistema de streaming compatível com Kafka, foram encontrados 2 problemas de disponibilidade no próprio Bufstream e 3 problemas de segurança; todos os 5 foram corrigidos na versão 0.1.3
- Os testes foram baseados no Java Kafka Client 3.8.0 e nos testes Jepsen já existentes para Kafka/Redpanda, usando configurações com prioridade para segurança como
acks = all,enable.idempotence = true,enable.auto.commit = falseeread_committed - Os problemas do Bufstream incluem paralisação de consumidor e produtor, resposta incorreta com offset
0, perda de commit de transação e perda de escrita reconhecida causada por um bug de filtragem do tamanho de resposta da API fetch - Durante a investigação, também vieram à tona problemas no Kafka Java client e no protocolo de transações do Kafka, como bloqueio indefinido em
Consumer.close(), offset de consumer imprevisível e problemas de aborted read·lost write·torn transaction - O Jepsen considera que, como o protocolo de transações do Kafka não garante explicitamente a ordem das requisições do cliente nem o número da transação, a segurança transacional do Kafka e de sistemas compatíveis com Kafka pode ser quebrada ao usar o cliente Java oficial
Estrutura do Bufstream e escopo da verificação
- O Kafka é um sistema de streaming que fornece logs append-only com replicação e sharding, e o Bufstream é uma implementação alternativa ao Kafka voltada a governança de dados e eficiência de custos em ambientes de nuvem
- Assim como o Kafka, o Bufstream oferece topics e partitions e funciona com clientes Kafka padrão
- o producer faz append de records com
producer.send() - o consumer se vincula a uma partition com
consumer.assign()ouconsumer.subscribe()e depois lê records comconsumer.poll() - um consumer group divide entre si o processamento dos records de um conjunto de topics
- o producer faz append de records com
- Quando integrado ao Buf Schema Registry, é possível inspecionar records em Protocol Buffer para oferecer validação de records, controle de acesso em nível de campo e conversão de formato de dados para outros sistemas
- Diferentemente do Kafka, que usa disco local e seu próprio protocolo de replicação, o Bufstream grava os dados diretamente em object storage
- a ideia é reduzir custos aproveitando a estrutura de custos do tráfego de replicação do object storage
- os nós do Bufstream podem operar como VMs stateless com auto scaling
- O Bufstream é composto por três subsistemas
- agent: serviço stateless que fornece a API Kafka
- object store: armazena chunks de records e os fornece aos leitores
- coordination service: atualmente usa etcd e define quais chunks foram commitados e a ordem dos records
- Em outubro de 2024, o Bufstream havia sido distribuído apenas para alguns clientes, e a documentação o apresentava como um “drop-in replacement do Apache Kafka” e destacava compatibilidade com transações Kafka e exactly-once semantics, mas sem muitas afirmações concretas sobre segurança
Configuração do cliente e premissas transacionais
- Como em testes anteriores de sistemas compatíveis com Kafka, o Jepsen ajustou a configuração dos clientes para obter um comportamento mais seguro
-
Configuração do producer
- usa o padrão
acks = all - no Bufstream,
acks = 0pode reconhecer uma escrita sem esperar pelo storage, o que pode levar à perda de escritas commitadas - tanto
acks = 1quantoacks = allbloqueiam até que o Bufstream tenha certeza da persistência durável - para evitar append duplicado nas tentativas automáticas do producer do Kafka, foi usado o valor padrão
enable.idempotence = true
- usa o padrão
-
Configuração do consumer
- como há documentos alertando que auto-commit pode levar à perda de dados, em geral foi usado
enable.auto.commit = false - quando não há offset commitado, o padrão
auto.offset.resetcomeça do offset mais recente, o que não garante entrega at-least-once - para que o consumer possa observar todo o log, foi usado
auto.offset.reset = earliest - uma transação no Kafka é composta pelo conjunto de records enviados pelo producer e pelo mapa dos offsets máximos por partition obtidos pelo consumer via poll
- somente quando a transação é commitada os records enviados se tornam duráveis e acabam visíveis para consumidores
read_committed, e o offset commitado também avança até pelo menos o offset especificado na transação - se a transação não for commitada, o offset commitado não avança, e a visibilidade da escrita pode variar conforme a configuração do consumer
- quando um consumer
read_uncommittedlê valores de uma transação abortada, isso é classificado como aborted read (G1a) - a documentação do Kafka afirma que
read_committedimpede G1a e, em certa medida, garante a propriedade de que todas as escritas de uma transação aparecem ou nenhuma aparece, mas nos testes do Jepsen com Kafka, Redpanda e Bufstream foram observados ciclos de escrita (fenômeno semelhante a G0) e algumas formas de G1c
- como há documentos alertando que auto-commit pode levar à perda de dados, em geral foi usado
Projeto de teste
- O Jepsen testou o Bufstream 0.1.0 até 0.1.3, além de vários builds release candidate
- O harness de teste usa o Bufstream test harness, a biblioteca de testes Jepsen e o Java Kafka Client 3.8.0
-
Ambiente de execução
- foram usados de 3 a 5 nós Debian Bookworm tanto em containers LXC quanto em VMs EC2
- 1 nó foi usado para o etcd, 1 nó para o Minio e o restante como agents do Bufstream
- producer, consumer e admin client foram inicializados com apenas um único nó em
bootstrap_servers, mas a descoberta inteligente de clientes não foi bloqueada
-
Principais configurações de segurança
- auto-commit false
acks = all- retries 1.000
- idempotence enabled
- nível de isolamento
read_committed auto_offset_reset = earliest- criação automática de topics no lado do servidor desativada
- a injeção de falhas incluiu pausa de processo (
SIGSTOP), crash (SIGKILL), desvio de relógio (clock_settime) e partição de rede (iptables) - como o Bufstream é dividido em agent, object store e coordination service, foi criada uma nova ferramenta do Jepsen para injetar falhas mirando apenas subsistemas específicos
- por exemplo, alternando ao longo do tempo entre derrubar apenas nós do Bufstream ou pausar apenas o coordenador etcd
Workload de fila e workload de abort
- O workload de fila analisa a segurança de acordo com o modelo de dados do Kafka
- Cada processo lógico executa producer, consumer e admin client
- A chave numérica identifica uma topic-partition específica
- As chaves são escolhidas com frequência exponencial, de modo que algumas são acessadas com frequência e outras raramente
- São usadas três operações básicas
crash: encerra um processo lógico e o substitui por um novo clientsubscribeouassign: altera o conjunto de topics ou partitions que o consumer irá consultar compolltxn,poll,send: executa uma sequência de micro-operaçõespollousend
- No workload não transacional, cada
sendoupollinclui exatamente uma micro-operação - No workload transacional, várias micro-operações são encapsuladas em uma transação Kafka
- A análise constrói um mapeamento de offset para valor por chave e então procura erros
- Se vários valores aparecem no mesmo offset, há offset inconsistente
- Se o mesmo valor aparece em vários offsets, há erro de duplicação
- Se um registro confirmado nunca é observado, ele é classificado como perdido ou não visto
- Se um
pollretorna um valor enviado por uma operação abortada, há leitura abortada - Também verifica se a transação observa suas próprias gravações
- Após o teste principal, as falhas são resolvidas e começa a etapa de leituras finais
- Cada processo lê todas as topic-partitions a partir do offset 0 e faz
pollaté o maior offset gravado conhecido - Se as leituras finais expiram e um registro confirmado ainda não é observado, ele é classificado como não visto
- Cada processo lê todas as topic-partitions a partir do offset 0 e faz
- O workload de abort foi adicionado para rastrear o comportamento do offset de
pollapós o aborto de transações- O topic é limitado a uma única partition, processo, producer e consumer
- Depois que a transação faz
pollde um registro, ela é abortada intencionalmente e, em seguida, o offset depollé classificado como advance, rewind, rewind-further ou other
5 problemas encontrados no Bufstream
-
Consumidores travados (#1)
- Da versão 0.1.0 até a 0.1.3-rc.8, a etapa final de leitura frequentemente travava
consumer.poll()retornava imediatamente um resultado vazio, mas ainda restavam milhares de records reconhecidos nos logs- Esse estado persistia de dezenas de segundos até mais de 1 hora
- Em um teste, 691 records reconhecidos foram enviados nos primeiros 120 segundos, e no início das leituras finais 40 deles não haviam sido observados por nenhum poller
- Depois disso,
consumer.poll()não retornou resultados por mais de 1 hora, e o teste expirou por timeout - A causa foi que um nó do Bufstream reiniciado podia retornar um valor em cache desatualizado para o last stable offset e o high watermark
- Algumas bibliotecas cliente concluíam que não havia records mais à frente e entravam em stall, e o Bufstream aplicou um patch na 0.1.3-rc.6 para atualizar o cache na inicialização
-
Produtores e consumidores travados (#2)
- Mesmo na 0.1.3-rc.6, continuaram sendo observados problemas de escrita não vista após pause, crash e partition no coordinator, storage e nó do Bufstream
- Em alguns casos, após uma pausa do coordinator, o cliente entrava em estado de timeout aguardando
InitProducerId, mesmo com todos os nós do Bufstream em execução - Em outros casos,
listOffsetsfalhava comnode ... being disconnectedoutimed out waiting for a node assignment, epollconcluía sem retornar resultados - Matar e reiniciar o nó do Bufstream resolvia o problema
- A causa estava relacionada ao lease do etcd
- O agente do Bufstream usa leases do etcd para rastrear agentes ativos
- Por causa de pauses ou partitions curtas, o etcd apagava chaves vinculadas ao lease do agente, mas a atualização de exclusão podia não ser entregue ao agente
- O agente ficava sem saber que havia perdido o próprio lease
- A equipe do Bufstream adicionou lógica extra de polling, e na 0.1.3-rc.8 as escritas não vistas foram em grande parte resolvidas
-
Offsets zero espúrios (#3)
- Da 0.1.0 até a 0.1.3-rc.2, um valor enviado podia receber o offset
0e depois aparecer em um offset real mais alto - Isso ocorria mesmo quando o offset
0já havia sido atribuído muito antes - Apenas o remetente observava o offset 0; o poller observava um offset mais alto
- Em um teste de 2 minutos com um único nó do Bufstream e pause do processo etcd, 6 escritas receberam offset
0e depois apareceram em offsets mais altos - A causa foi a ausência de um campo necessário na resposta de erro do Bufstream
- O Bufstream enviava ao etcd uma solicitação para commit do log, e o etcd a processava, mas por causa de pause ou partition o Bufstream podia expirar por timeout esperando a resposta
- O Bufstream enviava ao cliente um código de erro, mas não definia o offset do record enviado como
-1, que é o sinal de erro - O cliente Java do Kafka interpretava isso como uma resposta de sucesso com offset
0 - O Franz-go, usado pela suíte de testes do Bufstream, interpretava essa mensagem como erro, então esse problema não apareceu nos testes
- O Bufstream corrigiu isso na 0.1.3-rc.6, e o Jepsen não conseguiu reproduzi-lo novamente depois disso
- Da 0.1.0 até a 0.1.3-rc.2, um valor enviado podia receber o offset
-
Escritas de transação perdidas (#4)
- Na 0.1.2, perdas de escrita em que alguns records de transações commitadas desapareciam e não eram mais observados ocorriam com frequência
- Em um teste, 240 records escritos por transações commitadas foram perdidos ao longo de 100 segundos e 6.761 transações de escrita
- No exemplo, o valor
141da chave5foi retornado como escrito com sucesso no offset274, mas todoconsumer.poll()pulava esse offset - A causa foi um bug no mecanismo de segurança de concorrência adicionado na 0.1.2
- Esse mecanismo atribui um número único a cada transação dentro de um epoch do producer para mitigar a falta de idempotência do protocolo de transações do Kafka
- Por causa de um bug na lógica de rastreamento do número de transação, quando várias transações eram commitadas ao longo de vários epochs, alguns commits eram incorretamente ignorados
- Uma transação que parecia ter sido commitada podia na prática ser abortada, ou o contrário
- O Jepsen descobriu esse bug porque configurou o timeout da transação para um valor baixo, de 1 segundo
- O Bufstream identificou o problema poucas horas após o lançamento da 0.1.2, bloqueou o upgrade dos clientes, e os clientes não fizeram upgrade para a 0.1.2
- A correção foi incluída na 0.1.3-rc2
-
Escritas perdidas por causa do filtering no lado do servidor (#5)
- Na 0.1.3-rc.8, uma curta janela de perda de escrita aparecia com frequência após pequenas falhas, como pause do processo Bufstream ou do coordinator, ou uma partition entre os dois
- A perda de dados ocorria independentemente do uso de transações
- Em um teste de 5 minutos, entre 16.770 records, 22 foram reconhecidos mas nenhum consumidor conseguiu fazer poll deles
- Alguns records chegavam a ser vistos por pollers por um tempo, mas depois desapareciam do poll
- A causa foi a lógica de limitação do tamanho da resposta da fetch API adicionada na 0.1.3-rc.8 para contornar um bug de uma popular GUI web do Kafka
- Um bug nessa lógica de filtering escondia records de consumidores atrasados, fazendo o problema parecer perda de escrita
- O Bufstream corrigiu isso na 0.1.3-rc.12
Problemas no cliente Java do Kafka e no protocolo Kafka
-
KIP-588:
ProducerFencedExceptionenganosa- Durante os testes, o erro
ProducerFencedException: There is a newer producer with the same transactionalId which fences the current one.ocorreu com frequência - Esse erro também apareceu em testes nos quais todos os producers recebiam IDs transacionais únicos, o que atrasou a identificação da causa
- O KIP-588 diz que
ProducerFencedExceptiontambém pode ser lançada em timeouts de transação - O cliente Java do Kafka usa
TimeoutExceptiondedicada para a maioria dos timeouts, mas neste caso lançaProducerFencedException - Mesmo sem existir um producer conflitante de fato, a mensagem de erro diz que existe uma segunda instância de producer
- O KIP-588 está aberto há dois anos, e o Jepsen recomenda que a equipe do Kafka mude a mensagem de erro
- Durante os testes, o erro
-
KAFKA-17734:
Consumer.close()pode bloquear por tempo indefinido- Tanto nos testes do Bufstream quanto nos do Kafka, os testes travavam a cada poucas horas por causa de um bug no cliente Java
Consumer.close()bloqueia em network IO por padrão- O parâmetro de timeout de
close()deveria impedir bloqueio indefinido, mas não funcionou - A abordagem de chamar
consumer.wakeup()em uma thread separada para interromper o consumer preso em IO também não teve efeito - O Jepsen considera que programas de longa execução devem conseguir liberar recursos como cliente, conexão, thread e memória em tempo razoável mesmo com erros de rede, e abriu o KAFKA-17734
-
KAFKA-17582: offset do consumer fica imprevisível após falha de transação
- A documentação oficial do Kafka quase não explica como o offset do consumer deveria se comportar quando o commit de uma transação falha
- A documentação de design do Kafka da Confluent diz que, quando uma transação é abortada, a posição do consumer volta ao valor anterior, mas o cliente Java real nem sempre se comporta assim
- Nos resultados da workload de abort, mesmo em um cluster saudável, o comportamento após abortar era difícil de prever
- A maioria dos pares de transação avançava para offsets posteriores
- Alguns voltavam para offsets anteriores
- Todos os rewinds estavam relacionados a eventos de rebalance, e todos os avanços ocorreram sem rebalance
- Segundo a resposta do lado do Kafka, esse comportamento é intencional
- O consumer continua avançando
- Se ocorrer um rebalance, ele pode voltar para um ponto arbitrário de acordo com o committed offset
- O usuário precisa fazer rewind manual da posição do consumer quando a transação é abortada
- O Jepsen abriu o KAFKA-17582 e propôs documentar esse comportamento e avaliar a mudança do rewind padrão em abort de transação
- A workload de fila também foi modificada para fazer rewind explícito do consumer
-
KAFKA-17754: perda de escrita, leitura abortada, transação fragmentada
- No Bufstream 0.1.0~0.1.3, foram observadas leitura abortada, perda de escrita e violação de atomicidade apenas com pausa do processo do Bufstream, pausa do coordinator, crash e particionamento de rede
- A análise levou a uma falha fundamental no protocolo de transações do Kafka
- No exemplo, o cliente executou uma transação com o ID transacional único
jt1234e envioucommitted = falseemEndTxnpara abortar, mas 15 chamadas depoll()observaram escritas da transação abortada - Outras escritas da mesma transação não foram observadas por nenhum poller
- Ao combinar packet capture com logs do Bufstream, a causa foi uma mensagem de commit atrasada
- Um
EndTxnde commit enviado algumas transações antes foi processado com atraso em um nó - O cliente já estava avançando para as transações seguintes
- O commit atrasado foi aplicado à transação atual, fazendo com que apenas a parte inicial da transação fosse commitada, enquanto o restante foi tratado como uma transação separada e abortado
- O protocolo Kafka foi projetado para permitir que o cliente envie requisições por várias conexões TCP e para vários nós, mas não há números de sequência que determinem a ordem das requisições do mesmo cliente
- Também não existe o conceito de número de transação, então, quando o servidor recebe uma mensagem de commit ou abort, ele não consegue saber qual transação o cliente pretendia encerrar
- Como resultado, as seguintes situações podem ocorrer
- uma transação que parece ter sido commitada na verdade é abortada
- uma transação abortada na verdade é commitada
- apenas parte das escritas da transação é preservada e parte se perde, gerando uma transação fragmentada
- O cliente oficial Java do Kafka trata timeouts como retryable e pode enviar automaticamente várias mensagens
EndTxn, então o problema pode ocorrer mesmo que o usuário chame commit ou abort apenas uma vez por transação - O Jepsen também observou leitura abortada e transação fragmentada no Kafka com pausa de processo e abriu o KAFKA-17754
- Engenheiros do Kafka acreditam que o KIP-890 pode corrigir esse problema
- O KIP-890 muda o protocolo de transações para elevar o producer epoch a cada transação
- Como o servidor rejeita mensagens de epoch anterior, isso pode impedir que mensagens de commit de transações passadas vazem para transações posteriores
- No Bufstream 0.1.3, foi adicionado um mecanismo para reduzir a frequência usando a revisão do etcd como clock lógico, mas isso ainda não impede reordenação entre o cliente e o Bufstream
- O Jepsen continuou observando leitura abortada, perda de escrita e transação fragmentada no 0.1.3, e considera que é necessária uma solução do lado do cliente
Resumo geral dos resultados
- Todos os 5 problemas do próprio Bufstream foram corrigidos
- #1: o consumer ficava travado por causa de um highest stable offset atrasado, sem necessidade de falha, corrigido na 0.1.3-rc.6
- #2: producer/consumer ficavam travados por causa do vencimento de lease do etcd, exigia pause, corrigido na 0.1.3-rc.8
- #3: offsets zero espúrios, exigia pause, corrigido na 0.1.3-rc.6
- #4: writes de transação perdidos, sem necessidade de falha, corrigido na 0.1.3-rc.2
- #5: writes perdidos por causa de filtragem no lado do servidor, exigia pause, corrigido na 0.1.3-rc.12
- Problemas relacionados ao Kafka ainda permanecem
- KIP-588: mensagem de erro incorreta em transaction timeout, não resolvido
- KAFKA-17734:
ConsumerClient.close()pode bloquear indefinidamente, não resolvido - KAFKA-17582: depois de falha de transaction, o offset do consumer fica imprevisível, não resolvido
- KAFKA-17754: perda de write, leitura abortada, transação fragmentada, não resolvido
- O Jepsen alerta que verificação experimental de segurança pode provar a existência de bugs, mas não sua ausência
- Em especial, considera que o KAFKA-17754 dificulta determinar se há outros casos de perda de write no Bufstream
Recomendações para usuários e operação do Bufstream
- Usuários que usam transactions do Bufstream com o cliente oficial Java Kafka devem considerar que, no momento, as transactions podem não ser seguras
- uma transaction abortada pode na prática ser commitada
- uma transaction commitada pode na prática ser abortada
- uma transaction pode ser partida ao meio e apenas parte dos seus efeitos ser preservada
- O Bufstream considera que o cliente Franz-go é menos vulnerável a esse problema, mas o Jepsen não testou o Franz-go com técnicas como as deste trabalho
- Outros clients podem ou não ser vulneráveis
- Usuários anteriores ao Bufstream 0.1.3 podem enfrentar os seguintes problemas
producer.send()pode retornar incorretamente o offset0em vez do offset real- problema metastável de disponibilidade em que o client fica travado
- O Jepsen recomenda upgrade para a 0.1.3
- A arquitetura geral do Bufstream parece sound
- a forma de definir a ordem de chunks de dados imutáveis com um serviço de coordenação como o etcd é uma abordagem relativamente simples, com precedentes em sistemas OLTP e de streaming
- Do ponto de vista operacional, duas melhorias são recomendadas
- se a solicitação de arquivo compartilhado do storage falhar na inicialização, o cluster pode crashar; foi recomendado adicionar retry, e o Bufstream adicionou uma camada de retry
- quando uma dependency fica unavailable, recomenda-se que o agent continue rodando em vez de morrer imediatamente, oferecendo backpressure e status do sistema e se recuperando de forma mais suave
- Na 0.1.3, o Bufstream adicionou lógica extra de retry para o etcd, mas ainda requer supervisão constante para permanecer online
- Usuários devem garantir que exista um supervisor de processos e testar se ele continua operando sem desistir mesmo durante outages prolongadas
Necessidade de documentar e corrigir o protocolo de transactions do Kafka
- A documentação oficial do Kafka quase não fala sobre transactions, então os usuários precisam combinar várias sources ambíguas e conflitantes
- O Jepsen recomendou à equipe do Kafka criar um documento central que organize com clareza a semântica de transactions e mencionou o KAFKA-17671
- Esse documento deveria ao menos especificar o seguinte
- quando o consumer observa offsets monotonicamente crescentes
- quando o consumer pode pular records que foram acknowledgeados
- se um rebalance pode afetar uma transaction no meio do processo
- quando os offsets de write do producer aumentam monotonicamente
- quando G0, G1a, G1b, G1c, fractured read e leitura dos writes da própria transaction são legais
- que significado têm o valor retornado por
poll()e os offsets após uma transaction abortada - como tratar erro de transaction, erro durante abort e erro durante rewind
- O Jepsen aponta que a documentação da Confluent repete que o padrão do Kafka fornece entrega at-least-once, mas isso aparentemente não é verdade
auto.offset.reset = latestpode fazer records não processados parecerem “commitados”- a documentação da Confluent sobre gerenciamento de offsets também fala do risco de perda do progresso de mensagens em caso de crash com auto-commit padrão
- a documentação de que o consumer é rebobinado quando uma transaction é abortada também difere da prática
- O Jepsen considera que o protocolo de transactions do Kafka precisa ser corrigido de forma fundamental
- o protocolo assume implicitamente entrega ordenada e confiável, mas existem pause de processo, falta de confiabilidade de rede, latência não nula e entrega desordenada entre vários sockets TCP
- o protocolo do Kafka distribui mensagens entre vários nodes e sockets TCP, e o client faz retry automático das mensagens
- não há sequence number para restaurar a ordem das mesmas mensagens do client, nem transaction number para verificar o alvo da transaction
- O KIP-890 tenta garantir uma ordem mais estrita elevando o epoch a cada commit de transaction
- A library cliente também poderia ajudar reinicializando o producer para elevar o epoch quando uma mensagem não é acknowledgeada
- O Java Kafka Client 3.8.0 é vulnerável a esse problema
- O Jepsen considera que o Franz-go pode mitigar ou evitar o problema ao fazer reinicialização em caso de timeout, mas não investigou outras libraries cliente
Trabalhos futuros
- Muitos usuários dependem da “exactly-once semantics” da API Kafka Streams em vez de lidar diretamente com transactions, então no futuro seria possível investigar a correção de aplicações Streams
- O Jepsen também encontrou um unseen write no Kafka enquanto investigava o KAFKA-17754, mas não conseguiu analisá-lo por limitações de tempo
- unseen write pode ser um sinal de hanging transaction, stuck consumer ou data loss
- Também permanece a dúvida se uma mensagem
Produceatrasada pode entrar em uma transaction futura e violar as garantias de transaction - Também há a suspeita de que o Kafka Java Client reutilize o sequence number quando há timeout de request, o que pode fazer com que um write seja acknowledged, mas descartado silenciosamente
- Quando ocorre um rebalance event, a consumer position pode avançar ou recuar, mas as regras disso não estão claras
- Se o Kafka documentar o comportamento pretendido, o Jepsen gostaria de validá-lo
- O Jepsen explica que, por ser um processo aleatório, é difícil explorar anomalies raras
- Problemas que acontecem uma única vez são muito difíceis de debugar e reproduzir
- O Bufstream também usa o Antithesis, que executa todo o sistema distribuído em um hypervisor determinístico e em uma simulated network
- Combinar a geração de workload e a verificação de history do Jepsen com o ambiente determinístico e reproduzível do Antithesis pode aumentar a reprodutibilidade dos testes
1 comentários
Comentários do Hacker News
Se, ao investigar issues como KAFKA-17754, também foram encontradas escritas invisíveis no Kafka, parece que está na hora de a Jepsen voltar a investigar o Kafka a fundo
A última investigação foi em 2013 (https://aphyr.com/posts/293-call-me-maybe-kafka, Kafka 0.8 beta), e agora parece que estamos em uma fase em que vários problemas estão apenas começando a ser descobertos no próprio Kafka
Algo como “uma escrita pode ser confirmada e depois silenciosamente descartada” é bem assustador
É muito surpreendente a parte em que, com o padrão
enable.auto.commit=true, um consumidor Kafka pode commitar offsets independentemente de a aplicação realmente os ter processadoNunca entendi o auto commit dessa forma, e acho que esse padrão não faz sentido se for assim
A explicação da documentação não é totalmente clara, mas, no geral, eu lia como se o offset só fosse commitado depois que o processamento terminasse
Eu entendia que ajustar o intervalo de auto commit ajudava a reduzir a janela de processamento duplicado, não a perda de mensagens, como se espera em processamento pelo menos uma vez (at-least-once)
Se você não commita explicitamente, o Kafka não tem como saber se a mensagem foi processada
O Kafka presume que a mensagem entregue foi processada imediatamente
Auto commit é parecido com entregar uma casquinha de sorvete e virar as costas imediatamente presumindo que a pessoa a comeu. Algumas pessoas podem derrubá-la assim que recebem e não conseguir comer nem uma mordida
Se você quer essa garantia, precisa fazer um acknowledgement explícito
Por exemplo, se tudo o que você faz é gravar a mensagem em um banco de dados, ela é tratada como confirmada no momento em que chega ao callback do handler do cliente
Mas, na prática, é bem provável que você queira que ela só seja confirmada depois que a inserção no DB tiver sucesso
Se o DB ficar inacessível por causa de rede, Kubernetes, configuração de firewall etc., e nesse meio-tempo um engenheiro tentar reiniciar algo e o cliente cair, é fácil acabar com mensagens não processadas
Outro sistema pode determinar se houve falha, e esse recurso permite mover o limite superior para reduzir reprocessamento
Mas, se o timing bater e ocorrer uma falha, você precisa assumir que, após o restart, pode receber de novo parte do que já foi processado
O problema é quando esse processamento não existe antes do auto commit
Pelo que dá para ler, parece que a intenção é commitar bastante tempo depois do processamento, mas também soa contraditório que, sendo auto commit, ele deva commitar apenas itens de alguns milissegundos antes do momento do auto commit
polle depois processa as mensagens de forma durávelO ponto confuso é que a verificação do auto commit não acontece de forma assíncrona depois de um timeout, mas sim na próxima chamada a
pollPortanto, só deveria ser possível perder escritas nos casos em que, antes de chamar
pollnovamente, você não processa as mensagens de forma durável e apenas as armazena — por exemplo, ao usar processamento assíncrono, atrasos, filas etc.Isso é com base no comportamento documentado da biblioteca cliente Java (https://kafka.apache.org/32/javadoc/org/apache/kafka/clients...); se a implementação atual realmente faz isso é outra questão
O protocolo Kafka fica preso entre alto e baixo nível e não faz nenhum dos dois especialmente bem
Auto commit é um recurso de alto nível que ajuda a criar aplicações simples com facilidade, mas naturalmente pode falhar se não for usado da forma esperada
Hoje em dia, acho que usuários finais deveriam usar uma implementação de nível mais alto que trate corretamente os detalhes, em vez de usar diretamente o cliente Kafka. Para casos de dados, algo como um mecanismo de processamento de streams; para casos de aplicação, algo como um mecanismo de execução contínua
Olhando a página do produto (https://buf.build/product/bufstream), fico curioso para saber como a descrição de que ele “roda apenas dentro da sua VPC na AWS ou GCP e não se comunica externamente” convive com a cobrança baseada em uso de “US$ 0,002 por GiB antes da compressão”
Imagino que eles não operem todo o negócio no sistema de honra
Claro que há risco de abuso, mas pode ser uma troca valiosa para atrair certos clientes
Se o código-fonte não está público, nunca se deve acreditar na afirmação de que “não se comunica externamente”
“O protocolo de transações do Kafka está fundamentalmente quebrado e precisa ser revisado” soa doloroso
Ainda assim, como sempre, a investigação e o texto são excelentes
Fico curioso se Kyle já avaliou o NATS JetStream. Tenho curiosidade sobre o que ele acharia
Algumas pessoas sugeriram que isso seria… como dizer… interessante :-)
Não consigo encontrar o projeto bufstream no GitHub; alguém sabe onde está?
Estranhamente, porém, ele também não tem licença
Lendo posts de blog e documentos relacionados, parece que a “entrega exatamente uma vez” do Kafka é definida como uma propriedade de uma operação ler-processar-escrever em que um worker lê do tópico 1 e escreve no tópico 2, com ambos os tópicos dentro do mesmo sistema Kafka lógico
Se for isso mesmo, acho que seria melhor chamar isso de transação
Mas há duas formas de enxergar “exatamente uma vez”
Uma é no sentido de que, como em transações de banco de dados, os efeitos não devem ser duplicados nem desaparecer
A outra se parece mais com uma propriedade de grafo de fluxo de dados sobre relações entre mensagens atravessando tópicos-partições, e fica um pouco mais próxima da consistência em ACID
Assim como um sistema de transações serializáveis garante uma certa consistência em nível de domínio, é possível usar transações para chegar a essa propriedade de fluxo de dados
Por exemplo, serializabilidade garante que invariantes preservados quando se olha cada transação isoladamente também sejam preservados em históricos concorrentes
Pode-se dizer que é dessa forma que o Kafka tenta chegar à sua “semântica exatamente uma vez”
Não confundir com https://www.warpstream.com/
Errata: “Transactions may observe none, part, or all” provavelmente deveria ser “Consumers may observe none, part, or all”
A semântica de consumidores fora de uma transação é mais nebulosa
Todas as leituras dessa carga de trabalho ocorrem em contexto transacional e passam pelo caminho de commit de offsets transacionais
Fico curioso para saber onde este software é usado. Instrumentação? Caixa-preta?
Lágrimas de alegria, claro. Receber a atenção da Jepsen já é, por si só, uma conquista