标题提到的Flink 1.16在本文写作时(2022年6月)尚未正式发布(预计2022年10月发布)。本文基于Flink 1.14/1.15的使用经验,结合1.16版本的预期特性,深入剖析Flink的底层原理。

我用Flink做实时计算已经好几年了,从1.9用到1.15。最开始只是写SQL和DataStream API,能跑通就行。后来遇到各种问题:反压、状态膨胀、检查点失败、数据延迟……才发现,不理解底层原理,根本解决不了这些问题。

本文深入剖析Flink的底层机制,包括流处理模型、状态管理、检查点机制、容错机制、内存管理、调度机制等。帮你从底层理解Flink,而不只是会写SQL。

一、Flink是什么

先简单介绍一下Flink。

1. 定义

Apache Flink是一个分布式流处理引擎,用于有界流和无界流的处理。

  • 有界流:有固定大小的数据集,比如历史数据、文件
  • 无界流:源源不断的数据,比如日志、消息队列

Flink的核心思想是"流是核心,批是流的特例"。也就是说,Flink把批处理也当作流处理来做,只是这个流是有界的。

2. Flink的特点

  • 低延迟:毫秒级延迟
  • 高吞吐:每秒能处理百万级数据
  • 精确一次(Exactly-Once):保证数据不丢不重
  • 事件时间(Event Time):支持基于事件发生时间的处理
  • 状态管理:内置状态管理,支持大状态
  • 容错机制:检查点机制,故障自动恢复

3. Flink vs Spark Streaming

很多人会拿Flink和Spark Streaming对比。

  • Spark Streaming是微批处理,把流切成小批次,延迟在秒级
  • Flink是真流处理,逐条处理,延迟在毫秒级
  • Flink的状态管理更完善,支持更大的状态
  • Flink的事件时间支持更好,处理乱序数据更优雅

当然,Spark也推出了Structured Streaming,在向真流处理靠拢。但目前来说,Flink在流处理领域还是领先的。

二、流处理模型

Flink的核心是流处理模型。理解了这个模型,就理解了Flink的一半。

1. 数据流和转换

Flink的基本编程模型是:

  • 数据源(Source):读取数据
  • 转换(Transformation):处理数据
  • 数据汇(Sink):输出结果
Source -> Transformation -> Sink

每个转换操作,都会把一个数据流变成另一个数据流。比如map、filter、keyBy、window等。

2. 并行数据流

Flink是分布式的,每个算子都可以并行运行。

一个数据流,会被分成多个并行的子流(Sub-Stream),每个子流由一个并行实例处理。

并行度(Parallelism)可以设置,默认是集群的CPU核数。

3. 数据分区

数据在算子之间传输时,需要分区。常见的分区策略:

  • Forward:一对一,上游一个子任务只发给下游一个子任务
  • Shuffle:随机分发,负载均衡
  • Rebalance:轮询分发,负载均衡
  • Rescale:局部轮询,只在同一个TaskManager内轮询
  • KeyBy:按key哈希分发,相同key到同一个子任务
  • Broadcast:广播,发给所有下游子任务

不同的分区策略,适用于不同的场景。比如聚合操作必须用KeyBy,广播变量用Broadcast。

三、时间和窗口

Flink的时间和窗口机制,是它区别于其他流处理框架的重要特性。

1. 三种时间

Flink支持三种时间:

  • 事件时间(Event Time):事件发生的时间,由数据中的时间戳决定
  • 摄入时间(Ingestion Time):数据进入Flink的时间
  • 处理时间(Processing Time):算子处理数据的时间

事件时间最重要,因为它能处理乱序数据和延迟数据。比如一条日志,10:00产生,10:05才到达Flink,用事件时间就能正确处理。

2. Watermark

用事件时间,就需要Watermark(水位线)。

Watermark是一个时间戳,表示"这个时间之前的数据都已经到了"。Flink根据Watermark来判断窗口是否可以触发计算。

比如,Watermark是10:05,就表示10:05之前的数据都到了,10:00-10:05的窗口可以触发计算了。

Watermark的生成策略:

  • 单调递增:数据时间戳单调递增,Watermark就是当前最大时间戳
  • 固定延迟:Watermark = 当前最大时间戳 - 延迟时间(比如允许5分钟乱序)

3. 窗口

窗口是把无界流切成有界的片段,然后在每个片段上做计算。

常见的窗口类型:

  • 滚动窗口(Tumbling):窗口不重叠,比如每5分钟一个窗口
  • 滑动窗口(Sliding):窗口可以重叠,比如每5分钟计算过去10分钟的数据
  • 会话窗口(Session):按会话间隔切分,超过间隔就开新窗口

窗口还可以按时间或按数量划分:

  • 时间窗口:按时间划分
  • 计数窗口:按数据条数划分

