Apache Storm拓扑提交异常:组件订阅不存在流问题求助
解决Storm拓扑提交时的流订阅错误
从你给出的错误日志来看:
603 [main] WARN b.s.StormSubmitter - Topology submission exception: Component: [escribirFichero] subscribes from non-existent stream: [default] of component [buscamosEnKlout] Exception in thread "main" java.lang.RuntimeException: InvalidTopologyException(msg:Component: [escribirFichero] subscribes from non-existent stream: [de...
核心问题非常明确:你的Storm拓扑里,escribirFichero组件尝试订阅buscamosEnKlout组件的默认流(default),但这个流并没有在buscamosEnKlout组件中被定义或发射。
下面给你几个具体的排查和解决步骤:
1. 核对buscamosEnKlout的流定义与发射逻辑
- 先确认
buscamosEnKlout这个Spout/Bolt是否真的存在default流:如果它没有显式声明自定义流,Storm会默认提供default流,但前提是代码里确实有往这个流发射数据的逻辑。 - 比如在Bolt的
execute方法中,collector.emit(outputTuple)是往default流发射数据;但如果用了collector.emit("custom_stream", outputTuple)这种带流名称的调用,那default流就不存在。
2. 修正escribirFichero的订阅配置
- 如果
buscamosEnKlout实际只发射了自定义流(比如名为klout_data_stream),那你需要修改escribirFichero的订阅代码,指定正确的流名称:// 错误的默认流订阅 builder.setBolt("escribirFichero", new EscribirFicheroBolt()) .shuffleGrouping("buscamosEnKlout"); // 正确的自定义流订阅 builder.setBolt("escribirFichero", new EscribirFicheroBolt()) .shuffleGrouping("buscamosEnKlout", "klout_data_stream"); - 反之,如果
buscamosEnKlout本应发射default流,那就要检查它的emit调用,确保没有写错流名称,或者遗漏了数据发射的逻辑。
3. 验证拓扑的流依赖关系
- 提交拓扑前,可以通过打印拓扑元数据或者借助Storm UI的可视化功能(如果能正常启动集群的话),核对每个组件的输出流和订阅关系是否匹配,避免出现流名称不对应的情况。
额外提醒
- 别忘了在Bolt的
declareOutputFields方法里,明确声明每个流对应的字段,确保流的定义是清晰的:@Override public void declareOutputFields(OutputFieldsDeclarer declarer) { // 声明default流的字段 declarer.declare(new Fields("user_id", "content")); // 声明自定义流的字段 declarer.declareStream("klout_data_stream", new Fields("user_id", "klout_score")); }
内容的提问来源于stack exchange,提问作者Krakenudo
相关产品推荐
相关产品推荐

