每次savepoint执行结束后内存小幅上涨,累计产生整体内存增量是什么原因?
Flink Savepoint 执行后内存持续上涨问题排查与优化
常见诱因
- 未释放的savepoint临时状态缓存:savepoint执行过程中为了保证快照一致性,Flink会临时缓存状态块元数据、未写入完成的状态片段,若状态后端配置不当、异步快照线程池资源回收逻辑异常,这些缓存不会被完全释放,每次执行就会累计少量增量。
- 状态引用泄漏:若自定义状态序列化器、UDF中持有状态快照相关的对象引用,或是RocksDB的native内存分配出现泄漏,会导致每次savepoint后残留无法被GC回收的内存。
- 快照元数据残留:Flink默认会保留一定数量的已完成快照元数据,若配置的保留数量过大,存放在JobManager堆内存的这些元数据会随着savepoint执行次数增加逐步上涨。
- 堆外内存碎片:如果使用RocksDB作为状态后端,savepoint过程中产生的native内存碎片没有被及时整理,也会表现为进程内存占用逐步上涨。
排查步骤
- 先区分内存上涨区域:通过监控拆分JobManager堆内存、TaskManager堆内存、TaskManager堆外/native内存的涨幅,确认具体是哪部分内存出现累计上涨。
- 排查元数据占用:如果是JobManager堆内存上涨,核对
jobmanager.checkpoints.num-retained配置,查看当前保留的savepoint/checkpoint数量,验证元数据内存占用和保留数量的相关性。 - 排查GC情况:抓取TaskManager的GC日志,确认每次savepoint后Old区回收后的内存是否逐步上涨,如果是大概率为堆内存对象泄漏,可通过heap dump分析快照相关对象的引用链定位泄漏点。
- 排查native内存:如果使用RocksDB状态后端,开启RocksDB的native内存统计,对比savepoint前后的内存块分配、碎片率变化,确认是否为RocksDB层面的内存泄漏或碎片问题。
- 排查自定义逻辑:检查自定义状态序列化器、
CheckpointListener实现类中是否存在未释放的对象、缓存资源。
优化方案
- 调整快照保留策略:根据业务需要设置合理的savepoint保留数量,不需要永久保留的历史savepoint及时手动清理,避免元数据持续占用内存。
- 优化状态后端配置:使用RocksDB作为状态后端时,开启
state.backend.rocksdb.memory.managed配置托管RocksDB的内存分配,避免无限制的native内存上涨;使用Heap状态后端时,调整异步快照的线程池大小,限制临时缓存的最大占用。 - 定期重启作业:如果是内存碎片导致的缓慢上涨,可在业务低峰期定期从最近的savepoint重启作业,一次性释放累积的内存碎片和残留对象。
- 修复泄漏逻辑:如果排查到是UDF或自定义序列化器的引用泄漏,修复对应代码的资源释放逻辑,对实现了
CheckpointListener的类,在notifyComplete方法中及时清理临时缓存。
内容的提问来源于stack exchange,提问作者lumi
相关产品推荐
相关产品推荐

