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个月运行正常,问题始于一周前
排查方向建议
主题元数据与存储异常排查
- 检查
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验证
- 检查
Kafka集群资源瓶颈排查
- 监控broker的JVM内存使用:当前
KAFKA_HEAP_OPTS设置为-Xmx1G -Xms256m,对于运行18个月的生产集群可能不足,检查GC日志是否存在频繁Full GC导致的服务停顿 - 监控broker的CPU、网络负载:是否存在突发流量或进程占用过高资源的情况
- 检查控制器节点状态:KRaft模式下控制器是否稳定,是否存在控制器选举频繁的情况(查看controller日志)
- 监控broker的JVM内存使用:当前
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状态
版本兼容性排查
- 当前Kafka集群为v3.3.1,客户端(Streams/Connect)为v3.4.1,虽然跨小版本理论兼容,但需验证是否存在已知的版本兼容问题,考虑将集群同步升级至v3.4.1
内容的提问来源于stack exchange,提问作者donnie
相关产品推荐
相关产品推荐

