标题提到的Flink 1.16在本文写作时尚未正式发布(预计2022年下半年发布),本文基于Flink 1.14/1.15的使用经验,以及1.16预览版的测试体验,总结Flink使用中常见的坑和解决方法。

Flink是目前最流行的实时计算引擎之一。它功能强大,但也很复杂。用Flink的这几年,我踩了不少坑,有些坑让我熬了好几个通宵。

今天把这些坑整理出来,包括Checkpoint、反压、状态管理、内存配置、序列化、SQL等方面的问题。希望能帮你少走弯路,少熬夜。

一、Checkpoint相关的坑

Checkpoint是Flink容错机制的核心,也是最容易出问题的地方。

坑一:Checkpoint超时

现象: Checkpoint经常超时,任务频繁重启。

原因:

  1. 状态太大,Checkpoint时间长
  2. 反压导致Barrier对齐慢
  3. 磁盘IO慢,写Checkpoint慢
  4. 网络带宽不够,State Backend同步慢

解决方法:

  1. 增大Checkpoint超时时间(execution.checkpointing.timeout)
  2. 优化状态大小,清理不需要的状态
  3. 用RocksDB State Backend,支持增量Checkpoint
  4. 用分布式文件系统(HDFS、S3)存Checkpoint,不要用本地磁盘
  5. 排查反压问题(后面会讲)
  6. 开启Unaligned Checkpoint(Flink 1.11+支持),减少Barrier对齐时间

坑二:Checkpoint失败导致任务重启

现象: Checkpoint连续失败几次后,任务自动重启。

原因: Flink默认如果连续10次Checkpoint失败,就会重启任务。

解决方法:

  1. 先解决Checkpoint失败的根本原因
  2. 可以临时调大失败容忍次数(execution.checkpointing.tolerable-failed-checkpoints),但这只是治标
  3. 监控Checkpoint成功率,提前发现问题

坑三:从Savepoint恢复失败

现象: 从Savepoint恢复任务时,报错找不到状态。

原因:

  1. 改了算子的UID,导致状态无法匹配
  2. 改了状态的结构(比如ListState改成MapState)
  3. Savepoint版本和Flink版本不兼容

解决方法:

  1. 给每个算子设置明确的UID(.uid("xxx")),不要让Flink自动生成
  2. 改代码时,不要随便改UID
  3. 状态结构变更时,用State Processor API迁移状态
  4. 升级Flink版本前,先看兼容性说明

二、反压相关的坑

反压是Flink最常见的性能问题。

坑四:反压导致数据延迟

现象: Web UI上显示反压(Backpressure)为HIGH,数据处理延迟大。

原因: 下游算子处理速度跟不上上游发送速度,导致数据积压。

排查方法:

  1. 看Web UI的Backpressure页面,找到反压的根源
  2. 看每个算子的Busy时间,找到最慢的算子
  3. 看TaskManager的CPU、内存、磁盘IO、网络
  4. 看GC日志,是否频繁Full GC

解决方法:

  1. 增加慢算子的并行度
  2. 优化慢算子的逻辑(比如减少外部IO、优化算法)
  3. 增加TaskManager的资源
  4. 如果是数据倾斜,先解决数据倾斜

坑五:数据倾斜

现象: 某个SubTask的负载远高于其他SubTask,导致反压。

原因: Key分布不均匀,某些Key的数据量特别大。

解决方法:

  1. 加盐(Salt):给Key加随机前缀,把大Key拆成多个小Key
  2. 两阶段聚合:先局部聚合,再全局聚合
  3. 热点Key单独处理:把大Key拆出来,用单独的流处理
  4. 优化Key的选择,尽量选择分布均匀的Key

三、状态管理相关的坑

Flink的状态管理很强大,但也很容易出问题。

坑六:状态越来越大,OOM

现象: 任务运行一段时间后,状态越来越大,最后OOM。

原因:

  1. 没有设置状态的TTL,过期状态没有清理
  2. Key太多,每个Key都有状态
  3. 用了Heap State Backend,状态都在内存里

解决方法:

  1. 给状态设置TTL(StateTtlConfig),自动清理过期状态
  2. 用RocksDB State Backend,状态存在磁盘上,支持大状态
  3. 定期清理不需要的状态
  4. 监控状态大小,提前发现问题

坑七:RocksDB性能差

现象: 用了RocksDB State Backend后,性能比Heap差很多。

原因:

  1. RocksDB配置不合理
  2. 磁盘IO慢
  3. State访问模式不好(比如随机读写太多)

