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

Kafka Streams频繁连接超时问题排查请求

Kafka集群数据管道超时与重平衡异常排查求助

我们有多套Kafka集群用于数据管道,其中一套生产集群已稳定运行18个月,数据流转路径为:source topic -> Kafka Streams -> destination topics -> Kafka Connect。近期突发异常:

  • Kafka Connect初始化失败
  • Kafka Streams无法处理数据
  • 消费者数分钟后持续陷入重平衡状态

初始报错为包装了TimeoutException的IllegalStateException,已知该包装问题在v3.4.1和v3.5版本中已修复,但升级后问题仍存在。

已尝试的操作

  • 调整客户端参数:
    max.poll.records = 500
    max.poll.interval.ms = 600000
    session.timeout.ms = 100000
    api.timeout.ms = 300000
    
  • 将Kafka Streams临时存储迁移至持久化卷,支持重启时增量重建RocksDB存储,无效
  • 将Kafka Connect和Kafka Streams升级至v3.4.1,仍出现「运行态→暂停→关闭」的状态循环

仅删除并重建超时涉及的主题可暂时缓解问题,但生产环境无法频繁执行此操作。推测主题本身(如运行时长、数据量)可能存在问题,但根因未明。

关键错误日志

[event-app-t1-5a919f07-826b-4527-8329-047d96fce7f6-StreamThread-1-producer] INFO org.apache.kafka.clients.Metadata - [Producer clientId=event-app-t1-5a919f07-826b-4527-8329-047d96fce7f6-StreamThread-1-producer] Resetting the last seen epoch of partition event-app-t1-KSTREAM-AGGREGATE-STATE-STORE-0000000090-changelog-1 to 0 since the associated topicId changed from null to AlNoTK-HQE-_qTVdA2Dvcg
[event-app-t1-5a919f07-826b-4527-8329-047d96fce7f6-StreamThread-1] ERROR org.apache.kafka.streams.processor.internals.TaskManager - stream-thread [event-app-t1-5a919f07-826b-4527-8329-047d96fce7f6-StreamThread-1] Error flushing caches of dirty task 1_1
java.lang.IllegalStateException: org.apache.kafka.common.errors.TimeoutException: Timeout of 300000ms expired before the position for partition some_messages-1 could be determined
    at org.apache.kafka.streams.processor.internals.StreamTask.findOffset(StreamTask.java:427)
    at org.apache.kafka.streams.processor.internals.StreamTask.committableOffsetsAndMetadata(StreamTask.java:456)
    at org.apache.kafka.streams.processor.internals.StreamTask.prepareCommit(StreamTask.java:402)
    at org.apache.kafka.streams.processor.internals.TaskManager.closeTaskDirty(TaskManager.java:1244)
    at org.apache.kafka.streams.processor.internals.TaskManager.closeAndCleanUpTasks(TaskManager.java:1391)
    at org.apache.kafka.streams.processor.internals.TaskManager.lambda$shutdown$3(TaskManager.java:1289)
    at org.apache.kafka.streams.processor.internals.TaskManager.executeAndMaybeSwallow(TaskManager.java:1812)
    at org.apache.kafka.streams.processor.internals.TaskManager.shutdown(TaskManager.java:1287)
    at org.apache.kafka.streams.processor.internals.StreamThread.completeShutdown(StreamThread.java:1173)
    at org.apache.kafka.streams.processor.internals.StreamThread.run(StreamThread.java:582)
 Caused by: org.apache.kafka.common.errors.TimeoutException: Timeout of 300000ms expired before the position for partition some_messages-1 could be determined
 [event-app-t1-5a919f07-826b-4527-8329-047d96fce7f6-StreamThread-1] INFO org.apache.kafka.streams.processor.internals.StreamTask - stream-thread [event-app-t1-5a919f07-826b-4527-8329-047d96fce7f6-StreamThread-1] task [1_1] Suspended from RUNNING

环境信息

  • Kafka集群:3节点,KRaft模式,v3.3.1
  • Kafka Streams:v3.4.1,单实例,3个流线程
  • Kafka Connect:v3.4.1,分布式模式,单实例,多连接器(Elasticsearch、BigQuery、Google Pub/Sub、HTTP)

Kafka集群核心配置

