Kubernetes部署Flink作业Checkpoint随机失败及吞吐量异常问题咨询
关于Async duration的含义
Flink Checkpoint的Async阶段是指算子完成同步快照(屏障对齐、状态时间点冻结)后,异步将状态数据持久化到远端Checkpoint存储的阶段,该阶段耗时占比高通常对应状态读写、序列化或存储侧存在IO瓶颈。
具体问题解答
问题1:为何首次运行时吞吐量极高,后续会骤降?
首次运行时状态为空且体量极小,RocksDB读写都在内存memtable中完成,无磁盘IO开销,算子处理延迟极低。随着运行时间增加,状态数据持续累积,memtable写满后触发flush落盘,后续读请求也需要访问磁盘上的SSTable文件甚至触发compaction操作,磁盘IO开销陡增;同时Checkpoint需要异步同步的状态数据量变大,抢占业务处理的CPU和IO资源,直接导致吞吐量骤降。问题2:为何首次运行仍出现Checkpoint失败?
你使用的Flink 1.13.2版本存在增量Checkpoint的已知缺陷:RocksDB的SSTable文件上传逻辑会和本地compaction逻辑冲突,当状态快速增长时,compaction会删除正在上传的SSTable文件,导致上传重试甚至失败。15分钟的Checkpoint间隔下,第三次Checkpoint时状态体量已经达到触发高频compaction的阈值,文件上传耗时超过10分钟的超时阈值就会触发失败。如果你的远端Checkpoint存储仍使用原有NFS,也可能存在上传带宽瓶颈。问题3:取消首次运行的作业后重新提交,为何无法恢复原有的高吞吐量?
如果重新提交作业是从上次失败的Checkpoint或Savepoint恢复,需要先把全量状态数据加载到RocksDB本地磁盘,大量随机读IO会直接打满磁盘带宽;且状态恢复后直接处于大体量状态的运行状态,memtable缓存命中率低,compaction持续运行,自然无法恢复初始空状态下的高吞吐量。就算未从快照恢复,只要你的聚合逻辑是按key全量累计,消费到之前已处理的偏移量位置时,状态体量仍会增长到之前的水平,吞吐量同样上不去。问题4:当前作业运行稳定无Checkpoint失败,且全部采用本地SSD部署,为何吞吐量仍然很低?
首先排查RocksDB配置合理性:Flink 1.13默认的RocksDB参数为通用配置,未针对大状态场景优化,比如write_buffer_size过小会导致频繁flush,level0_file_num_compaction_trigger阈值太低会触发频繁compaction,大量CPU和IO被RocksDB后台线程占用,业务处理线程拿不到足够资源。其次排查是否存在热点key,单个并行子任务的状态远超其他子任务,成为整个作业的性能瓶颈,调高整体并行度也无法解决。另外可检查Kafka消费/写入配置,确认是否存在消费限速、Sink侧Kafka分区不足导致写入限流的问题。问题5:为何Checkpoint大小波动极大,时而很高时而很低?
增量Checkpoint的大小取决于两次Checkpoint之间RocksDB新增的SSTable文件大小。当RocksDB触发compaction时,会将多个旧的SSTable文件合并成新的SSTable文件,增量Checkpoint会将这些合并后的新SSTable全部上传,不会重复上传旧文件,所以如果某次Checkpoint周期内刚好触发大规模compaction,这次的Checkpoint大小就会陡增;如果周期内没有compaction,只有新增的memtable flush出来的小SSTable,Checkpoint大小就会很小。另外如果你的业务数据有潮汐特性,不同时段的新增状态量差异大,也会导致Checkpoint大小波动。
优化建议
- 升级Flink版本到1.14及以上,修复了增量Checkpoint与RocksDB compaction冲突的多个已知缺陷
- 调整RocksDB配置:增大
write_buffer_size到64MB~128MB,调高level0_file_num_compaction_trigger到10,增大compaction线程数到4,开启RocksDB统计日志定位具体瓶颈 - 对热点key做拆分,或开启Flink本地聚合做预聚合,减少下游聚合算子的状态压力
- 如不需要永久保留历史状态,配置合理的状态TTL,自动清理过期状态,控制状态总体量
- 调整Checkpoint间隔到3~5分钟,减小单次Checkpoint需要上传的文件大小,降低超时概率
内容的提问来源于stack exchange,提问作者Abhay

