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

Apache Storm KafkaSpout读取约3万条消息后停滞求助

Apache Storm KafkaSpout Stops Processing After ~30k Messages (No Tuples Reaching Bolt)

我遇到了Apache Storm与Kafka集成的问题:使用KafkaSpout从Kafka读取消息,再通过Bolt将数据写入数据库。初始阶段Spout可以正常从Kafka读取消息,Bolt也能正常处理保存,但在处理约30000条消息后,出现失败tuple,Bolt无法接收任何消息,整个拓扑彻底停滞。

查看worker.log的关键内容如下:

拓扑启动时的正常日志(成功读取分区与偏移量)

offset 9539 Read partition information from: /twitter_streaming_tweet_test/STREAMING_TWEET_WRITER_SPOUT/partition_2 --> {"partition":2,"offset":9539,"topology":{"name":"DATA_WRITER_TOPOLOGY","id":"DATA_WRITER_TOPOLOGY-67-1516077955"},"topic":"twitter_streaming_tweet_test","broker":{"port":9092,"host":"zoo1"}}
2018-01-16 17:05:57.510 o.a.s.k.PartitionManager Thread-11-STREAMING_TWEET_WRITER_SPOUT-executor[9 9] [INFO] 从Zookeeper读取最后提交的偏移量:9539;旧topology_id: DATA_WRITER_TOPOLOGY-67-1516077955 - 新topology_id: DATA_WRITER_TOPOLOGY-68-1516089922
2018-01-16 17:05:57.514 o.a.s.k.PartitionManager Thread-11-STREAMING_TWEET_WRITER_SPOUT-executor[9 9] [INFO] 从偏移量9539启动Kafka zoo1 Partition{host=zoo1:9092, topic=twitter_streaming_tweet_test, partition=2}
2018-01-16 17:05:57.518 o.a.s.k.ZkCoordinator Thread-11-STREAMING_TWEET_WRITER_SPOUT-executor[9 9] [INFO] 任务[3/3]完成刷新

处理约30000条消息时的正常日志(Bolt正常保存数据)

2018-01-16 17:06:39.732 TWLogger Thread-7-STREAMING_TWEET_WRITER_BOLT-executor[3 3] [INFO] Tweet ID 952850493570654209已保存至数据库
2018-01-16 17:06:39.739 TWLogger Thread-9-STREAMING_TWEET_WRITER_BOLT-executor[6 6] [INFO] Tweet ID 952850099335348224已保存至数据库
2018-01-16 17:06:39.742 TWLogger Thread-7-STREAMING_TWEET_WRITER_BOLT-executor[3 3] [INFO] Tweet ID 952850787981393920已保存至数据库
2018-01-16 17:06:39.753 TWLogger Thread-7-STREAMING_TWEET_WRITER_BOLT-executor[3 3] [INFO] Tweet ID 952850152573685760已保存至数据库
2018-01-16 17:06:39.754 TWLogger Thread-9-STREAMING_TWEET_WRITER_BOLT-executor[6 6] [INFO] Tweet ID 952850099578654721已保存至数据库
2018-01-16 17:06:39.763 TWLogger Thread-7-STREAMING_TWEET_WRITER_BOLT-executor[3 3] [INFO] Tweet ID 952850153173524481已保存至数据库
2018-01-16 17:06:39.768 TWLogger Thread-9-STREAMING_TWEET_WRITER_BOLT-executor[6 6] [INFO] Tweet ID 952850099989704705已保存至数据库
2018-01-16 17:06:39.776 TWLogger Thread-7-STREAMING_TWEET_WRITER_BOLT-executor[3 3] [INFO] Tweet ID 952850153232154624已保存至数据库
2018-01-16 17:06:39.779 TWLogger Thread-9-STREAMING_TWEET_WRITER_BOLT-executor[6 6] [INFO] Tweet ID 952850758289956864已保存至数据库
2018-01-16 17:06:39.787 TWLogger Thread-7-STREAMING_TWEET_WRITER_BOLT-executor[3 3] [INFO] Tweet ID 952850154436018176已保存至数据库

