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
相关产品推荐
相关产品推荐