解决方法:

  1. 调优RocksDB配置:增大Block Cache、用SSD、开启Compaction优化
  2. 用本地SSD,不要用机械硬盘
  3. 优化State的访问模式,尽量顺序读写
  4. 考虑用增量Checkpoint,减少Checkpoint时间

四、内存配置相关的坑

Flink的内存模型比较复杂,配置不好很容易出问题。

坑八:TaskManager OOM

现象: TaskManager频繁OOM,被Yarn/K8s杀掉。

原因:

  1. 内存配置不合理,Heap太小
  2. 状态太大,用了Heap State Backend
  3. 网络缓冲区配置太大
  4. 直接内存(Direct Memory)溢出

解决方法:

  1. 合理配置TaskManager内存:Framework Heap、Task Heap、Managed Memory、Network Memory
  2. 大状态用RocksDB,不要用Heap
  3. 监控内存使用,提前调整配置
  4. 用Flink 1.10+的统一内存配置,更简单清晰

坑九:Managed Memory不够

现象: 用RocksDB时,Managed Memory不够,导致性能下降或OOM。

原因: Managed Memory是RocksDB和批处理用的内存,配置太小。

解决方法:

  1. 增大taskmanager.memory.managed.size
  2. 或者用taskmanager.memory.managed.fraction配置比例(默认0.4)
  3. 监控Managed Memory使用情况

五、序列化相关的坑

Flink的序列化对性能影响很大。

坑十:序列化性能差

现象: 序列化/反序列化占用大量CPU,吞吐量上不去。

原因:

  1. 用了Kryo序列化,性能比Flink原生序列化差
  2. POJO类没有符合规范,Flink无法用POJO序列化器
  3. 数据结构太复杂,序列化开销大

解决方法:

  1. 尽量用Flink原生支持的类型(基本类型、String、Tuple、POJO)
  2. POJO类要符合规范:public类、public无参构造、字段都是public或有getter/setter
  3. 避免用复杂的嵌套结构
  4. 必要时自定义序列化器

坑十一:POJO序列化器回退到Kryo

现象: 日志里显示"Class ... cannot be used as POJO class, falling back to Kryo"。

原因: POJO类不符合规范,Flink无法用POJO序列化器,回退到Kryo。

解决方法:

  1. 检查POJO类是否符合规范
  2. 字段不要用final
  3. 不要用内部类(除非是public static的)
  4. 父类的字段也要符合规范

Flink SQL越来越流行,但坑也不少。

坑十二:SQL任务性能差

现象: 同样的逻辑,SQL API比DataStream API性能差很多。

原因:

  1. SQL优化器没有生成最优的执行计划
  2. 没有正确配置MiniBatch、LocalGlobal等优化
  3. UDF性能差

解决方法:

  1. 开启MiniBatch聚合(table.exec.mini-batch.enabled)
  2. 开启LocalGlobal聚合(table.exec.mini-batch.enabled + table.optimizer.agg-phase-strategy)
  3. 优化UDF,避免复杂逻辑
  4. 用EXPLAIN查看执行计划,找到性能瓶颈

坑十三:Watermark问题

现象: 窗口不触发,或者数据延迟很大。

原因:

  1. Watermark生成策略不合理
  2. 数据乱序太严重
  3. 某个分区没有数据,导致Watermark不前进

解决方法:

  1. 合理设置Watermark的延迟时间
  2. 用WatermarkStrategy.withIdleness处理空闲分区
  3. 监控Watermark的推进情况
  4. 必要时用Allowed Lateness处理迟到数据

七、其他常见的坑

坑十四:并行度配置不合理

现象: 并行度太大,资源浪费;并行度太小,性能不够。

解决方法:

  1. 根据数据量和资源合理配置并行度
  2. 不同算子可以设置不同的并行度
  3. 用压测找到最优并行度
  4. 注意KeyBy后的并行度,避免数据倾斜

坑十五:Time Characteristic没设置

现象: 用了事件时间,但窗口不按预期触发。

原因: 没有设置Time Characteristic为EventTime。

解决方法:

  1. 在StreamExecutionEnvironment设置setStreamTimeCharacteristic(TimeCharacteristic.EventTime)
  2. Flink 1.12+默认就是EventTime,不需要手动设置

坑十六:重启策略配置不合理

现象: 任务出问题后,频繁重启,或者不重启。

解决方法:

  1. 合理配置重启策略:固定延迟、失败率、指数延迟
  2. 设置重启次数和延迟时间
  3. 监控重启次数,及时发现问题

八、我的排查经验

总结一下我排查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任务都能稳定运行,不用熬夜。