在Trident拓扑中配置Bolt读取其他Bolt的特定流
实现TridentTopology中Bolt订阅其他Bolt的特定流
我来帮你搞定这个Trident拓扑里Bolt订阅特定流的问题!下面我会结合你的初始代码,一步步展示具体实现方式,关键细节也会给你讲清楚。
核心思路
要让一个Bolt订阅另一个Bolt的特定流,需要做两件关键事:
- 发射特定流的Bolt,必须在
declareOutputFields里明确声明这个流的字段信息 - 订阅的Bolt,要通过上游流的
getStream("流名称")方法获取到目标特定流,再绑定自己的处理逻辑
完整代码示例
1. 初始拓扑与Spout设置
首先补全你给出的初始代码:
import org.apache.storm.trident.TridentTopology; import org.apache.storm.trident.Stream; import org.apache.storm.trident.spout.FixedBatchSpout; import org.apache.storm.tuple.Fields; import org.apache.storm.tuple.Values; import org.apache.storm.Config; import org.apache.storm.LocalCluster; import org.apache.storm.topology.base.BaseRichBolt; import org.apache.storm.topology.OutputFieldsDeclarer; import org.apache.storm.tuple.Tuple; import org.apache.storm.task.OutputCollector; import java.util.Map; import java.util.HashMap; public class TridentStreamExample { public static void main(String[] args) { TridentTopology tridentTopology = new TridentTopology(); // 初始化你的FixedBatchSpout FixedBatchSpout spout = new FixedBatchSpout(new Fields("sentence"), 2, new Values("the cow jumped over the moon"), new Values("the man went to the store and bought some candy"), new Values("four score and seven years ago"), new Values("how many apples can you eat")); // 将Spout添加到拓扑,得到初始输入流 Stream inputStream = tridentTopology.newStream("sentence-spout", spout);
2. 定义发射特定流的Bolt
这里我们创建一个SplitBolt,它会拆分句子,同时向默认流发射所有单词,向**特定流filtered_words**发射长度大于3的单词:
// 自定义SplitBolt,发射默认流和特定流 class SplitBolt extends BaseRichBolt { private OutputCollector collector; @Override public void prepare(Map<String, Object> conf, TopologyContext context, OutputCollector collector) { this.collector = collector; } @Override public void execute(Tuple tuple) { String sentence = tuple.getString(0); String[] words = sentence.split(" "); for (String word : words) { // 发射到默认流 collector.emit(new Values(word)); // 过滤长单词,发射到特定流"filtered_words" if (word.length() > 3) { collector.emit("filtered_words", new Values(word)); } } collector.ack(tuple); } @Override public void declareOutputFields(OutputFieldsDeclarer declarer) { // 声明默认流的字段 declarer.declare(new Fields("word")); // 声明特定流"filtered_words"的字段 declarer.declareStream("filtered_words", new Fields("long_word")); } }
3. 绑定SplitBolt并获取特定流
把SplitBolt添加到拓扑,然后通过getStream拿到它的特定流:
// 将SplitBolt绑定到输入流,得到默认输出流 Stream splitDefaultStream = inputStream.each(new Fields("sentence"), new SplitBolt(), new Fields("word")); // 获取SplitBolt发射的特定流"filtered_words" Stream filteredLongWordsStream = splitDefaultStream.getStream("filtered_words");
4. 定义订阅特定流的Bolt
创建一个CountBolt,用来统计特定流里的长单词出现次数:
// 自定义CountBolt,订阅特定流进行计数 class CountBolt extends BaseRichBolt { private OutputCollector collector; private Map<String, Integer> wordCountMap = new HashMap<>(); @Override public void prepare(Map<String, Object> conf, TopologyContext context, OutputCollector collector) { this.collector = collector; } @Override public void execute(Tuple tuple) { String longWord = tuple.getString(0); int count = wordCountMap.getOrDefault(longWord, 0) + 1; wordCountMap.put(longWord, count); // 发射统计结果 collector.emit(new Values(longWord, count)); collector.ack(tuple); } @Override public void declareOutputFields(OutputFieldsDeclarer declarer) { declarer.declare(new Fields("long_word", "count")); } }
5. 让CountBolt订阅特定流
最后把CountBolt绑定到我们拿到的特定流上:
// 将CountBolt绑定到特定流,完成订阅 filteredLongWordsStream.each(new Fields("long_word"), new CountBolt(), new Fields("long_word", "count")); // 提交拓扑到本地集群测试 Config config = new Config(); LocalCluster cluster = new LocalCluster(); cluster.submitTopology("trident-specific-stream-topology", config, tridentTopology.build()); } }
关键细节说明
- 声明特定流:必须在Bolt的
declareOutputFields中用declareStream方法指定流名称和对应的字段,否则Trident无法识别这个流 - 获取特定流:上游流处理后,通过
getStream("流名称")可以精准拿到目标流,而不是默认的输出流 - 订阅逻辑:后续的Bolt操作(比如
each)直接绑定到这个特定流上,就只会处理该流发送的Tuple
内容的提问来源于stack exchange,提问作者Swetha
相关产品推荐
相关产品推荐

