Flink S3 Sink偶发RemoteTransportException致集群重启问题排查
Flink S3 Sink触发TaskManager失联导致集群重启问题
问题现象
使用Flink S3 Sink时,偶尔会出现TaskManager失联异常,触发集群重启,异常日志如下:
Task - - Sink ConsumptionWriter (207/250)#0 (5a254b0622e8571f30567c39872691d1) switched from RUNNING to FAILED with failure cause: org.apache.flink.runtime.io.network.netty.exception.RemoteTransportException: Lost connection to task manager 'ip-10-0-121-109.us-west-2.compute.internal/10.0.121.109:39219'. This indicates that the remote task manager was lost.
异常发生前,TaskManager内存监控日志显示Direct Memory已超出总容量,怀疑内存溢出是TaskManager挂掉的直接原因:
TaskManagerRunner - - Memory usage stats: [HEAP: 109905/206464/206464 MB, NON HEAP: 310/329/496 MB (used/committed/max)] 2022-09-01 17:08:08,422 INFO TaskManagerRunner - - Direct memory stats: Count: 1311709, Total Capacity: 43023060620, Used Memory: 43023060624 2022-09-01 17:08:08,422 INFO TaskManagerRunner - - Off-heap pool stats: [Code Cache: 123/124/240 MB (used/committed/max)], [Metaspace: 187/205/256 MB (used/committed/max)] 2022-09-01 17:08:08,422 INFO TaskManagerRunner - - Garbage collector stats: [G1 Young Generation, GC TIME (ms): 176063, GC COUNT: 2118], [G1 Old Generation, GC TIME (ms): 0, GC COUNT: 0] TaskManagerRunner - - Memory usage stats: [HEAP: 111441/206464/206464 MB, NON HEAP: 310/329/496 MB (used/committed/max)] 2022-09-01 17:08:18,423 INFO TaskManagerRunner - - Direct memory stats: Count: 1311709, Total Capacity: 43023061644, Used Memory: 43023061648 2022-09-01 17:08:18,423 INFO TaskManagerRunner - - Off-heap pool stats: [Code Cache: 123/124/240 MB (used/committed/max)], [Metaspace: 187/205/256 MB (used/committed/max)] 2022-09-01 17:08:18,423 INFO TaskManagerRunner - - Garbage collector stats: [G1 Young Generation, GC TIME (ms): 176063, GC COUNT: 2118], [G1 Old Generation, GC TIME (ms): 0, GC COUNT: 0]
集群与作业配置
节点规格
使用r5d.24xlarge云实例:
r5d.24xlarge 96 vCore, 768 GiB memory, 3600 SSD GB storage EBS Storage:none
Flink作业提交命令
flink run \ -m yarn-cluster \ -d \ -yjm 491000 \ -ytm 485000 \ -yD taskmanager.numberOfTaskSlots=40 \ -yD jobmanager.io-pool.size=64 \ -yD jobmanager.future-pool.size=64 \ -yD jobmanager.memory.heap.size=460g \ -yD jobmanager.memory.jvm-overhead.max=40g \ -yD taskmanager.memory.process.size=440g \ -yD taskmanager.memory.jvm-overhead.max=50g \ -yD taskmanager.memory.network.max=50g \ -yD taskmanager.memory.network.min=40g \ -yD taskmanager.memory.managed.size=150g \ -yD taskmanager.memory.task.off-heap.size=4g \ -yD web.timeout=120000 \ -yD akka.ask.timeout=600s \ -yD heartbeat.timeout=600000 \ -yD rest.bind-port=8081 \ -yD akka.framesize="25006612b" \ -yD cluster.evenly-spread-out-slots=true \ -yD taskmanager.network.memory.buffer-debloat.enabled=true \ -yD state.checkpoints.dir=hdfs://ha-nn-uri/flink/checkpoints \ -yD state.savepoints.dir=hdfs://ha-nn-uri/flink/savepoints \ -yD high-availability.zookeeper.client.session-timeout=300000 \ -yD high-availability.zookeeper.client.connection-timeout=60000 \ -yD high-availability.zookeeper.client.max-retry-attempts=10 \ -yD high-availability.zookeeper.client.tolerate-suspended-connections=true \ -yD state.backend.local-recovery=true \ -yD taskmanager.debug.memory.log=true \ -yD taskmanager.debug.memory.log-interval=10000 \ -yD env.java.opts="-XX:+UnlockExperimentalVMOptions -Xloggc:/tmp/taskman.gc.log -XX:+PrintGCApplicationStoppedTime -XX:+PrintGCDetails -XX:+PrintGCDateStamps -XX:+UseGCLogFileRotation -XX:NumberOfGCLogFiles=10 -XX:GCLogFileSize=10M -XX:+PrintPromotionFailure -XX:+PrintGCCause -XX:+UseG1GC" \ $SAVEPOINT_CMD ATVPlaybackStateMachineFlinkJob-1.0-super-1.2.3.jar --stage "prod" --cell-name "StreamProcessor-Cell1-us-east-1"
问题分析与解决方案
核心原因
Direct Memory超限是TaskManager失联的关键诱因:S3 Sink依赖的AWS SDK会大量使用Direct Memory处理上传缓冲区,当内存占用超出JVM分配的Direct Memory上限时,会触发OOM或被系统内核的OOM Killer强制终止,进而导致TaskManager与JobManager失联。
具体优化措施
显式限制Direct Memory大小
在env.java.opts中添加-XX:MaxDirectMemorySize=64g(可根据业务调整至60-80g),避免无限制占用内存。当前节点总内存768G,TaskManager配置的process.size为440g,预留足够Direct Memory空间可避免溢出。优化S3 Sink参数
- 降低批量写入阈值:调整
batch.size或bucket.checkpoint.interval,减少单次上传的内存占用; - 启用异步上传:确保使用支持异步写入的S3 Sink实现,避免同步上传阻塞内存释放;
- 调整AWS SDK缓冲区:通过SDK配置减少单次请求的缓冲区大小,降低Direct Memory消耗。
- 降低批量写入阈值:调整
调整Flink内存分配
- 下调
taskmanager.memory.managed.size:当前配置150g,若状态后端为RocksDB,Managed Memory主要用于排序/哈希表,若业务状态不大可降至80-100g,释放内存给Direct Memory; - 压缩JVM Overhead:
taskmanager.memory.jvm-overhead.max=50g过高,调整至20-30g即可满足需求,节省内存空间。
- 下调
GC调优
当前Young GC次数过多(2118次)、耗时过长(176秒),易引发心跳超时:- 添加
-XX:G1NewSizePercent=20和-XX:G1MaxNewSizePercent=40,扩大Young区,减少Minor GC频率; - 添加
-XX:ParallelGCThreads=32和-XX:ConcGCThreads=8,匹配96核的CPU规格,提升GC效率。
- 添加
系统层面配置
检查节点OOM Killer规则,调整/etc/security/limits.conf中的内存限制,或降低TaskManager进程的OOM优先级,避免被内核优先终止。
内容的提问来源于stack exchange,提问作者priyadhingra19
相关产品推荐
相关产品推荐

