Flink任务从Checkpoint重启时卡住问题求助
背景信息
我们使用Flink运行多个流处理任务,从Kafka读取数据、执行SQL转换后写入Kafka。任务部署在Kubernetes集群,含2个JobManager和多个TaskManager,采用RocksDB做Checkpoint存储,Checkpoint文件写入AWS S3存储桶。
近期将Flink从1.13.1升级到1.15.2,通过Savepoint机制迁移任务。两个Kubernetes集群迁移初期正常,但一段时间后(第一个集群约1个月,第二个2-3天)出现异常。
问题描述
部分任务启动失败:执行图中Source任务一直处于CREATED状态,下游ChangelogNormalize、Writer等任务已RUNNING。任务定期重启并抛出超时错误(堆栈信息简化):
java.lang.Exception: Cannot deploy task Source: source_consumer -> *anonymous_datastream_source$81*[211] (1/8) (de8f109e944dfa92d35cdc3f79f41e6f) - TaskManager (<address>) not responding after a rpcTimeout of 10000 ms at org.apache.flink.runtime.executiongraph.Execution.lambda$deploy$5(Execution.java:602) ... Caused by: java.util.concurrent.TimeoutException: Invocation of [RemoteRpcInvocation(TaskExecutorGateway.submitTask(TaskDeploymentDescriptor, JobMasterId, Time))] at recipient [akka.tcp://flink@<address>/user/rpc/taskmanager_0] timed out. This is usually caused by: 1) Akka failed sending the message silently, due to problems like oversized payload or serialization failures. In that case, you should find detailed error information in the logs. 2) The recipient needs more time for responding, due to problems like slow machines or network jitters. In that case, you can try to increase akka.ask.timeout. at org.apache.flink.runtime.jobmaster.RpcTaskManagerGateway.submitTask(RpcTaskManagerGateway.java:60) at org.apache.flink.runtime.executiongraph.Execution.lambda$deploy$4(Execution.java:580) at java.base/java.util.concurrent.CompletableFuture$AsyncSupply.run(Unknown Source) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Unknown Source) at java.base/java.util.concurrent.FutureTask.run(Unknown Source) at java.base/java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(Unknown Source) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(Unknown Source) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(Unknown Source) at java.base/java.lang.Thread.run(Unknown Source) Caused by: akka.pattern.AskTimeoutException: Ask timed out on [Actor[akka.tcp://flink@<address>/user/rpc/taskmanager_0#1723317240]] after [10000 ms]. Message of type [org.apache.flink.runtime.rpc.messages.RemoteRpcInvocation]. A typical reason for `AskTimeoutException` is that the recipient actor didn't send a reply.
JobManager日志中还出现超大消息提示:
Discarding oversized payload sent to Actor[akka.tcp://flink@<address>/user/rpc/taskmanager_0#1153219611]: max allowed size 10485760b bytes, actual size of encoded class org.apache.flink.runtime.rpc.messages.RemoteRpcInvocation was 83938405 bytes.
将akka.framesize设为100MB后,超时消失,任务进入INITIALIZING状态,但会长时间停滞,偶尔启动成功,偶尔抛出OOM错误:
java.lang.OutOfMemoryError: Java heap space
增加TaskManager内存仅对部分任务有效,部分任务还出现S3连接重置问题:
Caused by: org.apache.flink.runtime.state.BackendBuildingException: Failed when trying to restore operator state backend ... Caused by: java.lang.IllegalStateException: Connection pool shut down
最新发现(2023-02-08): 异常任务的Checkpoint中_metadata文件超大(最大168MB),且每次从Checkpoint恢复后该文件大小翻倍(重启后首次Checkpoint生成后大小稳定)。
问题解答
1. 提交任务时发送超大Akka消息的原因是什么?
核心原因是TaskDeploymentDescriptor(TDD)携带了过度膨胀的状态元数据。从_metadata文件翻倍的现象来看,Flink在恢复过程中,会将历史Checkpoint的元数据重复嵌入到新生成的Checkpoint元数据中,导致元数据体积指数级增长。当JobManager向TaskManager发送TDD时,这个超大元数据会被序列化到Akka RPC消息中,直接触发消息大小超限和RPC超时。
此外,RocksDB状态后端的元数据在恢复时未正确清理,或者跨版本Savepoint迁移时的元数据解析逻辑异常,也会导致TDD携带的状态信息过度膨胀。
2. Flink 1.13与1.15版本间的哪些变更可能导致这些问题?
- Checkpoint元数据存储逻辑变更:Flink 1.14+重构了Checkpoint元数据的序列化方式,增加了更详细的状态信息记录,但恢复流程中若处理不当,会导致元数据重复累加。
- SQL算子状态管理重构:1.15版本对SQL流处理的状态管理(如ChangelogNormalize算子)做了优化,部分算子的状态元数据在恢复时未正确去重或清理,引发元数据膨胀。
- Akka RPC机制升级:1.15版本升级了Akka依赖,对消息序列化、传输的校验逻辑更严格,原本1.13中可容忍的中等大小消息,在新版本中更容易触发超限或超时。
- Savepoint兼容性问题:1.13到1.15的跨版本Savepoint迁移,存在状态元数据解析的兼容性缺口,导致恢复时元数据被重复写入。
3. 如何排查堆内存占用过高的原因?
- 抓取并分析TaskManager堆快照:使用
jmap生成堆Dump文件,通过VisualVM或MAT工具分析内存占用TOP对象,重点关注CheckpointMetadata、TaskDeploymentDescriptor相关实例,确认是否为元数据堆积导致OOM。 - 开启Flink内存监控:启用Flink Metrics,跟踪TaskManager堆内存、非堆内存的实时变化趋势,重点观察任务启动、恢复阶段的内存突增节点。
- 解析Checkpoint元数据内容:使用Flink的
fsck工具或直接解析S3上的_metadata文件,检查是否存在重复的状态条目,定位元数据膨胀的具体来源。 - 开启DEBUG级日志:在TaskManager日志中开启DEBUG级别,重点查看状态恢复阶段的日志,确认是否存在重复加载状态元数据的行为。
- 对比空状态启动:尝试从空状态启动任务,对比内存占用情况,验证是否为Checkpoint恢复流程导致的内存问题。
内容的提问来源于stack exchange,提问作者C.S.

