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原生
相关产品推荐
相关产品推荐

