
Isolamento de Backpressure em Kafka Connect Sink para Alto Volume
Isolamento de Backpressure em Kafka Connect Sink evita picos de lag ao gravar no banco. Recupere SLOs de streaming e reduza o desperdício de memória.
Isolamento de Backpressure em Kafka Connect Sink para Alto Volume
Isolamento de Backpressure em Kafka Connect Sink falhou durante um pico de 400% na carga CDC, empurrando o lag para 1,8 milhão de mensagens e derrubando os nós do cluster. Quando o banco de dados de destino sofre com contenção de locks ou gargalos de I/O, os conectores de sink continuam consumindo dados dos brokers até estourar os limites de memória JVM. O comportamento padrão do Kafka Connect—fazer polling de lotes na velocidade máxima da rede enquanto confia em buffers em memória sem limite—transforma lentidões pontuais em interrupções sistêmicas do cluster. Entender como delimitar o consumo de memória, ajustar o polling do consumidor e implementar gatilhos de backpressure explícitos no ciclo de execução da task é essencial para garantir SLOs de baixa latência.
Análise de Causa Raiz: Esgotamento de JVM Heap e Tempestades de Rebalanceamento
Quando um destino como PostgreSQL, Elasticsearch ou Snowflake perde capacidade de escrita, o worker do Kafka Connect entra em um estado crítico de execução. O framework utiliza threads internas de consumo que buscam registros dos brokers via KafkaConsumer.poll() e os entregam ao método SinkTask.put(). Se o banco de destino reduz sua taxa de 10.000 gravações por segundo para 200 registros devido a bloqueios na tabela, a task passa a levar muito mais tempo processando a chamada put() ou nos ciclos de flush do lote.
Enquanto a thread do SinkTask.put() permanece bloqueada aguardando o retorno do soquete TCP ou a liberação de conexões do pool, o loop de polling do consumidor fica paralisado. Se essa pausa ultrapassar o limite definido em max.poll.interval.ms (que por padrão é de 300.000 milissegundos ou 5 minutos), o coordenador do grupo no broker considera o consumidor morto. Uma tempestade de rebalanceamento é disparada no cluster. Durante o rebalanceamento, as partições atribuídas são revogadas, lotes em transação são desfeitos e offsets não gravados resultam em duplicidades no destino. Quando a reatribuição termina, os novos workers tentam processar o mesmo lote gigante, repetindo o travamento e gerando outro rebalanceamento em ciclo contínuo.
Paralelamente, a pressão sobre a memória JVM atinge níveis críticos. Caso o conector acumule registros em listas internas antes de executar operações de inserção em massa (bulk operations), o Garbage Collector da JVM não consegue liberar memória de objetos de vida curta. Pausar o processo por conta de Stop-The-World (STW) GC impede o envio de pings do thread de heartbeat definido em heartbeat.interval.ms. Com isso, o nó é expulso da topologia e causa falhas em cascata em outros conectores que compartilham a mesma infraestrutura.
Delimitando a Memória com Configurações Finas de Consumo
Para eliminar o estouro de memória, é indispensável controlar a ingestão no nível do consumidor Kafka antes mesmo que as mensagens sejam alocadas no código da aplicação. As principais alavancas para ajustar essa taxa são o limite de max.poll.records e o dimensionamento correto do heap da JVM em relação ao tamanho máximo das mensagens.