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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 07:17:18