You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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);
    }
}

但该方案无效,元空间仍持续增长,类加载器依旧无法被回收,剩余强引用显示仍有其他引用链指向自定义输入类。

疑问

  1. 这是否是Flink的Bug?
  2. 是否存在可调整的Flink配置,以强制任务完成后进行正确清理?
  3. 是否有其他可行的解决办法,比如类似注销MBean的方案?

内容的提问来源于stack exchange,提问作者Konstantin

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.10 17:17:28