4. 窗口函数

窗口触发后,用窗口函数计算:

  • ReduceFunction:增量聚合,来一条算一条
  • AggregateFunction:增量聚合,比Reduce更灵活
  • ProcessWindowFunction:全量聚合,窗口触发后一次性计算所有数据

增量聚合性能好,但拿不到窗口的元信息;全量聚合性能差,但能拿到所有数据和窗口信息。

四、状态管理

状态管理是Flink的核心能力之一。

1. 什么是状态

状态,就是算子在处理数据时,需要保存的中间结果。

比如:

  • 聚合操作:需要保存当前的聚合值(sum、count等)
  • 窗口操作:需要保存窗口内的所有数据
  • 去重操作:需要保存已经见过的key
  • 机器学习:需要保存模型参数

Flink的状态,是分布式的,每个并行实例保存自己的那部分状态。

2. 状态的分类

Flink的状态分为两类:

  • Keyed State:按key分区的状态,每个key有自己的状态。只能在KeyedStream上使用。
  • Operator State:算子级别的状态,每个并行实例有自己的状态。不按key分区。

Keyed State更常用,包括:

  • ValueState:单个值
  • ListState:列表
  • MapState:键值对
  • ReducingState:聚合值
  • AggregatingState:聚合值(更灵活)

3. 状态后端

状态存在哪里?由状态后端(State Backend)决定。

Flink支持三种状态后端:

  • MemoryStateBackend:状态存在内存中,快但不稳定,适合测试
  • FsStateBackend:状态存在内存中,检查点存到文件系统,适合中等状态
  • RocksDBStateBackend:状态存在RocksDB中(本地磁盘),支持大状态,适合生产环境

生产环境,大状态用RocksDBStateBackend,小状态可以用FsStateBackend。

4. 状态TTL

状态可以设置TTL(Time To Live),过期自动清理。

比如,去重的状态,只需要保留最近1小时的数据,就可以设置TTL为1小时。

不设置TTL,状态会一直增长,最终导致状态膨胀,检查点变慢,甚至OOM。

五、检查点机制

检查点(Checkpoint)是Flink容错机制的核心。

1. 什么是检查点

检查点,就是把所有算子的状态,定期保存到持久化存储(比如HDFS、S3)。

当任务失败时,可以从最近的检查点恢复,保证数据不丢不重(Exactly-Once)。

2. 检查点的原理

Flink的检查点,基于Chandy-Lamport算法。

原理是:

  1. JobManager向所有Source发送Checkpoint Barrier(屏障)
  2. Source收到Barrier后,保存自己的状态,然后把Barrier发给下游
  3. 下游算子收到所有上游的Barrier后,保存自己的状态,然后把Barrier发给更下游
  4. 所有Sink都保存完状态后,检查点完成

Barrier是一种特殊的消息,它把数据流分成"检查点之前"和"检查点之后"两部分。

3. 对齐和非对齐

检查点有两种模式:

  • 对齐检查点(Aligned):下游算子要等所有上游的Barrier都到了,才保存状态。保证Exactly-Once,但会导致反压
  • 非对齐检查点(Unaligned):下游算子不等所有Barrier,先保存状态,把还没到的数据也存下来。减少反压,但检查点更大

Flink 1.11开始支持非对齐检查点,1.15之后越来越成熟。对于反压严重的作业,非对齐检查点很有用。

4. 检查点的配置

重要的配置参数:

  • interval:检查点间隔,比如1分钟
  • timeout:检查点超时时间
  • minPauseBetweenCheckpoints:两次检查点之间的最小间隔
  • tolerableCheckpointFailureNumber:允许失败的检查点次数
  • mode:Exactly-Once或At-Least-Once

六、容错机制

有了检查点,Flink就能实现容错。

1. 故障恢复

当任务失败时:

  1. 停止所有算子
  2. 从最近的检查点恢复状态
  3. 重置Source的消费位置(比如Kafka的offset)
  4. 重新开始处理

这样,从检查点之后的数据会被重新处理,检查点之前的数据不会重复处理,保证Exactly-Once。

2. 重启策略

Flink支持多种重启策略:

  • 固定延迟重启:失败后等待固定时间重启,最多重启N次
  • 失败率重启:在时间窗口内,失败率超过阈值就停止
  • 无重启:失败后直接停止

生产环境,一般用固定延迟重启,允许重启几次。

3. Savepoint

Savepoint(保存点)和Checkpoint类似,但它是手动触发的,用于版本升级、集群迁移等场景。

Savepoint和Checkpoint的区别:

  • Checkpoint是自动的,定期触发,用于故障恢复
  • Savepoint是手动的,按需触发,用于运维操作

