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

Airflow任务报错仍标记DAG成功问题求助

问题

使用BashOperator触发服务器上Spark任务的Airflow DAG,该Spark任务读取按日分区的S3 Bucket数据并执行操作。当Bucket无对应分区数据时,Spark会抛出'path does not exist'的ERROR级异常,但Airflow将所有Spark日志都以INFO级别打印,导致任务实际报错时,Airflow仍将DAG标记为成功运行。相关日志如下:

[2022-11-17, 08:46:37 IST] {subprocess.py:92} INFO - 2022-11-17, 08:46:37 IST [main] WARN  org.apache.hadoop.util.NativeCodeLoader - Unable to load native-hadoop library for your platform... using builtin-java classes where applicable
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO - 2022-11-17, 08:46:41 IST [main] ERROR com.newjs.preprocessors.DeletedProfileContactsEligibleForRetPreProcessor - Error in Incremental Data Preprocessor
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO - org.apache.spark.sql.AnalysisException: Path does not exist;
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.sql.execution.datasources.DataSource$$anonfun$org$apache$spark$sql$execution$datasources$DataSource$$checkAndGlobPathIfNecessary$1.apply(DataSource.scala:558)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.sql.execution.datasources.DataSource$$anonfun$org$apache$spark$sql$execution$datasources$DataSource$$checkAndGlobPathIfNecessary$1.apply(DataSource.scala:545)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at scala.collection.TraversableLike$$anonfun$flatMap$1.apply(TraversableLike.scala:241)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at scala.collection.TraversableLike$$anonfun$flatMap$1.apply(TraversableLike.scala:241)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at scala.collection.immutable.List.foreach(List.scala:392)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at scala.collection.TraversableLike$class.flatMap(TraversableLike.scala:241)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at scala.collection.immutable.List.flatMap(List.scala:355)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.sql.execution.datasources.DataSource.org$apache$spark$sql$execution$datasources$DataSource$$checkAndGlobPathIfNecessary(DataSource.scala:545)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.sql.execution.datasources.DataSource.resolveRelation(DataSource.scala:359)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.sql.DataFrameReader.loadV1Source(DataFrameReader.scala:223)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:211)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.sql.DataFrameReader.parquet(DataFrameReader.scala:644)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.sql.DataFrameReader.parquet(DataFrameReader.scala:643)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at com.js.utils.ReadUtils.readParquetFromS3(ReadUtils.java:24)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at com.js.newjs.preprocessors.DeletedProfileContactsEligibleForRetPreProcessor.incrementalDataProcess(DeletedProfileContactsEligibleForRetPreProcessor.java:91)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at com.js.JsETLSparkTransformationsApplication.process(JsETLSparkTransformationsApplication.java:91)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at com.js.JsETLSparkTransformationsApplication.main(JsETLSparkTransformationsApplication.java:32)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at java.lang.reflect.Method.invoke(Method.java:498)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.deploy.JavaMainApplication.start(SparkApplication.scala:52)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.deploy.SparkSubmit.org$apache$spark$deploy$SparkSubmit$$runMain(SparkSubmit.scala:845)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.deploy.SparkSubmit.doRunMain$1(SparkSubmit.scala:161)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.deploy.SparkSubmit.submit(SparkSubmit.scala:184)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.deploy.SparkSubmit.doSubmit(SparkSubmit.scala:86)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.deploy.SparkSubmit$$anon$2.doSubmit(SparkSubmit.scala:920)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.deploy.SparkSubmit$.main(SparkSubmit.scala:929)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.deploy.SparkSubmit.main(SparkSubmit.scala)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO - Exception in thread "main" org.apache.spark.sql.AnalysisException: Path does not exist;
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.sql.execution.datasources.DataSource$$anonfun$org$apache$spark$sql$execution$datasources$DataSource$$checkAndGlobPathIfNecessary$1.apply(DataSource.scala:558)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.sql.execution.datasources.DataSource$$anonfun$org$apache$spark$sql$execution$datasources$DataSource$$checkAndGlobPathIfNecessary$1.apply(DataSource.scala:545)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at scala.collection.TraversableLike$$anonfun$flatMap$1.apply(TraversableLike.scala:241)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at scala.collection.TraversableLike$$anonfun$flatMap$1.apply(TraversableLike.scala:241)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at scala.collection.immutable.List.foreach(List.scala:392)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at scala.collection.TraversableLike$class.flatMap(TraversableLike.scala:241)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at scala.collection.immutable.List.flatMap(List.scala:355)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.sql.execution.datasources.DataSource.org$apache$spark$sql$execution$datasources$DataSource$$checkAndGlobPathIfNecessary(DataSource.scala:545)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.sql.execution.datasources.DataSource.resolveRelation(DataSource.scala:359)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.sql.DataFrameReader.loadV1Source(DataFrameReader.scala:223)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:211)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.sql.DataFrameReader.parquet(DataFrameReader.scala:644)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.sql.DataFrameReader.parquet(DataFrameReader.scala:643)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at com.js.utils.ReadUtils.readParquetFromS3(ReadUtils.java:24)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at com.js.newjs.preprocessors.DeletedProfileContactsEligibleForRetPreProcessor.incrementalDataProcess(DeletedProfileContactsEligibleForRetPreProcessor.java:91)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at com.js.JsETLSparkTransformationsApplication.process(JsETLSparkTransformationsApplication.java:91)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at com.js.JsETLSparkTransformationsApplication.main(JsETLSparkTransformationsApplication.java:32)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at java.lang.reflect.Method.invoke(Method.java:498)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.deploy.JavaMainApplication.start(SparkApplication.scala:52)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.deploy.SparkSubmit.org$apache$spark$deploy$SparkSubmit$$runMain(SparkSubmit.scala:845)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.deploy.SparkSubmit.doRunMain$1(SparkSubmit.scala:161)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.deploy.SparkSubmit.submit(SparkSubmit.scala:184)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.deploy.SparkSubmit.doSubmit(SparkSubmit.scala:86)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.deploy.SparkSubmit$$anon$2.doSubmit(SparkSubmit.scala:920)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.deploy.SparkSubmit$.main(SparkSubmit.scala:929)
[2022-11-17, 08:46:41 IST] {subprocess.py:92} INFO -    at org.apache.spark.deploy.SparkSubmit.main(SparkSubmit.scala)

解决方案

1. 让Spark任务返回非零退出码

Airflow的BashOperator通过命令的退出码判断任务状态,只要Spark任务在失败时返回非零值,Airflow就会标记任务失败。如果是Java编写的Spark应用,在捕获到AnalysisException时调用System.exit(1):

try {
    // 数据读取与处理逻辑
} catch (AnalysisException e) {
    log.error("数据路径不存在或分析错误", e);
    System.exit(1); // 关键:返回非零退出码
}

2. 在Bash命令中添加错误校验

如果无法修改Spark代码,可在执行spark-submit的命令中增加错误检查逻辑:

spark-submit --class com.js.JsETLSparkTransformationsApplication your-application.jar > spark_task.log 2>&1
# 检查命令退出码
if [ $? -ne 0 ]; then
    exit 1
fi
# 或检查日志中是否存在ERROR关键字(需确保该ERROR是任务失败的标志)
if grep -q "ERROR" spark_task.log; then
    exit 1
fi

3. 使用Airflow原生

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 23:55:19