一、先搞懂两个工具的基础定位

很多做实时数据处理的开发者,经常会在Storm和Flink之间纠结选哪个,尤其是接手老系统的时候,总怕选错了踩坑。其实不用把这俩想成复杂的“黑科技”,可以先把它们当成两种不同的“流水线加工厂”——都是用来把源源不断的原材料(实时数据),加工成成品(可用的业务结果),但两者的加工逻辑、效率、灵活性完全不一样。

先给大家说下最基础的区别:Storm是很早之前(2011年左右)由Twitter开源的实时处理框架,主打“低延迟”,也就是加工速度快;而Flink是后来(2014年左右)由Apache基金会推出的,主打“精准一次”的加工保证,也就是加工过程中不会丢数据、不会重复加工,同时也兼顾了低延迟。

为了让大家更直观,这里先放一个最基础的示例,统一用Java作为技术栈,展示两个工具的核心处理逻辑。 技术栈:Java 8+、Apache Storm 2.4.0、Apache Flink 1.17.0

1.1 Storm的基础示例:单词计数(实时)

Storm的核心概念是“拓扑(Topology)”,就像流水线的整体布局,由“Spout(数据源,负责取原材料)”和“Bolt(加工单元,负责处理数据)”组成。

// 自定义Spout:模拟源源不断的数据源,比如从Kafka取数据
public class WordSpout extends BaseRichSpout {
    private SpoutOutputCollector collector;
    // 模拟数据源:随机生成句子
    private static final String[] SENTENCES = {
        "hello world", "storm flink", "hello flink", "storm hello"
    };

    @Override
    public void open(Map conf, TopologyContext context, SpoutOutputCollector collector) {
        this.collector = collector;
    }

    @Override
    public void nextTuple() {
        // 随机选一个句子,模拟实时数据到来
        String sentence = SENTENCES[new Random().nextInt(SENTENCES.length)];
        // 发射数据:将句子传给下游Bolt
        collector.emit(new Values(sentence));
        // 模拟数据间隔:每1秒来一条
        try { Thread.sleep(1000); } catch (InterruptedException e) { e.printStackTrace(); }
    }

    @Override
    public void declareOutputFields(OutputFieldsDeclarer declarer) {
        // 定义输出字段名,方便下游识别
        declarer.declare(new Fields("sentence"));
    }
}

// 自定义Bolt1:拆分句子为单词
public class SplitBolt extends BaseRichBolt {
    private OutputCollector collector;

    @Override
    public void prepare(Map conf, TopologyContext context, OutputCollector collector) {
        this.collector = collector;
    }

    @Override
    public void execute(Tuple tuple) {
        // 从上游获取句子
        String sentence = tuple.getStringByField("sentence");
        // 拆分句子为单词
        String[] words = sentence.split(" ");
        // 每个单词单独发射
        for (String word : words) {
            collector.emit(new Values(word));
        }
    }

    @Override
    public void declareOutputFields(OutputFieldsDeclarer declarer) {
        declarer.declare(new Fields("word"));
    }
}

// 自定义Bolt2:统计单词数量
public class CountBolt extends BaseRichBolt {
    private OutputCollector collector;
    // 用Map存储每个单词的计数
    private Map<String, Integer> countMap = new HashMap<>();

    @Override
    public void prepare(Map conf, TopologyContext context, OutputCollector collector) {
        this.collector = collector;
    }

    @Override
    public void execute(Tuple tuple) {
        // 从上游获取单词
        String word = tuple.getStringByField("word");
        // 计数加1
        countMap.put(word, countMap.getOrDefault(word, 0) + 1);
        // 发射结果
        collector.emit(new Values(word, countMap.get(word)));
        // 打印结果
        System.out.println("Storm统计结果:" + word + " -> " + countMap.get(word));
    }

    @Override
    public void declareOutputFields(OutputFieldsDeclarer declarer) {
        declarer.declare(new Fields("word", "count"));
    }
}

