Airflow容器中Spark-Submit DAG无法定位GCS Connector文件求助
问题分析与解决方案
核心问题
你遇到的FileNotFoundException和ClassNotFoundException本质是Airflow与Spark容器的文件系统隔离导致的:
- 直接在Spark容器内提交任务时,JAR和凭证文件的路径是有效的,但通过Airflow的
SparkSubmitOperator提交时,任务的配置路径未对应到Spark容器内的实际位置; - 代码中硬编码的
setMaster('local')会覆盖Airflow Spark连接配置的集群Master,可能导致任务执行上下文路径混乱。
解决方案
1. 移除Spark代码中的硬编码配置,由Airflow统一管理
修改你的Spark任务代码,删除setMaster('local')和spark.jars配置(这些配置将由Airflow的SparkSubmitOperator传递):
from pyspark.sql import SparkSession from pyspark.conf import SparkConf from pyspark.context import SparkContext if __name__ == '__main__': # 仅保留应用名称和GCS认证基础配置 conf = SparkConf() \ .setAppName('test') \ .set("spark.hadoop.google.cloud.auth.service.account.enable", "true") sc = SparkContext(conf=conf.set("spark.files.overwrite", "true")) hadoop_conf = sc._jsc.hadoopConfiguration() hadoop_conf.set("fs.AbstractFileSystem.gs.impl", "com.google.cloud.hadoop.fs.gcs.GoogleHadoopFS") hadoop_conf.set("fs.gs.impl", "com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem") hadoop_conf.set("fs.gs.auth.service.account.enable", "true") spark = SparkSession.builder \ .config(conf=sc.getConf()) \ .getOrCreate() path = "gs://tfl-cycling/pq/" df_test = spark.read.option("recursiveFileLookup", "true").parquet(path) df_test.printSchema()
2. 配置SparkSubmitOperator传递JAR与凭证路径
在Airflow DAG中,通过SparkSubmitOperator明确指定Spark容器内的JAR路径、凭证文件路径,或使用Maven坐标自动拉取依赖:
Etl = SparkSubmitOperator( application = "/opt/airflow/dags/example.py", conn_id= 'spark_default', task_id='spark_submit_task_load', dag=spark_dag, # 方式1:指定Spark容器内JAR的绝对路径(与你直接提交时的有效路径一致) jars="/spark-lib/gcs-connector-hadoop3-2.2.5.jar", # 方式2:如果Spark集群能联网,用Maven坐标自动拉取(无需本地JAR) # packages="com.google.cloud.bigdataoss:gcs-connector:hadoop3-2.2.5", # 传递GCS凭证文件路径(Spark容器内的绝对路径) conf={ "spark.hadoop.google.cloud.auth.service.account.json.keyfile": "/path/to/your/credentials/in/spark/container/google_credentials.json" }, # 如果凭证文件在Airflow容器内,可通过files参数分发到Spark节点 # files="/opt/airflow/.google/credentials/google_credentials.json" )
3. 备选方案:将JAR放入Spark全局classpath
如果不想每次提交任务都指定JAR,可将gcs-connector-hadoop3-2.2.5.jar复制到Spark容器的$SPARK_HOME/jars目录下,Spark会自动加载该JAR,无需在任务中配置spark.jars。
关键注意事项
- 确保Airflow的Spark连接
conn_id指向的Spark Master节点,能访问到你指定的JAR和凭证文件; - 若使用
files参数分发凭证文件,Spark任务中需用SparkFiles.get("google_credentials.json")获取文件路径,而非硬编码相对路径。
内容的提问来源于stack exchange,提问作者Dawkins
相关产品推荐
相关产品推荐

