You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

排查Kafka事件发布耗时过长问题请求协助

Kafka消息发布耗时异常排查请求

问题现象

某环境下向Kafka发布单条消息耗时长达3秒,其他环境仅需20毫秒。

已完成排查

  • 执行路由追踪,未发现网络延迟问题
  • 分析Kafka客户端日志,明确看到消息发布耗时3秒:

关键日志片段:

2023-01-05T07:26:24.627Z DEBUG 40 --- [ad | producer-5] o.a.k.c.NetworkClient                    : [Producer clientId=producer-5] Using older server API v7 to send PRODUCE {acks=-1,timeout=30000,partitionSizes=[esd.mob.vocc.datacollector.sampio.plab01.raw-9=464]} with correlation id 1156 to node 46

2023-01-05T07:26:24.627Z TRACE 40 --- [ad | producer-5] o.a.k.c.p.i.Sender                       : [Producer clientId=producer-5] Sent produce request to 46: (type=ProduceRequest, acks=-1, timeout=30000, partitionRecords=({esd.mob.vocc.datacollector.sampio.plab01.raw-9=MemoryRecords(size=464, buffer=java.nio.HeapByteBuffer[pos=0 lim=464 cap=464])}), transactionalId=''

2023-01-05T07:26:27.700Z TRACE 40 --- [ad | producer-5] o.a.k.c.NetworkClient                    : [Producer clientId=producer-5] Completed receive from node 46 for PRODUCE with correlation id 1156, received {responses=[{topic=esd.mob.vocc.datacollector.sampio.plab01.raw,partition_responses=[{partition=9,error_code=0,base_offset=31376,log_append_time=-1,log_start_offset=31001}]}],throttle_time_ms=0}

2023-01-05T07:26:27.700Z TRACE 40 --- [ad | producer-5] o.a.k.c.p.i.Sender                       : [Producer clientId=producer-5] Received produce response from node 46 with correlation id 1156

日志显示请求发送时间为2023-01-05T07:26:24.627Z,响应接收时间为2023-01-05T07:26:27.700Z,耗时约3秒。

客户端代码片段

当前使用的Spring Kafka客户端代码如下:

CompletableFuture<SendResult<String, String>> sendResultCompletableFuture = kafkaTemplate.send(payloadMsg);

sendResultCompletableFuture.whenComplete((data, ex) -> {
    if (Objects.nonNull(ex)) {
        LOGGER.error("Exception while publishing message to kafka, {}", ExceptionUtils.getStackTrace(ex));
        throw new KafkaConnectorException(ExceptionUtils.getStackTrace(ex), topic, message);
    } else {
        LOGGER.info("Message handle published to kafka successfully for key ::");
    }
});

进一步排查思路

1. API版本兼容性检查

日志显示客户端使用旧版服务器API v7发送请求:

  • 确认异常环境Kafka集群版本,对比正常环境集群版本,排查是否存在跨大版本调用的兼容性问题
  • 核对客户端bootstrap.servers配置,确认是否指向了正确的集群节点

2. Broker端性能与状态排查

针对日志中的目标Broker(node 46)和分区(esd.mob.vocc.datacollector.sampio.plab01.raw-9):

  • 检查Broker节点的CPU、内存、磁盘IO使用率,排查是否存在资源瓶颈(如磁盘高IO等待、JVM频繁GC)
  • 查看该分区的ISR同步状态:由于acks=-1需要等待所有ISR节点确认,若ISR节点存在同步延迟会导致耗时增加
  • 检查分区日志段大小,若日志段过大可能导致刷盘耗时变长

3. 生产者配置核对

对比正常环境与异常环境的生产者配置:

  • 确认acks配置是否一致,若异常环境ISR节点数量过多,acks=-1会增加确认耗时
  • 检查linger.ms、batch.size是否存在异常配置,导致消息未及时发送
  • 核对request.timeout.ms与Broker端replica.lag.time.max.ms等配置是否存在冲突
  • 排查是否开启了事务、幂等性等额外特性,导致流程变慢

4. 深层网络排查

  • 通过tcpdump抓包分析客户端与Broker之间的TCP连接,排查是否存在丢包、重传
  • 检查防火墙、代理设备是否存在超时或限流规则
  • 确认DNS解析是否正常,排查是否存在域名解析延迟

5. 分区负载与Leader分布

  • 查看异常分区的Leader是否长期驻留在负载较高的Broker节点上
  • 检查主题的分区分布是否均匀,排查是否存在单Broker承载过多热点分区的情况

内容的提问来源于stack exchange,提问作者R Vigneswaran

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.05 23:10:27