// 主函数:构建Storm拓扑并提交
public class StormWordCountTopology {
    public static void main(String[] args) throws Exception {
        // 构建拓扑
        TopologyBuilder builder = new TopologyBuilder();
        // 设置Spout:并发数为1
        builder.setSpout("word-spout", new WordSpout(), 1);
        // 设置SplitBolt:并发数为2,上游是word-spout
        builder.setBolt("split-bolt", new SplitBolt(), 2).shuffleGrouping("word-spout");
        // 设置CountBolt:并发数为1,上游是split-bolt
        builder.setBolt("count-bolt", new CountBolt(), 1).fieldsGrouping("split-bolt", new Fields("word"));

        // 本地运行拓扑(生产环境需提交到集群)
        LocalCluster cluster = new LocalCluster();
        cluster.submitTopology("word-count-topology", new Config(), builder.createTopology());
        // 运行10秒后关闭
        Thread.sleep(10000);
        cluster.shutdown();
    }
}

1.2 Flink的基础示例:单词计数(实时)

Flink的核心概念是“流(Stream)”和“算子(Operator)”,核心处理逻辑是“流上的转换”,比Storm更强调“状态”的管理,也就是加工过程中临时存储的中间结果(比如上面示例中的计数)。

// 主函数:构建Flink流处理任务
public class FlinkWordCountJob {
    public static void main(String[] args) throws Exception {
        // 1. 创建Flink执行环境
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        // 开启精准一次保证(Flink核心特性)
        env.enableCheckpointing(1000); // 每1秒做一次检查点(状态备份)

        // 2. 模拟数据源:与Storm示例相同的随机句子
        DataStream<String> sentenceStream = env.addSource(new SourceFunction<String>() {
            private boolean isRunning = true;
            private static final String[] SENTENCES = {
                "hello world", "storm flink", "hello flink", "storm hello"
            };

            @Override
            public void run(SourceContext<String> ctx) throws Exception {
                while (isRunning) {
                    // 随机选一个句子,模拟实时数据
                    String sentence = SENTENCES[new Random().nextInt(SENTENCES.length)];
                    // 发射数据
                    ctx.collect(sentence);
                    // 模拟数据间隔
                    Thread.sleep(1000);
                }
            }

            @Override
            public void cancel() {
                isRunning = false;
            }
        });

        // 3. 流转换:拆分、计数
        DataStream<Tuple2<String, Integer>> countStream = sentenceStream
                // 拆分句子为单词
                .flatMap((FlatMapFunction<String, String>) (sentence, out) -> {
                    String[] words = sentence.split(" ");
                    for (String word : words) {
                        out.collect(word);
                    }
                })
                // 给每个单词初始计数1
                .map((MapFunction<String, Tuple2<String, Integer>>) word -> new Tuple2<>(word, 1))
                // 按单词分组
                .keyBy(value -> value.f0)
                // 统计计数:这里的sum是Flink自带的状态化算子
                .sum(1);

        // 4. 打印结果
        countStream.print("Flink统计结果:");

        // 5. 执行任务
        env.execute("Flink Word Count Job");
    }
}

从这两个示例就能看出最直观的区别:Storm需要手动管理中间状态(比如上面的countMap),而Flink自带状态管理,不需要开发者自己写Map来存计数,这也是Flink能实现“精准一次”的核心原因。

二、核心功能差异对比

接下来咱们从几个实际开发中最关心的点,来详细对比两者的功能差异,这些差异也是选型的核心依据。

2.1 数据处理保证:丢不丢数据、重不重复

