如何在Apache Storm拓扑中实现自定义对象在Bolt间传递?
解决方案:Storm Bolt间传递自定义统计对象并写入Redis
没问题,我来一步步帮你实现这个需求——让Bolt 1传递包含推文列表和对应出现次数的自定义对象,再由Bolt 2接收后写入Redis。
1. 定义可序列化的自定义数据对象
首先得搞一个能在Bolt之间传递的自定义类,Storm要求Tuple里的对象必须实现Serializable接口,不然没法跨节点序列化传输。我们就叫它TweetStats,包含两个核心列表:
import java.io.Serializable; import java.util.List; public class TweetStats implements Serializable { private List<String> tweets; private List<Integer> occurrenceCounts; // 构造方法 public TweetStats(List<String> tweets, List<Integer> occurrenceCounts) { this.tweets = tweets; this.occurrenceCounts = occurrenceCounts; } // Getter方法(Bolt 2需要读取这些数据) public List<String> getTweets() { return tweets; } public List<Integer> getOccurrenceCounts() { return occurrenceCounts; } }
2. Bolt 1:构建并发射自定义对象
在Bolt 1里,你需要先收集好推文列表和对应的出现次数,然后封装成TweetStats对象,通过Storm的OutputCollector发射出去。记得要在declareOutputFields里声明我们要传递的自定义对象的字段名,比如tweetStats:
import org.apache.storm.task.OutputCollector; import org.apache.storm.task.TopologyContext; import org.apache.storm.topology.OutputFieldsDeclarer; import org.apache.storm.topology.base.BaseRichBolt; import org.apache.storm.tuple.Fields; import org.apache.storm.tuple.Tuple; import org.apache.storm.tuple.Values; import java.util.List; import java.util.Map; public class TweetStatsEmitterBolt extends BaseRichBolt { private OutputCollector collector; @Override public void prepare(Map<String, Object> topoConf, TopologyContext context, OutputCollector collector) { this.collector = collector; } @Override public void execute(Tuple input) { // 这里假设你已经从上游获取/计算好了推文列表和对应的出现次数 // 示例数据,实际业务中替换成你的真实数据 List<String> tweets = List.of("storm is awesome", "redis is fast", "storm + redis"); List<Integer> counts = List.of(5, 12, 8); // 封装成自定义对象 TweetStats tweetStats = new TweetStats(tweets, counts); // 发射自定义对象到下游Bolt collector.emit(new Values(tweetStats)); // 确认Tuple处理完成 collector.ack(input); } @Override public void declareOutputFields(OutputFieldsDeclarer declarer) { // 声明输出字段,字段名随便取,只要Bolt 2能对应上就行 declarer.declare(new Fields("tweetStats")); } }
3. Bolt 2:接收对象并写入Redis
Bolt 2的任务就是接收Bolt 1发来的TweetStats对象,然后把推文和对应的次数写入Redis。这里推荐用Jedis连接池来管理Redis连接,避免频繁创建销毁连接:
import org.apache.storm.task.OutputCollector; import org.apache.storm.task.TopologyContext; import org.apache.storm.topology.OutputFieldsDeclarer; import org.apache.storm.topology.base.BaseRichBolt; import org.apache.storm.tuple.Tuple; import redis.clients.jedis.Jedis; import redis.clients.jedis.JedisPool; import redis.clients.jedis.JedisPoolConfig; import java.util.List; import java.util.Map; public class RedisWriterBolt extends BaseRichBolt { private OutputCollector collector; private JedisPool jedisPool; @Override public void prepare(Map<String, Object> topoConf, TopologyContext context, OutputCollector collector) { this.collector = collector; // 初始化Redis连接池,这里的配置可以根据你的实际环境调整 JedisPoolConfig poolConfig = new JedisPoolConfig(); poolConfig.setMaxTotal(10); poolConfig.setMaxIdle(5); // 替换成你的Redis地址和端口 jedisPool = new JedisPool(poolConfig, "localhost", 6379); } @Override public void execute(Tuple input) { try (Jedis jedis = jedisPool.getResource()) { // 从Tuple里取出自定义对象 TweetStats tweetStats = (TweetStats) input.getValueByField("tweetStats"); List<String> tweets = tweetStats.getTweets(); List<Integer> counts = tweetStats.getOccurrenceCounts(); // 遍历数据,写入Redis(这里用String类型存储,key是推文,value是次数) for (int i = 0; i < tweets.size(); i++) { String tweet = tweets.get(i); int count = counts.get(i); jedis.set(tweet, String.valueOf(count)); // 如果需要累加次数,可以用jedis.incrBy(tweet, count),根据你的需求调整 } collector.ack(input); } catch (Exception e) { // 处理异常,失败的话可以fail让Storm重试 collector.fail(input); e.printStackTrace(); } } @Override public void declareOutputFields(OutputFieldsDeclarer declarer) { // 这个Bolt不需要输出数据,所以不用声明 } @Override public void cleanup() { // 关闭Redis连接池 if (jedisPool != null) { jedisPool.close(); } } }
4. 拓扑配置:连接两个Bolt
最后需要把这两个Bolt组装成Storm拓扑,指定它们之间的数据流关系:
import org.apache.storm.Config; import org.apache.storm.LocalCluster; import org.apache.storm.topology.TopologyBuilder; import org.apache.storm.tuple.Fields; public class TweetTopology { public static void main(String[] args) { TopologyBuilder builder = new TopologyBuilder(); // 设置Bolt 1,这里假设你有一个Spout作为上游,没有的话可以用DummySpout或者直接从Bolt 1开始 builder.setBolt("tweet-stats-emitter", new TweetStatsEmitterBolt()) .shuffleGrouping("your-spout-id"); // 替换成你的上游Spout/组件ID // 设置Bolt 2,订阅Bolt 1的输出 builder.setBolt("redis-writer", new RedisWriterBolt()) .fieldsGrouping("tweet-stats-emitter", new Fields("tweetStats")); Config config = new Config(); config.setDebug(true); // 本地运行拓扑,生产环境提交到集群的话用StormSubmitter.submitTopology LocalCluster cluster = new LocalCluster(); cluster.submitTopology("tweet-redis-topology", config, builder.createTopology()); } }
注意事项
- 确保你的项目依赖里包含Storm和Jedis的相关jar包,比如Maven依赖:
<dependencies> <dependency> <groupId>org.apache.storm</groupId> <artifactId>storm-core</artifactId> <version>2.4.0</version> <scope>provided</scope> </dependency> <dependency> <groupId>redis.clients</groupId> <artifactId>jedis</artifactId> <version>4.3.1</version> </dependency> </dependencies>
- 如果你的自定义对象里有复杂类型,也要确保它们都实现了
Serializable接口; - Redis的配置(地址、端口、密码等)要根据你的实际环境调整,生产环境建议把配置放到Storm的
topoConf里,不要硬编码。
内容的提问来源于stack exchange,提问作者Ram Repaka
相关产品推荐
相关产品推荐

