别再死记硬背了!用Flink 1.18写一个词频统计,彻底搞懂批处理和流处理的区别
·
从词频统计实战看Flink批流一体本质:1.18版本API深度解析
第一次接触Flink时,很多人会被"批流一体"的概念弄得晕头转向——为什么同样的数据源既能当批处理又能当流处理?DataSet和DataStream到底该用哪个?本文将通过最经典的词频统计案例,用Flink 1.18最新API带你直击批流差异的核心。我们不会停留在表面语法对比,而是深入到执行引擎层面,看看同样的WordCount在批流模式下究竟经历了怎样不同的生命周期。
1. 环境准备与API演进
在Flink 1.18中,最显著的变化是DataSet API的全面弃用。这并非简单的API调整,而是标志着Flink彻底拥抱了"批是流的特例"这一设计哲学。我们先配置开发环境:
<!-- pom.xml关键依赖 -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java_2.12</artifactId>
<version>1.18.0</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-clients</artifactId>
<version>1.18.0</version>
</dependency>
注意:从1.12版本开始,Flink社区就推荐统一使用DataStream API处理所有场景,1.18版本则直接标记DataSet为
@Deprecated
2. 批处理模式的最后身影
让我们先用将被淘汰的DataSet API完成最后一次批处理词频统计,这有助于理解传统批处理的思维模式:
public class LegacyBatchWordCount {
public static void main(String[] args) throws Exception {
ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment();
DataSource<String> text = env.readTextFile("input.txt");
text.flatMap((line, out) -> {
for (String word : line.split("\\s")) {
out.collect(new Tuple2<>(word, 1));
}
}).returns(Types.TUPLE(Types.STRING, Types.INT))
.groupBy(0)
.sum(1)
.print();
}
}
批处理的核心特征:
- 有界数据集:处理前已知所有输入数据
- 全量计算:每次作业处理完整数据集
- 最终一致性:只有最终结果有意义
3. 流处理模式的本质突破
同样的词频统计改用DataStream API实现:
public class StreamingWordCount {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
DataStream<String> text = env.readTextFile("input.txt");
text.flatMap((line, out) -> {
for (String word : line.split("\\s")) {
out.collect(new Tuple2<>(word, 1));
}
}).returns(Types.TUPLE(Types.STRING, Types.INT))
.keyBy(0)
.sum(1)
.print();
env.execute("Streaming WordCount");
}
}
流处理的关键差异:
- 无界数据流:理论上输入永无止境
- 增量计算:来一条处理一条
- 持续更新:结果表不断随时间演进
4. 执行模型深度对比
4.1 批处理执行流程
graph LR
A[读取完整文件] --> B[全量分词]
B --> C[全局分组]
C --> D[聚合计算]
D --> E[输出最终结果]
4.2 流处理执行流程
graph LR
A[持续监控文件] --> B[逐行事件触发]
B --> C[键分区状态更新]
C --> D[持续输出更新]
关键发现:在Flink内部,批处理被当作一种特殊的流处理——有界流(Bounded Stream)
5. 批流一体实战技巧
Flink 1.18推荐使用RuntimeExecutionMode统一处理批流场景:
public class UnifiedWordCount {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 通过参数控制批流模式
env.setRuntimeMode(RuntimeExecutionMode.BATCH);
env.readTextFile("input.txt")
.flatMap(new Tokenizer())
.keyBy(0)
.sum(1)
.print();
env.execute("Unified Processing");
}
public static class Tokenizer implements FlatMapFunction<String, Tuple2<String, Integer>> {
@Override
public void flatMap(String value, Collector<Tuple2<String, Integer>> out) {
for (String word : value.split("\\s")) {
out.collect(new Tuple2<>(word, 1));
}
}
}
}
三种运行时模式对比:
| 模式 | 触发条件 | 状态管理 | 输出特性 |
|---|---|---|---|
| STREAMING | 默认模式 | 持续维护 | 增量更新 |
| BATCH | 明确设置 | 阶段释放 | 最终输出 |
| AUTOMATIC | 根据输入源推断 | 智能切换 | 自适应 |
6. 生产环境部署建议
对于刚接触Flink的团队,建议从以下路径逐步深入:
- 开发阶段:直接使用
getExecutionEnvironment()自动识别环境 - 本地测试:通过
createLocalEnvironmentWithWebUI()启用调试界面 - 生产部署:
# 以流模式运行 flink run -d -c com.YourClass yourJob.jar # 以批模式运行 flink run -d -Dexecution.runtime-mode=BATCH -c com.YourClass yourJob.jar
常见部署架构选择:
- 独立集群:适合小规模场景,部署简单
- YARN/K8S:需要资源动态扩展时的选择
- 会话模式vs应用模式:根据作业隔离需求决定
7. 调试与优化经验
在真实项目中处理词频统计时,有几个容易踩的坑:
- 并行度设置:流处理中
keyBy后的算子并行度需要保持一致 - 时间语义:流处理中注意
EventTime和ProcessingTime的区别 - 状态后端:大状态作业推荐配置RocksDBStateBackend
- 检查点配置:对于关键业务需要设置合理的检查点间隔
// 典型的生产配置示例
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.enableCheckpointing(5000); // 5秒检查点
env.setStateBackend(new RocksDBStateBackend("hdfs://checkpoints"));
env.getCheckpointConfig().setMinPauseBetweenCheckpoints(1000);
8. 从词频统计到复杂应用
掌握了基础词频统计后,可以逐步扩展到真实场景:
- 实时ETL:对接Kafka进行数据清洗
- 事件驱动:处理点击流、IoT设备数据
- 状态计算:实现会话窗口、用户行为分析
- 机器学习:结合Alink库进行实时预测
// 进阶示例:带窗口的流式词频
text.flatMap(new Tokenizer())
.keyBy(0)
.window(TumblingEventTimeWindows.of(Time.seconds(10)))
.sum(1)
.addSink(new ElasticsearchSink<>());
最终记住:批流差异不在于API表面,而在于对时间本质的理解。当你能用流处理思维处理所有场景时,才算真正掌握了Flink的精髓。
更多推荐
所有评论(0)