Airflow通过XCOM实现Python Operator向Bash Operator传值
Airflow XCOM跨任务多值传递方案
核心问题答复
- 单个Bash命令完全支持拉取多个XCOM变量,
bash_command字段支持完整Jinja2模板语法,可在同一命令内多次调用xcom_pull拉取不同task、不同key的值,无数量限制。 - 原有代码运行异常的原因有三点:
- 拉取XCOM时未指定push时设置的
count、recon独立key,直接拉取整个task的XCOM会返回包含所有推送值的字典,无法直接作为脚本入参被识别 - 代码从网页复制时带了HTML转义字符
&,实际shell命令需要的是&&,且单引号包裹的命令字符串随意折行会触发Python语法错误 - 拉取到的值没有按
app.py要求的count_check、recon_check入参格式传递,直接传入字典对象会导致脚本报错
- 拉取XCOM时未指定push时设置的
正确实现代码
from airflow import DAG from airflow.operators.python import PythonOperator from airflow.operators.bash import BashOperator from airflow.utils.dates import datetime args= {'owner': 'airflow', 'start_date': datetime(2022, 1, 1) } def take_args(ti): count_check = 'Y' recon_check = 'Y' ti.xcom_push(key='count', value=count_check) ti.xcom_push(key='recon', value=recon_check) with DAG(dag_id='data-validation-dag-python', default_args=args, schedule_interval='@daily', max_active_runs=1, catchup=False) as dag: task_1 = PythonOperator( task_id='Storing_Args', python_callable=take_args ) task_2 = BashOperator( task_id='task_validation_checks', # 分别拉取两个key的XCOM值作为入参,加引号避免特殊字符导致命令解析异常 bash_command='cd /root/airflow/dags && python3 app.py "{{ ti.xcom_pull(task_ids="Storing_Args", key="count") }}" "{{ ti.xcom_pull(task_ids="Storing_Args", key="recon") }}"', do_xcom_push=False ) task_1 >> task_2
补充说明
- 如果你的
app.py用argparse等方式解析关键字参数,推荐把参数部分写成显式关键字传参的格式,不需要严格依赖参数顺序,容错性更高:cd /root/airflow/dags && python3 app.py \ --count_check "{{ ti.xcom_pull(task_ids="Storing_Args", key="count") }}" \ --recon_check "{{ ti.xcom_pull(task_ids="Storing_Args", key="recon") }}" - XCOM传递的值默认会做模板转义,如果传递的参数包含空格、特殊字符,一定要给Jinja模板片段外层加双引号,避免shell把参数拆分成多个部分。
内容的提问来源于stack exchange,提问作者ojjasvi nirmal
相关产品推荐
相关产品推荐

