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,提问作者염경훈
相关产品推荐
相关产品推荐

