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

如何通过Apache Airflow远程运行远端Spark集群的PySpark任务

问题根因

你当前直接用PythonOperator在Airflow所在的abc服务器本地运行PySpark代码,报错是必然的,核心问题有三个:

  • Jar包/路径找不到:你代码里配置的spark.driver.extraClassPath、spark.executor.extraClassPath、spark.local.dir全是xyz服务器上的本地路径,默认情况下Spark任务的Driver进程会运行在提交任务的客户端节点也就是abc上,abc上根本没有这些路径和文件,自然报错。
  • 远程连接失败:要么是abc到xyz的7077端口(Spark Master RPC端口)、Spark Driver/Executor之间的随机通信端口没开防火墙,要么是abc上根本没装和xyz集群版本完全匹配的Spark客户端,硬写master地址初始化SparkSession根本完不成集群握手。
  • 你代码里写的spark.driver.supervise配置,只有cluster提交模式下在集群节点启动Driver时才生效,本地客户端提交时开这个参数也会抛错。
正确实现方案

生产环境推荐两种稳定落地的方案,优先选第一种:

方案1:使用SparkSubmitOperator(官方推荐,最规范)

不要用PythonOperator直接跑PySpark逻辑,Airflow官方自带SparkSubmitOperator,专门用来提交任务到远程Spark集群,不需要在业务代码里硬写master、Jar路径这类配置。

前置准备

  • 在abc服务器上安装和xyz集群版本完全一致的Spark二进制包,配置好SPARK_HOME环境变量,保证abc上执行spark-submit --master spark://xyz:7077 --version能正常返回集群信息,无连接报错。
  • 把你的PySpark业务代码(包括session_open方法)单独存成.py文件,放到abc服务器Airflow能访问的路径下,比如/opt/airflow/dags/spark_jobs/psv_etl.py。
  • 把mysql connector这类依赖Jar包,要么放到abc服务器Spark安装目录的jars/文件夹下,要么在提交的时候通过--jars参数分发到集群,不要在代码里硬编码xyz的本地路径。
  • 确认abc服务器和xyz集群之间的防火墙开放:7077(Spark Master)、8080(Spark UI)、Spark随机通信端口段(默认40000-60000,也可以在Spark配置里固定端口段)。

代码修改示例

首先修改PySpark业务代码的SparkSession初始化逻辑,去掉硬编码的master、本地路径、extraClassPath配置,这些参数统一交给spark-submit传入:

# psv_etl.py 存放在abc服务器的Airflow可访问路径
from pyspark.sql import SparkSession

def session_open():
    spark = SparkSession \
           .builder \
           .appName("Program_Single_View" ) \
           .config("spark.executor.memory", "14g") \
           .config("spark.cores.max", "6")  \
           .config("spark.driver.memory", "4g")  \
           .config("spark.executor.memoryOverhead", "384")  \
           .config("spark.default.parallelism","60") \
           .config("spark.sql.shuffle.partitions","36") \
           .config("spark.memory.offHeap.enabled","true") \
           .config("spark.memory.offHeap.size","1g") \
           .config("spark.network.timeout", "100000000")\
           .config("spark.executor.heartbeatInterval","10000000")\
           .config("spark.driver.maxResultSize","4g") \
           .config('spark.sql.session.timeZone', 'UTC') \
           .getOrCreate()
    return spark

# 下方写具体ETL业务逻辑
if __name__ == "__main__":
    spark = session_open()
    # 业务代码
    spark.stop()

然后修改DAG文件,替换原来的PythonOperator为SparkSubmitOperator:

from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator
import pendulum
from airflow.operators.dummy import DummyOperator
from datetime import datetime

dt = datetime
default_args = {
    'owner': 'airflow'
}