Savepoint的数据格式更稳定,跨版本兼容性更好。

七、内存管理

Flink的内存管理,是它能处理大状态的关键。

1. 内存模型

Flink的TaskManager内存分为:

  • 堆内存(Heap):JVM堆,用于运行时对象
  • 堆外内存(Off-Heap):直接内存,用于网络缓冲、RocksDB等
  • 托管内存(Managed Memory):Flink管理的内存,用于排序、哈希表、RocksDB等

Flink不是直接用JVM堆存数据,而是用自己的内存管理器,把数据序列化后存在堆外内存。这样做的好处:

  • 减少GC压力
  • 内存利用率高
  • 支持大状态

2. 序列化

Flink有自己的序列化框架,比Java原生序列化快很多,占用空间也小。

Flink支持多种序列化:

  • POJO:自动识别,性能好
  • Avro:结构化数据
  • Kryo:通用序列化
  • 原生类型:int、long等,性能最好

3. RocksDB内存

用RocksDBStateBackend时,RocksDB会用堆外内存。

RocksDB的内存配置很重要,配置不好会导致OOM或性能差。重要参数:

  • blockCacheSize:块缓存大小
  • writeBufferNumber:写缓冲区数量
  • writeBufferSize:写缓冲区大小
  • maxBackgroundCompactions:后台compaction线程数

八、调度机制

Flink的调度机制,决定了任务怎么分配到节点上。

1. 调度模型

Flink 1.13之后,默认用Reactive Mode(响应式模式)。

在Reactive Mode下:

  • JobManager把作业拆成Task
  • TaskManager主动拉取Task
  • 资源变化时,自动重新分配

这种模式更弹性,支持动态扩缩容。

2. 任务链

Flink会把可以链接的算子链接在一起,形成Task Chain。

链接的好处:

  • 减少线程切换
  • 减少网络传输
  • 减少序列化/反序列化

默认情况下,Forward分区的算子会自动链接。可以用startNewChain()或disableChaining()控制。

3. 槽位共享

Flink的TaskManager有多个Slot(槽位),每个Slot可以运行一个并行实例的Pipeline。

默认情况下,不同算子的子任务可以共享一个Slot(Slot Sharing)。这样,一个Slot里就能运行完整的Pipeline,数据不需要跨节点传输。

九、反压处理

反压(Backpressure)是流处理中常见的问题。

1. 什么是反压

反压,就是下游处理不过来,向上游传递压力,让上游放慢速度。

比如,Sink写数据库太慢,上游的数据就会堆积,最终Source也会放慢。

反压本身是一种保护机制,防止数据堆积导致OOM。但反压会导致延迟增加,需要找到瓶颈并优化。

2. 反压的原因

常见的反压原因:

  • 数据倾斜:某个key的数据特别多,导致某个并行实例处理不过来
  • 慢算子:某个算子处理慢,比如复杂的正则、外部IO
  • 资源不足:CPU、内存、网络不够
  • 状态太大:状态操作慢,比如RocksDB的compaction

3. 反压的排查

Flink Web UI可以看反压情况。每个算子的Backpressure指标,显示反压的程度。

排查步骤:

  1. 找到反压最严重的算子
  2. 看这个算子的CPU、内存、输入输出量
  3. 分析是数据倾斜、慢算子,还是资源不足
  4. 针对性优化

最后说说Flink 1.16的预期特性(基于目前的release plan,实际以发布为准)。

1. 更完善的异步检查点

1.16会继续优化检查点机制,支持更高效的异步检查点,减少检查点对处理的影响。

2. 状态管理优化

  • RocksDB状态后端的性能优化
  • 状态TTL的更细粒度控制
  • 大状态的增量检查点优化

3. SQL能力增强

  • 更多的SQL函数支持
  • 窗口表函数(Window TVF)的完善
  • 批流一体的SQL体验优化

4. 运维工具增强

  • Web UI的优化
  • 更完善的监控指标
  • 作业迁移和升级工具

这些特性,会让Flink更稳定、更高效、更易用。

十一、写在最后

Flink是一个强大但复杂的流处理引擎。

要真正用好Flink,不能只停留在写SQL和API的层面,需要理解它的底层原理:流处理模型、时间和窗口、状态管理、检查点、容错、内存管理、调度、反压等。

本文从原理层面剖析了Flink的核心机制,希望能帮你建立对Flink的系统性理解。

Flink 1.16即将发布,会带来更多优化和新特性。但核心原理是不变的,理解了原理,就能快速适应新版本。

最后,用一句话总结:"Flink的精髓,是把流处理做到极致。理解了流、状态、检查点这三个核心,就理解了Flink的80%。"

愿大家的Flink作业,都能稳定运行,低延迟,高吞吐。