Sim. O impacto depende do tempo de execução do callback e do valor de max.in.flight.requests.per.connection.
Por que callbacks podem reduzir a velocidade de envio
O producer Java do Kafka executa callbacks na thread de I/O interna. Se um callback executar tarefas demoradas — como gravar em um banco de dados ou fazer uma chamada HTTP —, ele bloqueará a thread de I/O e impedirá o producer de enviar novas mensagens até que o callback retorne.
Dois fatores determinam a gravidade do impacto:
Tempo de processamento do callback. Quanto maior a duração de cada callback, maior será o bloqueio da thread de I/O. Durante esse período, o producer não consegue enviar mensagens adicionais.
**
max.in.flight.requests.per.connection.** Esse parâmetro controla quantas requisições não confirmadas o producer pode enviar em uma única conexão antes de parar e aguardar. Antes da conclusão de um callback bloqueante, o producer envia, no máximo, essa quantidade de requisições adicionais. Ao atingir o limite, o envio é interrompido até que os callbacks terminem e liberem slots em trânsito.
Mantenha os callbacks rápidos
Para evitar a degradação do throughput:
-
Transfira tarefas pesadas para uma thread separada. Processe os resultados do callback de forma assíncrona em uma thread dedicado ou pool de threads, em vez de executar o trabalho inline.
ExecutorService callbackExecutor = Executors.newFixedThreadPool(4); producer.send(record, (metadata, exception) -> { // Hand off to a separate thread to keep the I/O thread unblocked callbackExecutor.submit(() -> { if (exception != null) { logger.error("Send failed for topic {}", metadata.topic(), exception); } else { // Perform time-consuming processing here persistOffset(metadata.topic(), metadata.partition(), metadata.offset()); } }); }); -
Processe callbacks em lote. Acumule uma quantidade específica de ACKs antes de processá-los em conjunto, em vez de agir sobre cada um individualmente.
AtomicInteger ackCount = new AtomicInteger(0); int batchThreshold = 100; producer.send(record, (metadata, exception) -> { if (ackCount.incrementAndGet() >= batchThreshold) { ackCount.set(0); // Process the accumulated batch flushMetrics(); } }); -
Mantenha os callbacks inline leves. Registrar logs ou incrementar um contador é aceitável. Evite chamadas de rede, I/O de arquivo ou qualquer operação que possa causar bloqueio.
// Lightweight callback -- safe to run on the I/O thread producer.send(record, (metadata, exception) -> { if (exception != null) { logger.error("Send failed", exception); } });