Flink DataStream API 入门:理解一条数据流的旅程

摘要:本文梳理 Flink DataStream API 的核心主线——Source(数据从哪来)、Transformation(数据怎么变)、Sink(数据到哪去),重点讲清楚聚合算子、UDF、分流合流、输出算子背后的设计逻辑,而不是罗列 API 参数。目标是建立一套"遇到需求该用哪个工具"的判断框架,而不是试图记住每一行代码怎么写。本文整理自尚硅谷《大数据技术之Flink》课程。


写在前面:为什么看完容易忘

DataStream API 这一章内容量大、API 种类多,很多人学完的感受是"每个方法都看懂了,合起来却不知道该怎么用"。原因很简单:API 的具体写法本来就不需要死记,真正需要内化的是"这一类问题该用哪个工具解决"的判断力。

所以下面的内容,故意把重点放在"为什么这么设计"“什么场景该用它”,具体的方法签名只作为示例贴出,忘了随时可以查文档,但"该用 keyBy 还是 connect"这种判断力,才是这一章真正要留下的东西。


一、一个 Flink 程序的骨架

不管业务多复杂,一段 Flink DataStream 代码的结构都逃不出这三部分:

获取执行环境 → Source(读数据) → Transformation(处理数据) → Sink(写数据) → env.execute()

有一点容易被忽略但很关键:写完 Sink 不代表程序结束。Flink 是事件驱动、延迟执行的——调用 env.execute() 之前,前面所有代码只是在"画图纸"(构建执行计划),数据还没有真正流动一滴。这也是为什么初学者常犯的错误是"忘了调用 execute() 导致程序直接退出,什么输出都没有"。


二、Source:数据从哪来

Source 解决的是"数据源怎么接入"的问题。Flink 1.12 之后统一用 env.fromSource(...) 这个新架构,常见的几种来源:

  • 从集合读取(测试用)
  • 从文件读取
  • 从 Socket 读取(教学演示常用,生产环境基本不用)
  • 从 Kafka 读取(生产环境最常见,几乎是实时数仓的标准输入)
  • 数据生成器(压测、演示用)

这部分不需要花太多时间纠结细节,记住"生产场景基本是 Kafka Source"这一条结论就够了,具体连接参数用的时候查官方文档即可。


三、Transformation:数据怎么变,这是重点

这一章内容最多,容易看花眼。按"处理粒度"分类反而更好理解:

1. 单条数据处理:map / filter / flatMap

这三个是最基础的算子,分别对应"转换一条变一条"" 筛选符合条件的"“转换一条变多条”。它们不涉及"跟其他数据的关系",每条数据独立处理,理解起来最简单,也是后面所有复杂算子的基础形态。

2. 需要"记住之前数据"的处理:聚合算子

这是概念上第一个真正的分水岭。聚合的本质是:这次计算的结果,依赖之前处理过的数据,不再是单条数据独立处理了。

要做聚合,必须先做一件事:keyBy(按键分区)

KeyedStream<WaterSensor, String> keyedStream = stream.keyBy(e -> e.id);

keyBy 做的事情是:把同一个 key 的数据,分配到同一个逻辑分区里去处理。为什么必须先 keyBy 才能聚合?因为海量数据不可能全部塞到一台机器上做汇总,必须先按 key 分组、分区并行处理,这样才有意义,也才有效率。

keyBy 之后得到的不再是普通的 DataStream,而是 KeyedStream——记住这个类型转换,它是聚合类算子的入口。

KeyedStream 基础上,才能用:

  • sum() / max() / min():简单聚合,只更新指定字段,其余字段保留最初的值
  • maxBy() / minBy():和 max/min 类似,区别是返回整条数据,而不是只更新一个字段
  • reduce():更灵活的归约聚合,自己定义"新数据和已有聚合结果怎么合并"的逻辑

这里有个必须理解的坑:聚合算子会为每一个 key 保存一份"状态",而且这个状态永远不会自动清空。所以聚合算子只应该用在 key 的取值范围有限的场景(比如按传感器ID聚合,传感器数量是有限的),如果 key 是无限增长的(比如按订单号聚合),状态会无限膨胀,这是实际项目里最容易踩的坑之一。

3. 需要自定义逻辑:UDF(用户自定义函数)

UDF 分三种写法,本质上是同一件事的三种表达方式:

  • 函数类:实现 MapFunction/FilterFunction 等接口,写成一个独立的类,适合逻辑复杂、需要复用的场景
  • 匿名类/Lambda:图省事,逻辑简单、一次性用的场景
  • 富函数类(Rich Function):比如 RichMapFunction,多了 open()/close() 生命周期方法

富函数类是三种里唯一值得单独强调的——它有生命周期,open() 在算子真正开始处理数据前调用一次,close() 在结束时调用一次。这个特性是"资源要复用"这类需求的标准解法:比如要连接数据库或者外部存储做维度关联,连接对象应该在 open() 里创建一次,而不是每处理一条数据就创建一次连接——这是从"能跑"到"跑得好"的关键分水岭。

