Flink 1.15.4 JobManager类加载泄漏问题求助
Flink 1.15.4 JobManager元空间持续增长问题排查与求助
问题描述
在Flink 1.15.4版本中,每次上传Jar包执行任务后,JobManager的元空间(metaspace)持续增长,经排查是任务对应的类加载器无法被垃圾回收(GC)导致的。
最初怀疑问题源于org.apache.flink.metrics.jmx.JMXReporter(集群已配置metrics.reporter.jmx.class: org.apache.flink.metrics.jmx.JMXReporter)。堆转储分析显示:
org.apache.flink.jobmanager.job.numRestarts指标通过org.apache.flink.metrics.jmx.JMXReporter$JmxGauge持有org.apache.flink.runtime.jobgraph.JobGraph的强引用- JobGraph又关联了自定义输入类的实例
- 任务完成后,对应的MBean未被正确注销
后续进一步排查发现,堆中存在更多未被回收的JobGraph实例。
复现代码
以下是可复现该问题的示例代码:
package io.debug.flink; import org.apache.flink.api.common.io.GenericInputFormat; import org.apache.flink.api.common.io.NonParallelInput; import org.apache.flink.api.common.io.RichOutputFormat; import org.apache.flink.api.java.ExecutionEnvironment; import org.apache.flink.configuration.Configuration; import org.apache.flink.core.io.GenericInputSplit; import java.time.ZonedDateTime; import java.util.Iterator; import java.util.UUID; import java.util.stream.Collectors; import java.util.stream.IntStream; import lombok.extern.slf4j.Slf4j; public class CustomInputJob { public static void main(String[] args) throws Exception { ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment(); env.createInput(new CustomInput()) .output(new CustomOutput()); env.execute("CustomInputJob(export=" + ZonedDateTime.now() + ")"); } public static class CustomInput extends GenericInputFormat<UUID> implements NonParallelInput { private Iterator<UUID> iterator; @Override public void open(GenericInputSplit split) { iterator = IntStream.range(0, 5).boxed().map(i -> UUID.randomUUID()).collect(Collectors.toList()).iterator(); } @Override public boolean reachedEnd() { return !iterator.hasNext(); } @Override public UUID nextRecord(UUID reuse) { return iterator.next(); } } @Slf4j public static class CustomOutput extends RichOutputFormat<UUID> { @Override public void configure(Configuration parameters) { } @Override public void open(int taskNumber, int numTasks) { } @Override public void writeRecord(UUID record) { log.info("received: {}", record); } @Override public void close() { } } }
日志信息
开启org.apache.flink的TRACE日志后未发现异常,日志显示任务已完成清理:
Job 52f460a9dcb068b1509137f12f28b061 has been registered for cleanup in the JobResultStore after reaching a terminal state. Cleanup for the job '52f460a9dcb068b1509137f12f28b061' has finished. Job has been marked as clean.
堆转储分析
堆转储的相关信息如下:
- 未被回收的类加载器:多个任务的类加载器残留于堆中,未被GC回收
- 未被回收类加载器的引用路径:引用链最终指向自定义输入类实例,关联到JMXReporter的JmxGauge
numRestarts指标的引用:该指标的JmxGauge持有JobGraph的强引用,导致关联的自定义类实例无法被回收
尝试过的方案
添加了JobListener手动注销任务相关的MBean,代码如下:
@Override public void onJobExecuted(@Nullable JobExecutionResult jobExecutionResult, @Nullable Throwable throwable) { var server = ManagementFactory.getPlatformMBeanServer(); if (jobExecutionResult == null) { log.warn("received null job execution result"); return; } try { log.info("stared cleaning up the jmx resources for job id '{}', job name '{}'", jobExecutionResult.getJobID(), jobName); Hashtable<String, String> ht = new Hashtable<>(3); ht.put("host", "user-export-cluster-jobmanager"); ht.put("job_id", replaceInvalidChars(jobExecutionResult.getJobID().toString())); ht.put("job_name", replaceInvalidChars(jobName)); for (var matchingName : server.queryNames(ObjectName.getInstance("org.apache.flink.jobmanager.job.*", ht), null)) { if (server.isRegistered(matchingName)) { log.debug("unregistering MBean for `{}`", matchingName); server.unregisterMBean(matchingName); log.info("unregistered MBean for `{}`", matchingName); } else { log.debug("no registered MBean for `{}`", matchingName); } } } catch (Exception e) { log.error("can't get access to jmx", e); } }
但该方案无效,元空间仍持续增长,类加载器依旧无法被回收,剩余强引用显示仍有其他引用链指向自定义输入类。
疑问
- 这是否是Flink的Bug?
- 是否存在可调整的Flink配置,以强制任务完成后进行正确清理?
- 是否有其他可行的解决办法,比如类似注销MBean的方案?
内容的提问来源于stack exchange,提问作者Konstantin
相关产品推荐
相关产品推荐