localtz = pendulum.timezone("Asia/Kolkata")
with DAG('PSV_all_server_ETL', default_args=default_args, description='etl', schedule_interval=None,
         start_date=dt(2022, 6, 21, tzinfo=localtz),
         catchup=False, tags=['PSV']) as dag:

    program_single_view_execution = SparkSubmitOperator(
        task_id="program_single_view_execution",
        application="/opt/airflow/dags/spark_jobs/psv_etl.py",
        master="spark://xyz:7077",
        deploy_mode="cluster", # cluster模式下Driver运行在xyz集群节点上,不依赖abc本地路径
        spark_home="/path/to/spark/on/abc", # 替换为abc服务器上的SPARK_HOME实际路径
        jars="/path/to/mysql-connector-java-8.0.17.jar", # 替换为abc上存放的Jar包实际路径,提交时会自动分发到集群
        conf={
            "spark.local.dir": "/datadrive",
            "spark.driver.supervise": "true"
        },
        driver_memory="4g",
        executor_memory="14g",
        cores_max=6
    )

    start = DummyOperator(task_id='Start')
    end = DummyOperator(task_id='End')

    start >> program_single_view_execution >> end

注意:如果用deploy_mode="cluster",Driver会运行在xyz集群的Worker节点上,所有配置的本地路径、Jar包只要在xyz集群节点上存在即可,不需要在abc上存业务依赖文件;如果用client模式,Driver跑在abc上,所有本地路径、Jar都要在abc上存对应文件。

方案2:用SSH远程提交(适合不想在abc装Spark客户端的场景)

如果不想在abc服务器部署Spark客户端,可以用SSHOperator直接ssh登录到xyz服务器,在xyz本地执行spark-submit命令提交任务,所有路径、Jar包都复用xyz本地的,原有PySpark代码基本不用改。

前置准备

  • 配置abc服务器上的Airflow运行用户到xyz服务器的SSH免密登录。
  • 把原有PySpark代码全部放到xyz服务器上的固定路径,比如/datadrive/PROD/psv_etl.py。

DAG代码修改示例

from airflow.providers.ssh.operators.ssh import SSHOperator
from airflow.operators.dummy import DummyOperator
import pendulum
from datetime import datetime

dt = datetime
default_args = {
    'owner': 'airflow'
}

localtz = pendulum.timezone("Asia/Kolkata")
with DAG('PSV_all_server_ETL', default_args=default_args, description='etl', schedule_interval=None,
         start_date=dt(2022, 6, 21, tzinfo=localtz),
         catchup=False, tags=['PSV']) as dag:

    program_single_view_execution = SSHOperator(
        task_id="program_single_view_execution",
        ssh_conn_id="ssh_xyz", # 提前在Airflow连接里配置xyz的SSH地址、用户名、认证密钥
        command="""
        source /etc/profile;
        cd /datadrive/PROD/;
        spark-submit --master spark://xyz:7077 \
        --executor-memory 14g \
        --conf spark.cores.max=6 \
        --driver-memory 4g \
        --conf spark.executor.memoryOverhead=384 \
        --conf spark.local.dir=/datadrive \
        --conf spark.default.parallelism=60 \
        --conf spark.sql.shuffle.partitions=36 \
        --conf spark.memory.offHeap.enabled=true \
        --conf spark.memory.offHeap.size=1g \
        --conf spark.network.timeout=100000000 \
        --conf spark.executor.heartbeatInterval=10000000 \
        --conf spark.driver.maxResultSize=4g \
        --conf spark.sql.session.timeZone=UTC \
        --conf spark.driver.supervise=true \
        --driver-class-path /datadrive/PROD/jars/mysql-connector-java-8.0.17.jar \
        --jars /datadrive/PROD/jars/mysql-connector-java-8.0.17.jar \
        psv_etl.py
        """
    )

    start = DummyOperator(task_id='Start')
    end = DummyOperator(task_id='End')

    start >> program_single_view_execution >> end
常见坑点提醒
  • 不要用PythonOperator直接在Airflow Worker进程里初始化SparkSession连远程集群,这种模式下Spark上下文跑在Airflow Worker进程里,稳定性极差,端口、路径、依赖问题极多,生产环境完全不推荐。
  • Spark集群所有节点和abc服务器的主机名要能互相解析,要么配内部DNS要么统一加hosts映射,不然会出现Worker连不上Driver的报错。
  • 如果用cluster模式提交任务,记得确认业务代码里所有写死的本地路径都是集群节点上真实存在的路径,不要引用abc上的本地路径。

内容的提问来源于stack exchange,提问作者subh

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 13:54:27