4. 控制数据往哪个并行任务发:物理分区算子

shuffle()(随机)、rebalance()(轮询)、rescale()(局部轮询)、broadcast()(广播到所有下游任务)、global()(全部发到第一个任务)、partitionCustom()(自定义规则)。

这类算子解决的是"数据在并行任务间怎么分配"的问题,和业务逻辑关系不大,更多是性能调优手段。日常开发用得最多的其实是 keyBy(它本质上也是一种分区),这几个显式的物理分区算子,记住"存在这几个工具,需要负载均衡或者广播的时候能想起来用"就够了。

5. 一条流拆成多条:分流

最朴素的做法是对同一条流反复调用 .filter(),但这样效率低——本质上是把原始数据复制了好几份,每份各自过滤一次。

更好的做法是用侧输出流(Side Output):一次遍历,根据条件把数据分别输出到不同的"标签"(OutputTag)里,主流之外的分支数据用 getSideOutput() 取出来。这是一次遍历完成多路分流,效率更高,也是实际项目中更常用的写法。

6. 多条流合成一条:合流

这里有两个层次,容易混淆,务必分清楚:

  • union(联合):只能合并数据类型完全相同的多条流,合并后就是简单的"数据都在一起了",不区分来源。适合"这几条流本质上是同一种数据,只是物理上分开了"的场景。
  • connect(连接):可以合并数据类型不同的两条流(注意只能两条,不像 union 可以传多个),但连接之后不能直接当一条流处理,必须用 CoMapFunction/CoProcessFunction 这类"各管各的"接口分别处理——map1() 处理第一条流的数据,map2() 处理第二条流的数据。

connect + keyBy 是一个非常重要的组合:两条流各自按同一个 key 做 keyBy 之后再 connect,就能保证"key 相同的数据会到同一个并行任务里去处理",这是实现"双流关联"(类似 SQL 的 join)的基础手法之一——比如一条是订单流、一条是支付流,按订单号 keyBy 之后 connect,在 CoProcessFunction 里用状态(比如 HashMap 缓存)分别记住两边来过的数据,匹配上了就输出,这就是简化版的双流 Join 逻辑。


四、Sink:数据到哪去

Sink 解决"计算结果往外部系统写"的问题。核心逻辑其实很统一:Flink 提供的官方连接器能覆盖大部分场景,真正需要自定义的时候不多

  • 输出到文件:FileSink,支持行编码/批量编码,按时间分桶,有滚动策略(按大小或时间切分文件)。有个细节:必须开启 Checkpoint,否则文件会一直停留在 .inprogress 状态不会真正落盘完成,这是很多人第一次用会踩的坑。
  • 输出到 Kafka:KafkaSink,如果要做到精确一次(Exactly-Once)写入,三个条件缺一不可——开启 Checkpoint、设置事务前缀、设置事务超时时间(要求介于 Checkpoint 间隔和 Kafka 最大 15 分钟之间)。这三条配置经常被漏掉一条导致语义达不到预期,需要对照检查。
  • 输出到 MySQL(JDBC):JdbcSink.sink() 需要四个参数——执行的 SQL、如何把数据填充到占位符里、执行选项(攒批大小、重试次数)、连接选项(地址账号密码)。这里的"攒批"(withBatchSize)是性能关键——不攒批、来一条写一条,数据库连接和事务开销会把吞吐量拖垮,这一点和写 ClickHouse 是同一个道理。
  • 自定义 Sink:当官方连接器覆盖不到的时候(比如要写入 HBase、ClickHouse 这类没有开箱即用连接器的系统),需要自己实现 SinkFunction 接口,重写 invoke() 方法。这也是很多实时数仓项目里"最后一公里"真正要动手写代码的地方——官方连接器解决了 90% 的常见存储,剩下 10% 靠自己实现。

五、这一章到底要留下什么

如果只能记住几条结论,应该是这些:

  1. 单条独立处理(map/filter) vs 需要历史数据(聚合),是这一章第一个重要的分界线,分界点就是 keyBy
  2. 聚合算子的状态永远不清空,只能用在 key 有限的场景,这是最容易踩的坑
  3. 富函数类的生命周期方法是解决"资源要复用而不是重复创建"这类需求的标准答案
  4. 分流用侧输出流,合流分 union(同类型)和 connect(不同类型,可结合 keyBy 做双流关联)
  5. Sink 的性能关键在"攒批",不管是写 MySQL 还是写其他外部存储,来一条写一条都是性能杀手

具体的方法名、参数顺序,查文档或者查自己以前写过的代码就行,不需要背。真正该问自己的是:“给我一个新需求,我知道该往这五条结论里的哪一条上靠”——如果这个判断力有了,这一章就算真的学明白了。

更多推荐