一、先搞懂两个工具的基础定位
很多做实时数据处理的开发者,经常会在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的,改造的步骤非常简单:
- 把Storm的Spout改成Flink的SourceFunction,也就是数据源部分,逻辑差不多,只是代码结构变了。
- 把Storm的Bolt改成Flink的FlatMap、Map、KeyBy、Sum等算子,也就是加工部分,Flink的算子比Storm的Bolt更简单,不需要手动管理状态。
- 把Storm的拓扑改成Flink的流处理任务,也就是整体的结构,Flink的结构更清晰。
改造后的Flink代码就是上面的示例,代码量比Storm的少很多,而且功能更全,性能更好。
五、优缺点及注意事项
5.1 Storm的优缺点
优点:延迟低、部署简单、适合简单的实时处理场景。 缺点:没有原生的状态管理、窗口处理、批流统一,开发复杂,维护成本高,不能实现精准一次保证。 注意事项:如果用Storm,一定要做好状态的备份和恢复,避免丢数据;如果需要精准一次保证,一定要做好事务控制,开发成本很高。
5.2 Flink的优缺点
优点:原生支持精准一次保证、状态管理、窗口处理、批流统一,开发简单,维护成本低,功能全,性能好。 缺点:部署比Storm复杂,对运维的要求高,延迟比Storm稍微高一点(但差距很小)。 注意事项:如果用Flink,一定要做好检查点的配置,避免状态丢失;如果数据量很大,一定要做好状态的存储优化,比如用RocksDB存储状态;如果是老系统改造,一定要做好灰度测试,避免上线出问题。
六、文章总结
Storm和Flink都是非常优秀的实时处理框架,Storm主打低延迟,适合简单的实时处理场景;Flink主打精准一次,适合复杂的实时处理场景。在存量系统中,要根据业务逻辑、数据量、性能要求、改造成本来选择,如果是新系统,优先选Flink,如果是老系统,要评估改造的成本,逐步替换。
Comments