标题提到的Flink 1.16在本文写作时尚未正式发布(预计2022年下半年发布),本文基于Flink 1.14/1.15的使用经验,以及1.16预览版的测试体验,总结Flink使用中常见的坑和解决方法。
Flink是目前最流行的实时计算引擎之一。它功能强大,但也很复杂。用Flink的这几年,我踩了不少坑,有些坑让我熬了好几个通宵。
今天把这些坑整理出来,包括Checkpoint、反压、状态管理、内存配置、序列化、SQL等方面的问题。希望能帮你少走弯路,少熬夜。
一、Checkpoint相关的坑
Checkpoint是Flink容错机制的核心,也是最容易出问题的地方。
坑一:Checkpoint超时
现象: Checkpoint经常超时,任务频繁重启。
原因:
- 状态太大,Checkpoint时间长
- 反压导致Barrier对齐慢
- 磁盘IO慢,写Checkpoint慢
- 网络带宽不够,State Backend同步慢
解决方法:
- 增大Checkpoint超时时间(execution.checkpointing.timeout)
- 优化状态大小,清理不需要的状态
- 用RocksDB State Backend,支持增量Checkpoint
- 用分布式文件系统(HDFS、S3)存Checkpoint,不要用本地磁盘
- 排查反压问题(后面会讲)
- 开启Unaligned Checkpoint(Flink 1.11+支持),减少Barrier对齐时间
坑二:Checkpoint失败导致任务重启
现象: Checkpoint连续失败几次后,任务自动重启。
原因: Flink默认如果连续10次Checkpoint失败,就会重启任务。
解决方法:
- 先解决Checkpoint失败的根本原因
- 可以临时调大失败容忍次数(execution.checkpointing.tolerable-failed-checkpoints),但这只是治标
- 监控Checkpoint成功率,提前发现问题
坑三:从Savepoint恢复失败
现象: 从Savepoint恢复任务时,报错找不到状态。
原因:
- 改了算子的UID,导致状态无法匹配
- 改了状态的结构(比如ListState改成MapState)
- Savepoint版本和Flink版本不兼容
解决方法:
- 给每个算子设置明确的UID(.uid("xxx")),不要让Flink自动生成
- 改代码时,不要随便改UID
- 状态结构变更时,用State Processor API迁移状态
- 升级Flink版本前,先看兼容性说明
二、反压相关的坑
反压是Flink最常见的性能问题。
坑四:反压导致数据延迟
现象: Web UI上显示反压(Backpressure)为HIGH,数据处理延迟大。
原因: 下游算子处理速度跟不上上游发送速度,导致数据积压。
排查方法:
- 看Web UI的Backpressure页面,找到反压的根源
- 看每个算子的Busy时间,找到最慢的算子
- 看TaskManager的CPU、内存、磁盘IO、网络
- 看GC日志,是否频繁Full GC
解决方法:
- 增加慢算子的并行度
- 优化慢算子的逻辑(比如减少外部IO、优化算法)
- 增加TaskManager的资源
- 如果是数据倾斜,先解决数据倾斜
坑五:数据倾斜
现象: 某个SubTask的负载远高于其他SubTask,导致反压。
原因: Key分布不均匀,某些Key的数据量特别大。
解决方法:
- 加盐(Salt):给Key加随机前缀,把大Key拆成多个小Key
- 两阶段聚合:先局部聚合,再全局聚合
- 热点Key单独处理:把大Key拆出来,用单独的流处理
- 优化Key的选择,尽量选择分布均匀的Key
三、状态管理相关的坑
Flink的状态管理很强大,但也很容易出问题。
坑六:状态越来越大,OOM
现象: 任务运行一段时间后,状态越来越大,最后OOM。
原因:
- 没有设置状态的TTL,过期状态没有清理
- Key太多,每个Key都有状态
- 用了Heap State Backend,状态都在内存里
解决方法:
- 给状态设置TTL(StateTtlConfig),自动清理过期状态
- 用RocksDB State Backend,状态存在磁盘上,支持大状态
- 定期清理不需要的状态
- 监控状态大小,提前发现问题
坑七:RocksDB性能差
现象: 用了RocksDB State Backend后,性能比Heap差很多。
原因:
- RocksDB配置不合理
- 磁盘IO慢
- State访问模式不好(比如随机读写太多)
解决方法:
- 调优RocksDB配置:增大Block Cache、用SSD、开启Compaction优化
- 用本地SSD,不要用机械硬盘
- 优化State的访问模式,尽量顺序读写
- 考虑用增量Checkpoint,减少Checkpoint时间
四、内存配置相关的坑
Flink的内存模型比较复杂,配置不好很容易出问题。
坑八:TaskManager OOM
现象: TaskManager频繁OOM,被Yarn/K8s杀掉。
原因:
- 内存配置不合理,Heap太小
- 状态太大,用了Heap State Backend
- 网络缓冲区配置太大
- 直接内存(Direct Memory)溢出
解决方法:
- 合理配置TaskManager内存:Framework Heap、Task Heap、Managed Memory、Network Memory
- 大状态用RocksDB,不要用Heap
- 监控内存使用,提前调整配置
- 用Flink 1.10+的统一内存配置,更简单清晰
坑九:Managed Memory不够
现象: 用RocksDB时,Managed Memory不够,导致性能下降或OOM。
原因: Managed Memory是RocksDB和批处理用的内存,配置太小。
解决方法:
- 增大taskmanager.memory.managed.size
- 或者用taskmanager.memory.managed.fraction配置比例(默认0.4)
- 监控Managed Memory使用情况
五、序列化相关的坑
Flink的序列化对性能影响很大。
坑十:序列化性能差
现象: 序列化/反序列化占用大量CPU,吞吐量上不去。
原因:
- 用了Kryo序列化,性能比Flink原生序列化差
- POJO类没有符合规范,Flink无法用POJO序列化器
- 数据结构太复杂,序列化开销大
解决方法:
- 尽量用Flink原生支持的类型(基本类型、String、Tuple、POJO)
- POJO类要符合规范:public类、public无参构造、字段都是public或有getter/setter
- 避免用复杂的嵌套结构
- 必要时自定义序列化器
坑十一:POJO序列化器回退到Kryo
现象: 日志里显示"Class ... cannot be used as POJO class, falling back to Kryo"。
原因: POJO类不符合规范,Flink无法用POJO序列化器,回退到Kryo。
解决方法:
- 检查POJO类是否符合规范
- 字段不要用final
- 不要用内部类(除非是public static的)
- 父类的字段也要符合规范
六、Flink SQL相关的坑
Flink SQL越来越流行,但坑也不少。
坑十二:SQL任务性能差
现象: 同样的逻辑,SQL API比DataStream API性能差很多。
原因:
- SQL优化器没有生成最优的执行计划
- 没有正确配置MiniBatch、LocalGlobal等优化
- UDF性能差
解决方法:
- 开启MiniBatch聚合(table.exec.mini-batch.enabled)
- 开启LocalGlobal聚合(table.exec.mini-batch.enabled + table.optimizer.agg-phase-strategy)
- 优化UDF,避免复杂逻辑
- 用EXPLAIN查看执行计划,找到性能瓶颈
坑十三:Watermark问题
现象: 窗口不触发,或者数据延迟很大。
原因:
- Watermark生成策略不合理
- 数据乱序太严重
- 某个分区没有数据,导致Watermark不前进
解决方法:
- 合理设置Watermark的延迟时间
- 用WatermarkStrategy.withIdleness处理空闲分区
- 监控Watermark的推进情况
- 必要时用Allowed Lateness处理迟到数据
七、其他常见的坑
坑十四:并行度配置不合理
现象: 并行度太大,资源浪费;并行度太小,性能不够。
解决方法:
- 根据数据量和资源合理配置并行度
- 不同算子可以设置不同的并行度
- 用压测找到最优并行度
- 注意KeyBy后的并行度,避免数据倾斜
坑十五:Time Characteristic没设置
现象: 用了事件时间,但窗口不按预期触发。
原因: 没有设置Time Characteristic为EventTime。
解决方法:
- 在StreamExecutionEnvironment设置setStreamTimeCharacteristic(TimeCharacteristic.EventTime)
- Flink 1.12+默认就是EventTime,不需要手动设置
坑十六:重启策略配置不合理
现象: 任务出问题后,频繁重启,或者不重启。
解决方法:
- 合理配置重启策略:固定延迟、失败率、指数延迟
- 设置重启次数和延迟时间
- 监控重启次数,及时发现问题
八、我的排查经验
总结一下我排查Flink问题的经验。
1. 先看Web UI
Flink的Web UI信息很丰富:
- Overview:看任务状态、并行度、重启次数
- Backpressure:看反压情况
- Checkpoints:看Checkpoint成功率、耗时、状态大小
- Metrics:看各种指标,找到瓶颈
2. 看日志
TaskManager和JobManager的日志,是排查问题的关键。
- 看ERROR和WARN级别的日志
- 看异常堆栈,找到根本原因
- 注意日志的时间,对应问题发生的时间
3. 看GC日志
很多性能问题和OOM,都和GC有关。
- 开启GC日志
- 看Full GC的频率和耗时
- 如果频繁Full GC,调整内存配置或优化代码
4. 看资源监控
看TaskManager的CPU、内存、磁盘IO、网络。
- CPU高:可能是计算密集或序列化问题
- 内存高:可能是状态大或内存泄漏
- 磁盘IO高:可能是RocksDB或Checkpoint问题
- 网络高:可能是数据倾斜或Shuffle问题
九、写在最后
Flink是一个强大但复杂的工具。用Flink的过程,就是不断踩坑、不断学习的过程。
本文总结了我在Flink使用中遇到的常见坑和解决方法,包括Checkpoint、反压、状态管理、内存配置、序列化、SQL等方面。希望能帮你少走弯路。
当然,Flink的坑远不止这些。不同的业务场景,会遇到不同的问题。关键是要掌握排查问题的方法,遇到问题不慌,一步步排查。
2022年了,Flink已经更新到1.15,1.16也即将发布。新版本修复了很多问题,也增加了很多新功能。建议大家尽量用新版本,少踩老版本的坑。
最后,用一句话总结:"Flink的坑,踩过了才知道。但只要掌握了排查方法,就没有解决不了的问题。"
祝大家的Flink任务都能稳定运行,不用熬夜。
评论(0)
暂无评论,快来抢沙发~
评论功能仅对会员开放,请先登录
登录