Todos os produtos
Search
Central de documentação

ApsaraMQ for Kafka:A configuração de callback no cliente Java afeta a velocidade de envio de mensagens?

Última atualização: Jun 27, 2026

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);
          }
      });