Scala编写的Spark Streaming任务在Airflow中运行失败求助
问题描述
我平时使用PySpark,现在需要处理一个Scala编写的Spark Streaming任务。直接在EMR上执行spark-submit可以正常运行,但通过Airflow执行时出现如下错误,不知道从何开始调试,恳请提供解决思路。
错误堆栈信息
org.apache.spark.SparkException: 在awaitResult中抛出异常: at org.apache.spark.util.ThreadUtils$.awaitResult(ThreadUtils.scala:226) at org.apache.spark.deploy.yarn.ApplicationMaster.runDriver(ApplicationMaster.scala:468) at org.apache.spark.deploy.yarn.ApplicationMaster.org$apache$spark$deploy$yarn$ApplicationMaster$$runImpl(ApplicationMaster.scala:305) at org.apache.spark.deploy.yarn.ApplicationMaster$$anonfun$run$1.apply$mcV$sp(ApplicationMaster.scala:245) at org.apache.spark.deploy.yarn.ApplicationMaster$$anonfun$run$1.apply(ApplicationMaster.scala:245) at org.apache.spark.deploy.yarn.ApplicationMaster$$anonfun$run$1.apply(ApplicationMaster.scala:245) at org.apache.spark.deploy.yarn.ApplicationMaster$$anon$3.run(ApplicationMaster.scala:779) at java.security.AccessController.doPrivileged(Native Method) at javax.security.auth.Subject.doAs(Subject.java:422) at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1844) at org.apache.spark.deploy.yarn.ApplicationMaster.doAsUser(ApplicationMaster.scala:778) at org.apache.spark.deploy.yarn.ApplicationMaster.run(ApplicationMaster.scala:244) at org.apache.spark.deploy.yarn.ApplicationMaster$.main(ApplicationMaster.scala:803) at org.apache.spark.deploy.yarn.ApplicationMaster.main(ApplicationMaster.scala) Caused by: com.typesafe.config.ConfigException$IO: available_application.properties -Dlog4j.configuration=log4j-yarn.properties: java.io.FileNotFoundException: available_application.properties -Dlog4j.configuration=log4j-yarn.properties (没有该文件或目录) at com.typesafe.config.impl.Parseable.parseValue(Parseable.java:183) at com.typesafe.config.impl.Parseable.parseValue(Parseable.java:170) at com.typesafe.config.impl.Parseable.parse(Parseable.java:227) at com.typesafe.config.ConfigFactory.parseFile(ConfigFactory.java:595) at com.typesafe.config.ConfigFactory.loadDefaultConfig(ConfigFactory.java:244) at com.typesafe.config.ConfigFactory.access$000(ConfigFactory.java:38) at com.typesafe.config.ConfigFactory$1.call(ConfigFactory.java:378) at com.typesafe.config.ConfigFactory$1.call(ConfigFactory.java:375) at com.typesafe.config.impl.ConfigImpl$LoaderCache.getOrElseUpdate(ConfigImpl.java:58) at com.typesafe.config.impl.ConfigImpl.computeCachedConfig(ConfigImpl.java:86) at com.typesafe.config.ConfigFactory.load(ConfigFactory.java:375) at com.typesafe.config.ConfigFactory.load(ConfigFactory.java:299) at com.typesafe.config.ConfigFactory.load(ConfigFactory.java:288) at com.nike.tdp.AvailabilityKafkaEvents$.main(AvailabilityKafkaEvents.scala:101) at com.nike.tdp.AvailabilityKafkaEvents.main(AvailabilityKafkaEvents.scala) at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.lang.reflect.Method.invoke(Method.java:498) at org.apache.spark.deploy.yarn.ApplicationMaster$$anon$2.run(ApplicationMaster.scala:684) Caused by: java.io.FileNotFoundException: available_application.properties -Dlog4j.configuration=log4j-yarn.properties (没有该文件或目录) at java.io.FileInputStream.open0(Native Method) at java.io.FileInputStream.open(FileInputStream.java:195) at java.io.FileInputStream.<init>(FileInputStream.java:138) at com.typesafe.config.impl.Parseable$ParseableFile.reader(Parseable.java:512) at com.typesafe.config.impl.Parseable.rawParseValue(Parseable.java:193) at com.typesafe.config.impl.Parseable.parseValue(Parseable.java:176) ... 19 more 22/10/26 19:14:02 INFO ShutdownHookManager: 调用关闭钩子
调试解决思路
- 定位核心问题:从错误堆栈可以看到,程序把
available_application.properties -Dlog4j.configuration=log4j-yarn.properties当作一个完整的文件名去读取,导致找不到文件。这说明参数传递时出现了格式错误,把JVM参数和配置文件路径混在了一起。 - 对比Spark-submit命令格式:把Airflow中执行的
spark-submit命令和EMR本地直接运行的命令做对比,检查-Dlog4j.configuration=log4j-yarn.properties这个JVM参数的位置是否正确。正确的做法是将这类参数放在--driver-java-options(给Driver)或--executor-java-options(给Executor)后面,不能和配置文件路径放在同一参数项中。 - 确认配置文件路径有效性:Airflow执行任务的工作目录和EMR本地shell的工作目录可能不同,确保
available_application.properties使用绝对路径,或者通过spark-submit的--files参数将该文件分发到集群所有节点的工作目录中。 - 检查Airflow任务环境变量:Airflow运行任务的环境变量(如
SPARK_HOME、HADOOP_CONF_DIR)可能和EMR本地不一致,这些变量会影响Spark的参数解析逻辑,需要确认环境变量是否正确设置。 - 查看Airflow完整执行日志:在Airflow的任务日志中找到完整的
spark-submit执行命令,确认参数是否被正确传递,有没有被Airflow的模板渲染、字符转义等操作篡改。
内容的提问来源于stack exchange,提问作者ndev
相关产品推荐
相关产品推荐

