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

如何在Shell脚本中判断Spark作业是否成功完成?

没问题,我来给你梳理下怎么实现这个需求——用Shell脚本提交Spark作业并根据执行结果分支处理,这在数据流水线里可太常见了!我分两种常用场景给你写具体实现,毕竟Spark的部署模式不同,处理方式也不一样~


场景1:Spark Client模式(直接等待作业完成)

在Client模式下,spark-submit命令会一直阻塞直到作业执行完毕,并且会返回对应的退出状态码(0代表成功,非0代表失败),这种情况处理起来最简单:

#!/bin/bash

# 定义Spark作业的完整提交命令
SPARK_SUBMIT_CMD="spark-submit \
  --class com.yourcompany.YourSparkMainClass \
  --master yarn \
  --deploy-mode client \
  --executor-memory 4G \
  --num-executors 5 \
  hdfs:///path/to/your/spark-job.jar \
  --input hdfs:///user/data/input \
  --output hdfs:///user/data/output"

# 提交作业并输出日志到文件(方便后续排查)
echo "启动Spark作业..."
$SPARK_SUBMIT_CMD > spark_job_exec.log 2>&1

# 捕获作业的退出状态码
SPARK_EXIT_CODE=$?

# 根据状态码分支处理
if [ $SPARK_EXIT_CODE -eq 0 ]; then
  echo "Spark作业执行成功!开始处理HDFS输出数据..."
  
  # 这里写你的后续HDFS操作,示例如下:
  # 1. 验证输出目录是否存在且有数据
  if hdfs dfs -test -d hdfs:///user/data/output && \
     [ $(hdfs dfs -count hdfs:///user/data/output | awk '{print $2}') -gt 0 ]; then
    # 2. 合并小文件(适合分区输出的场景)
    hdfs dfs -getmerge hdfs:///user/data/output hdfs:///user/data/merged/result.csv
    # 3. 修改输出目录权限,方便后续访问
    hdfs dfs -chmod -R 755 hdfs:///user/data/output
    echo "后续HDFS操作执行完成!"
  else
    echo "警告:输出目录存在但无数据,终止后续操作!"
    exit 1
  fi
else
  echo "Spark作业执行失败,退出码:$SPARK_EXIT_CODE"
  echo "可查看日志文件 spark_job_exec.log 排查问题"
  exit $SPARK_EXIT_CODE
fi

场景2:YARN Cluster模式(跟踪YARN应用状态)

如果用Cluster模式,默认情况下spark-submit提交后会立即返回(不会等待作业完成),这时候需要通过YARN命令来跟踪作业状态:

#!/bin/bash

# 定义Spark作业提交命令
SPARK_SUBMIT_CMD="spark-submit \
  --class com.yourcompany.YourSparkMainClass \
  --master yarn \
  --deploy-mode cluster \
  --executor-memory 4G \
  --num-executors 5 \
  hdfs:///path/to/your/spark-job.jar \
  --input hdfs:///user/data/input \
  --output hdfs:///user/data/output"

# 提交作业并提取YARN应用ID
echo "提交Spark作业到YARN..."
APP_ID=$($SPARK_SUBMIT_CMD 2>&1 | grep -oP 'application_\d+_\d+')

# 检查是否成功获取应用ID
if [ -z "$APP_ID" ]; then
  echo "提交失败:无法获取YARN应用ID,请检查Spark配置或Jar包路径!"
  exit 1
fi
echo "作业提交成功,YARN应用ID:$APP_ID"

# 等待作业执行完成
echo "等待作业执行完毕..."
yarn application -wait $APP_ID

# 获取作业最终状态
APP_STATUS=$(yarn application -status $APP_ID 2>&1 | grep -oP 'Final-State: \K\w+')

# 根据状态分支处理
if [ "$APP_STATUS" == "SUCCEEDED" ]; then
  echo "Spark作业执行成功!开始处理HDFS输出数据..."
  
  # 这里写你的后续HDFS操作,和Client模式的示例一致
  hdfs dfs -chmod -R 755 hdfs:///user/data/output
  echo "后续HDFS操作完成!"
else
  echo "Spark作业执行失败,最终状态:$APP_STATUS"
  echo "可执行命令查看日志:yarn logs -applicationId $APP_ID"
  exit 1
fi

关键细节与最佳实践
  • 日志捕获:把Spark作业的输出/错误日志重定向到文件(比如> spark_job_exec.log 2>&1),方便后续排查问题。
  • 幂等性设计:后续HDFS操作尽量保证幂等,比如执行前先检查输出目录是否存在,避免重复处理导致数据混乱。
  • 错误处理:不仅要判断Spark作业的状态,后续HDFS操作也要加入错误检查(比如用set -e让脚本在命令失败时自动退出,或者手动判断每个命令的退出码)。
  • YARN应用ID提取:不同Spark版本的spark-submit输出格式可能略有不同,如果grep正则匹配不到,可以调整正则表达式或者查看你的Spark版本输出格式。

内容的提问来源于stack exchange,提问作者Manoj Kumar Dhakad

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 07:05:35