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

Airflow中如何向SparkSubmitOperator触发的PySpark任务传递参数

问题根源

你在DAG中定义的spark是运行在Airflow进程内的对象,SparkSubmitOperator触发的是独立的Spark应用进程,二者内存空间完全隔离,不可能直接传递内存实例。application_args仅支持传递字符串类型的可序列化参数,你传入的['spark']只是普通字符串,不是你定义的SparkSession实例,自然会被识别为无效值。

SparkSession必须在你提交的PySpark脚本内部自行初始化,不需要从外部传入。如果需要传递业务参数(比如运行日期、表名、配置项等),才需要通过application_args传递字符串格式的参数值。

正确实现方案

1. 修正DAG代码

删除DAG中不必要的SparkSession初始化,application_args传入你实际需要的业务参数即可:

from airflow import DAG
from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator
from datetime import datetime

default_args = {
    'owner': 'airflow',
    'start_date': datetime(2024, 1, 1)
}

spark_config = {
    'conn_id': 'spark_default',
    'driver_memory': '1g',
    'executor_cores': 1,
    'num_executors': 1,
    'executor_memory': '1g'
}

dag = DAG(
    dag_id="spark_session_prgm",
    default_args=default_args,
    schedule_interval='@daily',
    catchup=False)

spark_submit_task1 = SparkSubmitOperator(
    task_id='spark_submit_task1',
    application='/home/airflow_home/dags/tmp_spark_1.py',
    # 此处传入你需要的业务参数,全部为字符串格式,比如执行日期、表名
    application_args=['{{ ds }}', 'dwd.user_info'],
    **spark_config, 
    dag=dag
)

2. 修正PySpark业务脚本(tmp_spark_1.py)

在脚本内部初始化SparkSession,通过sys.argv接收传入的参数:

import sys
from pyspark.sql import SparkSession

# 脚本内部自行初始化SparkSession
spark = SparkSession.builder.appName("tmp_spark_1").enableHiveSupport().getOrCreate()

# 按顺序接收传入的参数,sys.argv[0]为脚本本身路径,自定义参数从下标1开始
dt = sys.argv[1]
table_name = sys.argv[2]

# 后续业务逻辑
df = spark.table(table_name).filter(f"dt = '{dt}'")
df.show()
注意事项
  • application_args支持Airflow Jinja模板语法,可以传入动态参数,比如示例中的{{ ds }}会自动替换为DAG的执行日期
  • 传入参数的顺序要和脚本中sys.argv的取值顺序一一对应
  • 所有通过application_args传递的参数都必须是字符串格式,不可传递非序列化的内存对象

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 07:36:02