这是实时处理最核心的问题,比如你做订单统计,要是丢了一笔订单,或者重复统计了一笔,业务就会出大问题。

  • Storm:默认只能做到“至少一次”保证,也就是数据不会丢,但可能重复处理。如果要实现“精准一次”,需要开发者自己做很多额外工作,比如手动管理状态、自己做事务控制,非常麻烦,开发成本很高。
  • Flink:原生就支持“精准一次”保证,只需要开启检查点(Checkpoint)功能(上面示例中env.enableCheckpointing(1000)就是开启检查点),Flink会自动备份状态,就算任务挂了重启,也能从最近的检查点恢复,不会丢数据也不会重复。

2.2 状态管理:中间结果的存储

状态就是加工过程中需要临时存的中间数据,比如上面示例中的单词计数、用户的会话信息、订单的累计金额等。

  • Storm:没有原生的状态管理,所有状态都需要开发者自己实现,比如用HashMap存、或者存在外部数据库(比如Redis)里,不仅开发麻烦,性能也不好,还容易出问题。
  • Flink:原生支持状态管理,状态可以存在内存里(性能高),也可以存在磁盘里(数据量大的时候用),Flink会自动管理状态的备份、恢复,开发者只需要调用自带的算子(比如sum、count)就能用,非常方便。

2.3 延迟:加工速度

延迟就是从数据进来,到加工结果出来的时间,延迟越低,实时性越高。

  • Storm:延迟非常低,正常情况下可以做到毫秒级,适合对延迟要求极高的场景,比如实时告警、实时风控。
  • Flink:延迟也很低,大部分场景下也能做到毫秒级,虽然比Storm稍微高一点,但差距非常小,普通业务场景根本感知不到。

2.4 窗口处理:按时间/数量分组加工

窗口处理是实时处理中非常常用的功能,比如统计每5分钟的订单量、每10条消息的平均值,就需要用到窗口。

  • Storm:没有原生的窗口处理,需要开发者自己实现,比如手动记录时间、手动分组,开发难度大,容易出bug。
  • Flink:原生支持丰富的窗口处理,比如滚动窗口(按固定时间分组,不重叠)、滑动窗口(按固定时间分组,重叠)、会话窗口(按用户会话分组),只需要调用自带的算子就能实现,非常方便。

给大家举个Flink窗口处理的示例,比如统计每5分钟的单词计数: 技术栈:Java 8+、Apache Flink 1.17.0

// 基于上面的单词计数示例,增加窗口处理
DataStream<Tuple2<String, Integer>> windowCountStream = sentenceStream
        .flatMap((FlatMapFunction<String, String>) (sentence, out) -> {
            String[] words = sentence.split(" ");
            for (String word : words) {
                out.collect(word);
            }
        })
        .map((MapFunction<String, Tuple2<String, Integer>>) word -> new Tuple2<>(word, 1))
        .keyBy(value -> value.f0)
        // 增加滚动窗口:每5分钟统计一次
        .window(TumblingProcessingTimeWindows.of(Time.minutes(5)))
        // 统计窗口内的计数
        .sum(1);

这个示例就是按每5分钟的时间窗口,统计每个单词的出现次数,Flink只需要一行代码就能实现,要是用Storm的话,得写几十行代码来实现窗口逻辑。

2.5 批流统一:同时处理离线和实时数据

现在很多公司需要同时处理离线数据(比如每天的历史数据统计)和实时数据(比如当天的实时数据统计),如果用不同的框架,维护成本非常高。

  • Storm:只能处理实时数据,不能处理离线数据,要是需要处理离线数据,得用其他框架(比如MapReduce、Spark)。
  • Flink:支持批流统一,也就是同一个框架既能处理实时数据,也能处理离线数据,只需要改一下执行环境(比如用ExecutionEnvironment代替StreamExecutionEnvironment),就能实现离线数据处理,大大降低了维护成本。

三、应用场景对比

不同的场景适合不同的工具,咱们结合实际的业务场景来分析。

3.1 适合Storm的场景

  • 对延迟要求极高,且对数据处理保证要求不高的场景,比如实时日志采集、实时监控告警,只要日志不丢就行,重复一点也没关系。
  • 非常简单的实时处理场景,比如简单的消息过滤、简单的消息转发,不需要复杂的窗口、状态处理。
  • 老系统已经用了Storm,且业务逻辑非常简单,改造成本很高的场景。

