Oozie Spark2 Java Action:如何在Workflow终止时停止Spark应用
解决Oozie终止时Spark2集群模式任务无法停止的问题
我之前帮团队处理过类似的Spark+Oozie版本兼容问题,给你几个实用的解决方案,你可以根据自己的环境选择:
1. 结合Oozie回调与YARN API实现任务清理
因为你用的是Spark集群模式,驱动进程运行在YARN的ApplicationMaster中,直接杀Java Action容器不会影响到它。可以通过以下步骤关联Oozie工作流和Spark任务的生命周期:
- 在Java Action提交Spark任务时,给YARN应用添加关联Oozie工作流ID的标签,比如在spark-submit命令里加上:
这样YARN上的Spark应用会带上当前Oozie工作流的唯一标识。--conf spark.yarn.tags=oozie-wf-${wf:id()} - 在Oozie工作流中添加一个清理用的Shell Action,把它设置为
kill和fail节点的后续执行动作。Shell脚本里做两件事:- 通过YARN命令筛选出带指定标签的运行中应用:
APP_ID=$(yarn application -list -appStates RUNNING | grep "oozie-wf-${wf:id()}" | awk '{print $1}') - 如果找到应用ID,就执行终止命令:
if [ ! -z "$APP_ID" ]; then yarn application -kill $APP_ID fi
- 通过YARN命令筛选出带指定标签的运行中应用:
2. 让Spark驱动与Java Action进程绑定(客户端模式适用)
如果你的任务可以切换为Spark客户端模式,这个方法更直接:
- 修改Java Action的代码,不要通过Spark API直接初始化SparkSession,而是用
ProcessBuilder启动spark-submit命令,把Spark驱动作为Java进程的子进程运行。 - 添加JVM shutdown钩子,确保Java Action进程被Oozie终止时,能杀掉子进程:
这种方式下,Oozie终止Java Action的JVM时,shutdown钩子会触发,杀掉Spark驱动进程,进而终止整个Spark任务。public class SparkSubmitAction { public static void main(String[] args) throws Exception { ProcessBuilder pb = new ProcessBuilder( "spark-submit", "--class", "com.your.package.YourSparkJob", "--master", "yarn", "--deploy-mode", "client", "/path/to/your/spark-job.jar" ); pb.inheritIO(); Process sparkProcess = pb.start(); // 添加shutdown钩子,Oozie终止容器时触发 Runtime.getRuntime().addShutdownHook(new Thread(() -> { sparkProcess.destroyForcibly(); })); // 等待Spark任务完成 sparkProcess.waitFor(); } }
3. 升级Oozie版本(最彻底的方案)
Oozie 4.3.0及以上版本已经原生支持Spark2 Action,如果你所在的环境允许升级Oozie,这是一劳永逸的解决办法:
- 升级后直接在Oozie工作流中使用
spark2Action替代Java Action,Oozie会自动管理Spark任务的生命周期,工作流终止时会自动向YARN发送终止信号,杀掉对应的Spark应用。 - 注意:需要确认CDH 5.15.1与Oozie 4.3.0的兼容性,建议先在测试环境验证后再升级生产环境。
4. 监听Oozie工作流状态,调用Spark REST API终止任务
如果上述方法都不适用,你可以在Java Action中集成Oozie客户端API,实时监听工作流状态:
- 在Java代码中获取当前Oozie工作流的ID,通过Oozie客户端轮询工作流状态。
- 一旦检测到工作流状态变为
KILLED或FAILED,就调用Spark的REST API终止任务:
这种方式需要你维护Spark REST服务的可用性,并且处理好状态监听的线程逻辑,避免影响主任务执行。// 假设已经拿到Spark应用的submissionId(集群模式下可通过spark-submit的返回结果获取) String submissionId = "your-spark-submission-id"; URL killUrl = new URL("http://spark-master:6066/v1/submissions/kill/" + submissionId); HttpURLConnection conn = (HttpURLConnection) killUrl.openConnection(); conn.setRequestMethod("POST"); conn.getResponseCode();
内容的提问来源于stack exchange,提问作者sudharshan r
相关产品推荐
相关产品推荐