KAFKA_ENABLE_KRAFT: 'yes'
KAFKA_KRAFT_CLUSTER_ID: 'xxxxxxxxxxxxxxxx'
KAFKA_CFG_PROCESS_ROLES: broker,controller
KAFKA_CFG_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_CFG_LISTENERS: CONTROLLER://:9093,INSIDE://:9092,EXTERNAL://:9094
KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,INSIDE:PLAINTEXT,EXTERNAL:PLAINTEXT
KAFKA_CFG_CONTROLLER_QUORUM_VOTERS: 1@dpkafka01:9093,2@dpkafka02:9093,3@dpkafka03:9093
KAFKA_CFG_ADVERTISED_LISTENERS: INSIDE://:9092,EXTERNAL://_{HOSTIP}:9097
KAFKA_BROKER_ID: 1
KAFKA_CFG_NODE_ID: 1
KAFKA_INTER_BROKER_LISTENER_NAME: INSIDE
KAFKA_HEAP_OPTS: "-Xmx1G -Xms256m"
KAFKA_LOG_DIRS: /bitnami/kafka/kafka-logs
KAFKA_AUTO_CREATE_TOPICS_ENABLE: 'false'
KAFKA_LOG_RETENTION_MS: 7200000
KAFKA_LOG_SEGMENT_MS: 86400000
KAFKA_LOG_DELETE_RETENTION_MS: 7200000
KAFKA_LOG_RETENTION_CHECK_INTERVAL_MS: 300000
KAFKA_LOG_CLEANUP_POLICY: "compact,delete"
KAFKA_CFG_GROUP_INITIAL_REBALANCE_DELAY_MS: 12000
KAFKA_CFG_NUM_RECOVERY_THREADS_PER_DATA_DIR: 4
KAFKA_CFG_OFFSETS_TOPIC_REPLICATION_FACTOR: 2
KAFKA_CFG_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 2
KAFKA_CFG_TRANSACTION_STATE_LOG_MIN_ISR: 2
ALLOW_PLAINTEXT_LISTENER: 'yes'

补充说明

  • 已排除网络问题,异常在数据处理数分钟后触发
  • Kafka Streams会从运行态转为暂停,所有连接集群的客户端均出现相同超时
  • 集群此前18个月运行正常,问题始于一周前

排查方向建议

  1. 主题元数据与存储异常排查

    • 检查some_messages-1分区所在broker的磁盘状态:磁盘使用率、IO负载,是否存在磁盘瓶颈导致元数据读取或日志检索超时
    • 查看该主题的分区状态:是否存在ISR收缩、副本不同步,执行kafka-topics.sh --describe --topic some_messages --bootstrap-server <broker-addr>确认
    • 检查主题的日志段是否存在损坏:执行kafka-run-class.sh kafka.tools.DumpLogSegments --files <log-file-path> --verify验证
  2. Kafka集群资源瓶颈排查

    • 监控broker的JVM内存使用:当前KAFKA_HEAP_OPTS设置为-Xmx1G -Xms256m,对于运行18个月的生产集群可能不足,检查GC日志是否存在频繁Full GC导致的服务停顿
    • 监控broker的CPU、网络负载:是否存在突发流量或进程占用过高资源的情况
    • 检查控制器节点状态:KRaft模式下控制器是否稳定,是否存在控制器选举频繁的情况(查看controller日志)
  3. Kafka Streams状态存储与事务排查

    • 检查Kafka Streams的状态存储(RocksDB)是否存在损坏:查看RocksDB的日志文件,是否存在Corruption报错
    • 验证状态存储对应的changelog主题(如event-app-t1-KSTREAM-AGGREGATE-STATE-STORE-0000000090-changelog)是否存在异常:分区数、副本数、ISR状态,是否存在日志清理不及时导致的磁盘占用过高
    • 检查事务状态:是否存在未提交的事务导致offset提交阻塞,执行kafka-consumer-groups.sh --describe --group <streams-group-id> --bootstrap-server <broker-addr>查看消费组的offset状态
  4. 版本兼容性排查

    • 当前Kafka集群为v3.3.1,客户端(Streams/Connect)为v3.4.1,虽然跨小版本理论兼容,但需验证是否存在已知的版本兼容问题,考虑将集群同步升级至v3.4.1

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 18:45:14