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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 07:44:52