Apache Flink 是一个分布式流批一体数据处理框架,核心思想是:把一切数据都当成“流”来处理(批数据只是有界流)。下面从整体到细节说明 Flink 如何处理数据。
一个 Flink 应用通常经历以下步骤:
从外部系统读取数据:
DataStream<String> stream = env.socketTextStream("localhost", 9999);
Flink 提供丰富的算子对数据进行处理:
常见算子:
map:一对一转换filter:过滤flatMap:一对多keyBy:按 key 分组reduce / aggregate:聚合window:窗口计算connect / union:流合并stream
.flatMap((line, out) -> {
for (String word : line.split(" ")) {
out.collect(word);
}
})
.keyBy(word -> word)
.window(TumblingEventTimeWindows.of(Time.seconds(5)))
.sum(1);
将处理结果写入外部系统:
result.print();
Source → Map → KeyBy → Window → Sink
env.setParallelism(4);
Flink 强大的地方在于状态管理:
状态可:
Flink 支持三种时间语义:
| 时间类型 | 说明 |
|---|---|
| Event Time | 数据产生时间(最常用) |
| Ingestion Time | 进入 Flink 时间 |
| Processing Time | 算子处理时间 |
env.enableCheckpointing(5000);
在 Flink 中:
DataStream APIDataSet API(旧) / DataStream(新,统一)Flink 1.12+ 倡导:流批一体
StreamExecutionEnvironment env =
StreamExecutionEnvironment.getExecutionEnvironment();
env.socketTextStream("localhost", 9999)
.flatMap((String line, Collector<String> out) -> {
for (String w : line.split(" ")) out.collect(w);
})
.keyBy(w -> w)
.window(TumblingProcessingTimeWindows.of(Time.seconds(5)))
.sum(1)
.print();
env.execute("WordCount");
Flink 通过 Source 读取数据 → 算子流式转换 → 状态 + 窗口计算 → Sink 输出,并利用 Checkpoint 实现高可靠、低延迟的分布式数据处理。
如果你愿意,我可以进一步讲:
你更想了解哪一部分?
免责声明:本站发布的内容(图片、视频和文字)以原创、转载和分享为主,文章观点不代表本网站立场,如果涉及侵权请联系站长邮箱:is@yisu.com进行举报,并提供相关证据,一经查实,将立刻删除涉嫌侵权内容。