Airflow中XCOM值拉取返回None问题排查及解决
Airflow XCOM拉取返回None问题解决
我有一个简单的Airflow DAG,包含一个PythonOperator任务,用于从SWAPI API获取简单JSON数据,返回的身高整数为202。
我已确认该值已被正确推送为XCOM值:运行DAG后查看该任务日志和UI的XCOM面板,都能看到对应记录,截图如下:
我在调用API的Python函数中添加ti.xcom_push(key = 'height', value = height)代码后,也能在该任务的XCOM视图中看到height对应的值为202。
但我始终无法在其他任务中拉取到该值,比如我使用如下PythonOperator任务拉取:
def check_height(ti): height = ti.xcom_pull(key = 'height', task_ids=['get_data_darth_vader']) print(f"Height is: {height}")
我也试过不带key、指定key为'return_value'的拉取方式,均返回None,运行日志如下:
[2021-09-30 21:00:35,044] {logging_mixin.py:109} INFO - Height is: [None] [2021-09-30 21:00:35,047] {python.py:151} INFO - Done. Returned value was: None
初始DAG完整代码
from airflow import DAG from airflow.operators.python import PythonOperator, BranchPythonOperator from airflow.operators.bash import BashOperator from datetime import datetime import json import requests def get_darth_vader_height(ti): """ Get Darth Vader info from SWAPI """ response=requests.get('https://swapi.dev/api/people/4') data=json.loads(response.text) height=data['height'] print(f"DEBUG: {height}") ti.xcom_push(key="height", value=height) return height def check_height(ti): height=ti.xcom_pull(task_ids='task_one', key="height") print(f"Height is: {height}") print(str(height)) with DAG( 'my_dag', start_date = datetime(2021,1,1), schedule_interval="@daily", catchup=False, ) as dag: get_darth_vader_height = PythonOperator( task_id='task_one', python_callable=get_darth_vader_height ) check_darth_vader_height = PythonOperator( task_id='task_two', python_callable=check_height ) is_tall = BashOperator( task_id='task_three', bash_command="echo 'is tall!'" ) is_short = BashOperator( task_id='task_four', bash_command="echo 'is short!'" )
问题原因与修复方案
问题核心是没有定义任务之间的依赖关系,Airflow默认所有任务独立并行执行,拉取XCOM的task_two会在推送XCOM的task_one执行完成前就启动,此时XCOM还未生成,自然拉取到None。
只需添加任务依赖链,确保task_one执行完成后再运行task_two即可解决问题,修复后的完整可运行DAG代码如下:
from airflow import DAG from airflow.operators.python import PythonOperator, BranchPythonOperator from airflow.operators.bash import BashOperator from datetime import datetime import json import requests def get_darth_vader_height(ti): """ Get Darth Vader info from SWAPI """ response=requests.get('https://swapi.dev/api/people/4') data=json.loads(response.text) height=data['height'] print(f"DEBUG: {height}") ti.xcom_push(key="height", value=height) return height def check_height(ti): height=ti.xcom_pull(task_ids='task_one', key="height") print(f"Height is: {height}") if int(height) > 200: print('height is greater than 200') return 'task_three' print('height is less than 200') return 'task_four' with DAG( 'my_dag', start_date = datetime(2021,1,1), schedule_interval="@daily", catchup=False, ) as dag: get_darth_vader_height = PythonOperator( task_id='task_one', python_callable=get_darth_vader_height ) check_darth_vader_height = BranchPythonOperator( task_id='task_two', python_callable=check_height ) is_tall = BashOperator( task_id='task_three', bash_command="echo 'is tall!'" ) is_short = BashOperator( task_id='task_four', bash_command="echo 'is short!'" ) get_darth_vader_height >> check_darth_vader_height >> [is_tall, is_short]
内容的提问来源于stack exchange,提问作者Indrid
相关产品推荐
相关产品推荐

