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

Flink RichParallelSourceFunction的close/cancel函数在GCP DataProc强制终止时失效

问题场景

将Flink作业部署在GCP DataProc集群的工作流中,工作流设置了超时时间,到达时间后DAG被强制终止。此时自定义RichParallelSourceFunction(SampleSource)中的close和cancel方法完全没触发,日志里只有TaskExecutor相关组件的关闭信息,sf.close()根本没执行。

自定义Source代码

public class SampleSource extends RichParallelSourceFunction<Data> {

    private static final Logger logger = LoggerFactory.getLogger(SampleSource.class);
    private volatile boolean running;
    SourceFunction sf;

    @Override
    public void open(Configuration parameters) throws Exception {
        running = true;
        sf = new SourceFunction();
        sf.init();
    }

    @Override
    public void run(SourceContext<Data> ctx) throws Exception {
        while (running) {
            List<Data> data = sf.getData();
            data.forEach(ctx::collect);
            Thread.sleep(1000*10);
        }
    }

    @Override
    public void cancel() {
        try {
            logger.info("[Job Cancelled]");
            sf.updateTimestamp();
            this.close();
        } catch (Exception e) {
            logger.error(e.getMessage());
        }
    }

    @Override
    public void close() throws Exception {
        try {
            logger.info("[Job Closed]");
            running = false;
            sf.close();
            super.close();
        } catch (Exception e) {
            logger.error(e.getMessage());
        } finally {
            sf.updateTimestamp();
        }
    }
}

终止时的日志

2023-06-08 08:03:21,169 INFO state.TaskExecutorStateChangelogStoragesManager: Shutting down TaskExecutorStateChangelogStoragesManager.
2023-06-08 08:03:21,170 INFO state.TaskExecutorLocalStateStoresManager: Shutting down TaskExecutorLocalStateStoresManager.
2023-06-08 08:03:21,169 INFO blob.PermanentBlobCache: Shutting down BLOB cache
2023-06-08 08:03:21,170 INFO blob.TransientBlobCache: Shutting down BLOB cache
2023-06-08 08:03:21,170 INFO filecache.FileCache: removed file cache directory /tmp/flink-dist-cache-...
2023-06-08 08:03:21,173 INFO disk.FileChannelManagerImpl: FileChannelManager removed spill file directory /tmp/flink-netty-shuffle-...
解决方法
  • 修改工作流终止策略
    GCP DataProc超时后如果直接发SIGKILL杀进程,Flink完全没机会执行用户代码里的优雅关闭。得改成先给Flink JobManager发SIGTERM信号,触发作业的优雅取消流程,等30-60秒再强制终止,让框架有时间调用cancel和close方法。

  • 修复Source的生命周期逻辑
    不要在cancel里手动调用close(),Flink框架会在cancel执行完成后自动调用close,手动调用容易导致重复执行或者异常。另外,run方法里的Thread.sleep是阻塞操作,即使running设为false,也要等sleep结束才能退出,得用中断机制处理:

    @Override
    public void run(SourceContext<Data> ctx) throws Exception {
        try {
            while (running) {
                List<Data> data = sf.getData();
                data.forEach(ctx::collect);
                Thread.sleep(1000*10);
            }
        } catch (InterruptedException e) {
            running = false;
            logger.info("Source线程被中断,退出循环");
        }
    }
    
    @Override
    public void cancel() {
        try {
            logger.info("[Job Cancelled]");
            sf.updateTimestamp();
            running = false;
            // 中断run方法里的阻塞线程
            Thread.currentThread().interrupt();
        } catch (Exception e) {
            logger.error(e.getMessage());
        }
    }
    
  • 添加JVM关闭钩子
    万一还是遇到强制杀进程的情况,给JVM加个关闭钩子,确保资源能被清理:

    @Override
    public void open(Configuration parameters) throws Exception {
        running = true;
        sf = new SourceFunction();
        sf.init();
        // 注册关闭钩子
        Runtime.getRuntime().addShutdownHook(new Thread(() -> {
            try {
                logger.info("JVM关闭钩子触发,开始清理Source资源");
                sf.close();
                sf.updateTimestamp();
            } catch (Exception e) {
                logger.error("关闭钩子执行出错", e);
            }
        }));
    }
    
  • 检查Flink配置
    确认execution.attached设为false,避免JobManager和客户端绑定导致关闭逻辑失效;同时调整taskmanager.exit-on-out-of-memory等配置,防止因OOM直接杀死进程跳过关闭流程。

内容的提问来源于stack exchange,提问作者염경훈

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 22:30:41