标题提到的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算法。
原理是:
- JobManager向所有Source发送Checkpoint Barrier(屏障)
- Source收到Barrier后,保存自己的状态,然后把Barrier发给下游
- 下游算子收到所有上游的Barrier后,保存自己的状态,然后把Barrier发给更下游
- 所有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. 故障恢复
当任务失败时:
- 停止所有算子
- 从最近的检查点恢复状态
- 重置Source的消费位置(比如Kafka的offset)
- 重新开始处理
这样,从检查点之后的数据会被重新处理,检查点之前的数据不会重复处理,保证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指标,显示反压的程度。
排查步骤:
- 找到反压最严重的算子
- 看这个算子的CPU、内存、输入输出量
- 分析是数据倾斜、慢算子,还是资源不足
- 针对性优化
十、Flink 1.16的预期特性
最后说说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作业,都能稳定运行,低延迟,高吞吐。
评论(0)
暂无评论,快来抢沙发~
评论功能仅对会员开放,请先登录
登录