3.2 适合Flink的场景

  • 对数据处理保证要求极高的场景,比如实时订单统计、实时金融交易、实时用户画像,不能丢数据也不能重复。
  • 复杂的实时处理场景,比如需要窗口处理、状态管理、多流关联的场景,比如实时推荐、实时风控。
  • 需要同时处理离线和实时数据的场景,比如需要做数仓的实时数仓、离线数仓统一处理。
  • 新开发的实时处理系统,优先选Flink,因为Flink的功能更全,开发更简单,维护成本更低。

四、存量系统中的应用选择

很多开发者接手的都是老系统,这些老系统可能已经用了Storm,现在需要升级或者新增功能,这时候该怎么选呢?

4.1 先评估老系统的现状

首先要评估老系统的业务逻辑、数据量、性能要求、维护成本:

  • 如果老系统的业务逻辑非常简单,比如只是简单的消息过滤、转发,数据量不大,性能要求不高,那可以继续用Storm,不用改,改造成本太高。
  • 如果老系统的业务逻辑复杂,比如需要窗口处理、状态管理,或者数据量很大,性能要求高,或者经常出现丢数据、重复处理的问题,那就要考虑换成Flink。

4.2 改造的成本分析

改造的时候要考虑改造的成本,比如开发成本、测试成本、上线风险:

  • 如果老系统的代码量不大,业务逻辑简单,那可以直接重构,把Storm的代码改成Flink的代码,因为Flink的开发更简单,重构的成本不会太高。
  • 如果老系统的代码量很大,业务逻辑复杂,那可以采用灰度改造的方式,比如先把一部分功能改成Flink的,其他的继续用Storm,逐步替换,降低上线风险。

4.3 改造的示例:把Storm的单词计数改成Flink的

给大家举个简单的改造示例,比如把上面的Storm单词计数改成Flink的,改造的步骤非常简单:

  1. 把Storm的Spout改成Flink的SourceFunction,也就是数据源部分,逻辑差不多,只是代码结构变了。
  2. 把Storm的Bolt改成Flink的FlatMap、Map、KeyBy、Sum等算子,也就是加工部分,Flink的算子比Storm的Bolt更简单,不需要手动管理状态。
  3. 把Storm的拓扑改成Flink的流处理任务,也就是整体的结构,Flink的结构更清晰。

改造后的Flink代码就是上面的示例,代码量比Storm的少很多,而且功能更全,性能更好。

五、优缺点及注意事项

5.1 Storm的优缺点

优点:延迟低、部署简单、适合简单的实时处理场景。 缺点:没有原生的状态管理、窗口处理、批流统一,开发复杂,维护成本高,不能实现精准一次保证。 注意事项:如果用Storm,一定要做好状态的备份和恢复,避免丢数据;如果需要精准一次保证,一定要做好事务控制,开发成本很高。

5.2 Flink的优缺点

优点:原生支持精准一次保证、状态管理、窗口处理、批流统一,开发简单,维护成本低,功能全,性能好。 缺点:部署比Storm复杂,对运维的要求高,延迟比Storm稍微高一点(但差距很小)。 注意事项:如果用Flink,一定要做好检查点的配置,避免状态丢失;如果数据量很大,一定要做好状态的存储优化,比如用RocksDB存储状态;如果是老系统改造,一定要做好灰度测试,避免上线出问题。

六、文章总结

Storm和Flink都是非常优秀的实时处理框架,Storm主打低延迟,适合简单的实时处理场景;Flink主打精准一次,适合复杂的实时处理场景。在存量系统中,要根据业务逻辑、数据量、性能要求、改造成本来选择,如果是新系统,优先选Flink,如果是老系统,要评估改造的成本,逐步替换。