异常前的刷新日志与停滞现象

随后出现周期性的分区管理器刷新日志,但之后Kafka Spout尝试从Zookeeper读取分区信息时无结果,导致没有tuple被处理,拓扑陷入停滞:

2018-01-16 17:07:56.106 o.a.s.k.ZkCoordinator Thread-11-STREAMING_TWEET_WRITER_SPOUT-executor[9 9] [INFO] 任务[3/3]正在刷新分区管理器连接
2018-01-16 17:07:56.117 o.a.s.k.DynamicBrokersReader Thread-11-STREAMING_TWEET_WRITER_SPOUT-executor[9 9] [INFO] 从Zookeeper读取分区信息:GlobalPartitionInformation{topic=twitter_streaming_tweet_test, partitionMap={0=zoo2:9092, 1=zoo3:9092, 2=zoo1:9092}}
2018-01-16 17:07:56.117 o.a.s.k.KafkaUtils Thread-11-STREAMING_TWEET_WRITER_SPOUT-executor[9 9] [INFO] 任务[3/3]分配到[Partition{host=zoo1:9092, topic=twitter_streaming_tweet_test, partition=2}]
2018-01-16 17:07:56.117 o.a.s.k.ZkCoordinator Thread-11-STREAMING_TWEET_WRITER_SPOUT-executor[9 9] [INFO] 任务[3/3]已删除的分区管理器:[]
2018-01-16 17:07:56.117 o.a.s.k.ZkCoordinator Thread-11-STREAMING_TWEET_WRITER_SPOUT-executor[9 9] [INFO] 任务[3/3]新增的分区管理器:[]
2018-01-16 17:07:56.117 o.a.s.k.ZkCoordinator Thread-11-STREAMING_TWEET_WRITER_SPOUT-executor[9 9] [INFO] 任务[3/3]完成刷新
2018-01-16 17:09:54.150 o.a.s.k.ZkCoordinator Thread-11-STREAMING_TWEET_WRITER_SPOUT-executor[9 9] [INFO] 任务[3/3]正在刷新分区管理器连接
2018-01-16 17:09:54.160 o.a.s.k.DynamicBrokersReader Thread-11-STREAMING_TWEET_WRITER_SPOUT-executor[9 9] [INFO] 从Zookeeper读取分区信息:GlobalPartitionInformation{topic=twitter_streaming_tweet_test, partitionMap={0=zoo2:9092, 1=zoo3:9092, 2=zoo1:9092}}
2018-01-16 17:09:54.160 o.a.s.k.KafkaUtils Thread-11-STREAMING_TWEET_WRITER_SPOUT-executor[9 9] [INFO] 任务[3/3]分配到[Partition{host=zoo1:9092, topic=twitter_streaming_tweet_test, partition=2}]
2018-01-16 17:09:54.160 o.a.s.k.ZkCoordinator Thread-11-STREAMING_TWEET_WRITER_SPOUT-executor[9 9] [INFO] 任务[3/3]已删除的分区管理器:[]
2018-01-16 17:09:54.160 o.a.s.k.ZkCoordinator Thread-11-STREAMING_TWEET_WRITER_SPOUT-executor[9 9] [INFO] 任务[3/3]新增的分区管理器:[]
2018-01-16 17:09:54.160 o.a.s.k.ZkCoordinator Thread-11-STREAMING_TWEET_WRITER_SPOUT-executor[9 9] [INFO] 任务[3/3]完成刷新
2018-01-16 17:10:56.108 o.a.s.k.ZkCoordinator Thread-11-STREAMING_TWEET_WRITER_SPOUT-executor[9 9] [INFO] 任务[3/3]正在刷新分区管理器连接

恳请各位帮忙排查并解决这个问题,非常感谢!


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 08:41:55