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

Kafka测试消费者获取指定Topic消息延迟3-4秒,求排查原因

Kafka消费者延迟获取消息的原因分析

测试场景

测试中创建两个Kafka Topic:topicFunc1、2_1,流程如下:

  1. 测试生产者向topicFunc1发送消息;
  2. 系统消费者接收topicFunc1消息并执行业务操作后,由系统生产者向2_1发送消息;
  3. 测试消费者尝试获取2_1消息时需等待3-4秒,但topicFunc1的消息收发是即时完成的。

日志分析

测试日志

[2022.08.29 10:55:48.094] [DEBUG] [main] [httpApi.utils.kafka.KafkaProducerUtils.publishMessagesString() at line 47]        Send message "{\"subsId\":3}" to kafka topic: topicFunc1
[2022.08.29 10:55:48.156] [DEBUG] [main] [httpApi.utils.kafka.KafkaConsumerUtilsWithConstructor.getMessages() at line 43]       execute poll())
[2022.08.29 10:55:48.156] [DEBUG] [System_Processor] [kafkaApi.KafkaStringConsumer.processRecord() at line 24]      Received message: value = {\"subsId\":3}, topic = topicFunc1
[2022.08.29 10:55:48.235] [DEBUG] [System_Processor] [kafkaApi.KafkaSimpleProducer.produce() at line 46]        sent record( message={\"result\":\"1_1\"}) meta(partition=0, offset=0) to topic = 2_1
[2022.08.29 10:55:48.235] [DEBUG] [System] [kafkaApi.EpmKafkaConsumer.commitOffsets() at line 182]      Consumer : epm.1, committing the offsets : {topicFunc1-0=OffsetAndMetadata{offset=1, leaderEpoch=null, metadata=''}}
[2022.08.29 10:55:49.156] [DEBUG] [main] [httpApi.utils.kafka.KafkaConsumerUtilsWithConstructor.getMessages() at line 45]       look at list messages!)
[2022.08.29 10:55:49.156] [DEBUG] [main] [httpApi.utils.kafka.KafkaConsumerUtilsWithConstructor.getMessages() at line 43]       execute poll())
[2022.08.29 10:55:50.172] [DEBUG] [main] [httpApi.utils.kafka.KafkaConsumerUtilsWithConstructor.getMessages() at line 45]       look at list messages!)
[2022.08.29 10:55:50.172] [DEBUG] [main] [httpApi.utils.kafka.KafkaConsumerUtilsWithConstructor.getMessages() at line 43]       execute poll())
[2022.08.29 10:55:51.187] [DEBUG] [main] [httpApi.utils.kafka.KafkaConsumerUtilsWithConstructor.getMessages() at line 45]       look at list messages!)
[2022.08.29 10:55:51.187] [DEBUG] [main] [httpApi.utils.kafka.KafkaConsumerUtilsWithConstructor.getMessages() at line 43]       execute poll())
[2022.08.29 10:55:51.218] [DEBUG] [main] [httpApi.utils.kafka.KafkaConsumerUtilsWithConstructor.getMessages() at line 45]       look at list messages!)
[2022.08.29 10:55:51.218] [DEBUG] [main] [httpApi.utils.kafka.KafkaConsumerUtilsWithConstructor.lambda$getMessages$0() at line 49]      record.value() = {\"result\":\"1_1\"})

关键时间点:

  • 系统生产者于10:55:48.235成功发送消息到2_1;
  • 测试消费者直到10:55:51.218才收到该消息。

Wireshark抓包日志

№        Source           Destination   Protocol     length         info                          time utc
2638    ipMyServer      ipKafkaServer   Kafka        238        Kafka Produce v7 Request        2022-08-29 10:55:48,233929
2639    ipKafkaServer   ipMyServer      Kafka        109        Kafka Produce v7 Response       2022-08-29 10:55:48,238583
2891    ipMyServer      ipKafkaServer   Kafka        101        Kafka OffsetFetch v5 Request    2022-08-29 10:55:51,197174
2892    ipKafkaServer   ipMyServer      Kafka        101        Kafka OffsetFetch v5 Response   2022-08-29 10:55:51,201233
2903    ipMyServer      ipKafkaServer   Kafka        109        Kafka Offsets v5 Request        2022-08-29 10:55:51,210409
2904    ipKafkaServer   ipMyServer      Kafka        105        Kafka Offsets v5 Response       2022-08-29 10:55:51,213951
2905    ipMyServer      ipKafkaServer   Kafka        147        Kafka Fetch v11 Request         2022-08-29 10:55:51,214318
2907    ipKafkaServer   ipMyServer      Kafka        190        Kafka Fetch v11 Response        2022-08-29 10:55:51,218333

关键时间点:

  • Kafka生产者请求于10:55:48,233929发出并立即收到响应,说明消息已成功写入集群;
  • 测试消费者的Fetch Request直到10:55:51,214318才发出,与测试日志中多次打印"execute poll()"的时间点矛盾,说明客户端poll()调用并未立即触发网络请求。

延迟原因诊断

从日志对比可明确:测试消费者客户端在10:55:48.156到10:55:51.214期间,虽代码层面触发poll()调用,但未真正向Kafka集群发起消息拉取请求,直到10:55:51才完成必要前置操作(偏移量同步、组初始化等),最终发起Fetch请求并获取消息。

可能的诱因

  • 消费者客户端拉取配置不合理:
    • fetch.min.bytes设置过高:Kafka集群会等待积累足够字节数的消息才返回,测试场景仅一条消息无法满足阈值,导致消费者等待至fetch.max.wait.ms超时;
    • poll()方法超时参数设置过长:若测试代码中poll(Duration.ofSeconds(3)),会导致客户端每次调用最多等待3秒才返回。
  • 消费者组初始化/再平衡延迟:
    测试消费者刚加入组时,正处于组协调器分配分区的过程中,这段时间无法发起Fetch请求。抓包日志中OffsetFetch请求直到10:55:51才发起,说明消费者此时才完成组初始化流程。
  • Topic分区分配异常:
    2_1 Topic的分区数、副本数配置异常,或消费者分区分配策略(如RangeAssignor、RoundRobinAssignor)导致分区分配延迟,消费者需等待分配完成才能拉取消息。
  • 测试代码逻辑问题:
    测试中的getMessages()方法可能存在逻辑错误:虽日志打印"execute poll()",但实际未调用Kafka Consumer的poll()方法,或调用后被阻塞在本地循环、同步逻辑中,3秒后才处理返回结果。
  • 客户端偏移量同步延迟:
    消费者需先同步最新偏移量才能发起Fetch请求,抓包日志中OffsetFetch请求在10:55:51才发起,说明之前客户端未完成偏移量初始化或同步,导致无法触发拉取操作。

内容的提问来源于stack exchange,提问作者Dmitry Prokin

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 13:18:27