Airflow中读取dag_run.conf参数及使用xCom结果的问题
关于Airflow中根据触发参数选择Spark配置文件的解决方案
不能直接在非Operator的普通Python逻辑中使用xCom值——Airflow的DAG解析阶段(普通Python代码执行时)早于任务执行阶段,xCom是任务运行时才生成的,解析阶段根本获取不到这些运行时数据。
以下是两种可行的解决方案:
方案1:用BranchPythonOperator分支任务
通过BranchPythonOperator直接从触发上下文获取run_mode参数,分支到对应配置文件的Spark任务:
from airflow import DAG from airflow.operators.python import BranchPythonOperator from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator from datetime import datetime def pick_run_mode(**context): # 从触发的dag_run配置中提取run_mode,默认值设为'incremental' run_mode = context["dag_run"].conf.get("run_mode", "incremental") return f"spark_job_{run_mode}" with DAG( dag_id="dynamic_spark_config", start_date=datetime(2024, 1, 1), catchup=False, ) as dag: branch_task = BranchPythonOperator( task_id="decide_run_mode", python_callable=pick_run_mode, provide_context=True, ) # 全量运行任务,使用full配置 spark_full = SparkSubmitOperator( task_id="spark_job_full", application="/opt/spark/app/main.py", application_args=["--config", "/opt/configs/full.application.conf"], conf={"spark.executor.memory": "4g"}, ) # 增量运行任务,使用incremental配置 spark_incremental = SparkSubmitOperator( task_id="spark_job_incremental", application="/opt/spark/app/main.py", application_args=["--config", "/opt/configs/incremental.application.conf"], conf={"spark.executor.memory": "2g"}, ) branch_task >> [spark_full, spark_incremental]
方案2:用模板变量动态指定配置文件
如果不需要分支任务,可直接在SparkSubmitOperator的参数中使用Airflow模板变量,动态拼接配置文件路径:
from airflow import DAG from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator from datetime import datetime with DAG( dag_id="dynamic_spark_config_template", start_date=datetime(2024, 1, 1), catchup=False, ) as dag: spark_task = SparkSubmitOperator( task_id="run_spark_job", application="/opt/spark/app/main.py", # 用模板变量从dag_run.conf中取run_mode,默认用incremental配置 application_args=[ "--config", "/opt/configs/{{ dag_run.conf.get('run_mode', 'incremental') }}.application.conf" ], conf={"spark.executor.memory": "4g"}, )
补充说明
- 触发DAG时传入的
{"run_mode":"full"}参数,会被Airflow存在dag_run.conf中,可直接在Operator的上下文或模板中获取,无需通过xCom中转。 - 非Operator的普通Python代码属于DAG解析阶段逻辑,此时DAG还未触发运行,xCom、dag_run等运行时数据都不存在,因此无法直接使用。
内容的提问来源于stack exchange,提问作者Violeta Andreea Tudor
相关产品推荐
相关产品推荐

