Spark cluster mode下如何在spark submit提交节点调用shell script
Spark cluster模式无修改调用现有shell脚本解决方案
可以在不修改shell脚本代码的前提下实现cluster模式正常调用,核心是对齐client模式与cluster模式的运行环境、资源分发差异。
为什么client模式可正常运行,cluster模式失败
client模式下Spark driver进程运行在提交任务的本地机器上,调用shell脚本的代码如果运行在driver侧,可直接访问本地已存在的脚本、配套依赖、本地路径权限,因此可以正常执行。
cluster模式下driver进程会被集群调度器分配到任意工作节点启动,默认存在三个核心差异导致执行失败:
- 目标工作节点不存在对应的shell脚本文件
- 脚本没有可执行权限,或缺少运行依赖的命令、附属配置文件
- 脚本输出的文件写入driver节点本地路径,无法被提交任务的本地机器获取
具体操作步骤
不需要修改shell脚本代码,仅调整Spark提交配置和Java侧调用逻辑即可:
- 提交Spark任务时通过
--files参数将shell脚本加入资源分发列表,Spark会自动把脚本分发到driver运行节点的工作目录:
如果脚本依赖多个附属文件/压缩包,可以用spark-submit --class your.main.Class --master yarn --deploy-mode cluster --files ./your_business_script.sh your_spark_project.jar--archives参数传入压缩包自动解压使用,格式为--archives your_deps.tgz#deps_dir - Java代码中调用脚本时,优先用
sh解释器直接执行,无需额外配置执行权限,也避免权限不足问题:import java.io.BufferedReader; import java.io.InputStreamReader; public class ShellExecutor { public void runScript() throws Exception { // 直接调用分发到当前工作目录的脚本 Process process = Runtime.getRuntime().exec(new String[]{"/bin/sh", "your_business_script.sh"}); // 必须读取脚本的标准输出和错误流,避免缓冲区占满导致脚本阻塞挂起 BufferedReader stdOut = new BufferedReader(new InputStreamReader(process.getInputStream())); BufferedReader stdErr = new BufferedReader(new InputStreamReader(process.getErrorStream())); String line; while ((line = stdOut.readLine()) != null) { // 可处理输出日志 } while ((line = stdErr.readLine()) != null) { // 可处理错误日志 } int exitCode = process.waitFor(); if (exitCode != 0) { throw new RuntimeException("Shell script execution failed, exit code: " + exitCode); } } } - 提前确认集群所有工作节点已安装脚本运行依赖的所有系统命令、工具,如果有自定义二进制依赖也可以一并通过
--files参数分发,脚本中用相对路径调用即可 - 脚本生成的输出文件如果需要回传到提交任务的本地,可在脚本执行完成后将文件上传到HDFS等集群共享存储,再从共享存储拉取到本地,不要直接写入driver节点的本地绝对路径
如果脚本调用逻辑在executor侧
如果你的shell脚本调用逻辑写在RDD/DataSet的转换算子中(运行在executor进程),只需在提交任务时加上配置--conf spark.executor.extraFiles=your_business_script.sh,或者在代码中调用sparkContext.addFile("your_business_script.sh"),即可将脚本分发到所有executor的工作目录,调用逻辑和上面一致。
内容的提问来源于stack exchange,提问作者mukesh dewangan
相关产品推荐
相关产品推荐

