Airflow Taskflow API无法获取XCOM返回值,请求排查错误
Airflow Taskflow API传递字典XCom失败问题排查
你尝试使用Airflow Taskflow API返回字典并传递给下游任务,但greet_multiple任务报错无法找到get_fullname返回的firstname对应的XCom结果。以下是你的代码与任务日志:
原代码
from datetime import datetime, timedelta from airflow.decorators import dag, task default_args = { 'owner': 'ashish', 'retries': '2', 'retry_delay': timedelta(minutes=2) } @dag(dag_id='DAG_with_TaskFlow_API', default_args=default_args, description='Example for Task Flow API', start_date=datetime(2023, 8, 30, 2), schedule_interval='@daily') def hello(): # Setting the value normally @task() def get_name(): return 'Ashish' # Passing multiple values @task() def get_fullname(multiple_outputs=True): return {'firstname': 'Ashish', 'lastname': 'Inamdar'} @task() def greet(name): print(f"Hello {name}") @task() def greet_multiple(firstname, lastname): print(f"Hi {firstname} {lastname}") name = get_name() greet(name=name) fullname_dict = get_fullname() greet_multiple(fullname_dict['firstname'], fullname_dict['lastname']) taskflow_api_dag = hello()
任务日志
3a4feea079c *** Found local files: *** * /opt/airflow/logs/dag_id=DAG_with_TaskFlow_API/run_id=manual__2023-09-10T05:23:33.215377+00:00/task_id=greet_multiple/attempt=3.log [2023-09-10, 05:27:46 UTC] {taskinstance.py:1159} INFO - Dependencies all met for dep_context=non-requeueable deps ti=<TaskInstance: DAG_with_TaskFlow_API.greet_multiple manual__2023-09-10T05:23:33.215377+00:00 [queued]> [2023-09-10, 05:27:46 UTC] {taskinstance.py:1159} INFO - Dependencies all met for dep_context=requeueable deps ti=<TaskInstance: DAG_with_TaskFlow_API.greet_multiple manual__2023-09-10T05:23:33.215377+00:00 [queued]> [2023-09-10, 05:27:46 UTC] {taskinstance.py:1361} INFO - Starting attempt 3 of 3 [2023-09-10, 05:27:46 UTC] {taskinstance.py:1382} INFO - Executing <Task(_PythonDecoratedOperator): greet_multiple> on 2023-09-10 05:23:33.215377+00:00 [2023-09-10, 05:27:46 UTC] {standard_task_runner.py:57} INFO - Started process 16082 to run task [2023-09-10, 05:27:46 UTC] {standard_task_runner.py:84} INFO - Running: ['***', 'tasks', 'run', 'DAG_with_TaskFlow_API', 'greet_multiple', 'manual__2023-09-10T05:23:33.215377+00:00', '--job-id', '162', '--raw', '--subdir', 'DAGS_FOLDER/task_with_taskflow_api.py', '--cfg-path', '/tmp/tmptyfhicx8'] [2023-09-10, 05:27:46 UTC] {standard_task_runner.py:85} INFO - Job 162: Subtask greet_multiple [2023-09-10, 05:27:46 UTC] {task_command.py:415} INFO - Running <TaskInstance: DAG_with_TaskFlow_API.greet_multiple manual__2023-09-10T05:23:33.215377+00:00 [running]> on host f3a4feea079c [2023-09-10, 05:27:46 UTC] {abstractoperator.py:696} ERROR - Exception rendering Jinja template for task 'greet_multiple', field 'op_args'. Template: (XComArg(<Task(_PythonDecoratedOperator): get_fullname>, 'firstname'), XComArg(<Task(_PythonDecoratedOperator): get_fullname>, 'lastname')) Traceback (most recent call last): File "/home/airflow/.local/lib/python3.8/site-packages/airflow/models/abstractoperator.py", line 688, in _do_render_template_fields rendered_content = self.render_template( File "/home/airflow/.local/lib/python3.8/site-packages/airflow/template/templater.py", line 162, in render_template return tuple(self.render_template(element, context, jinja_env, oids) for element in value) File "/home/airflow/.local/lib/python3.8/site-packages/airflow/template/templater.py", line 162, in <genexpr> return tuple(self.render_template(element, context, jinja_env, oids) for element in value) File "/home/airflow/.local/lib/python3.8/site-packages/airflow/template/templater.py", line 158, in render_template return value.resolve(context) File "/home/airflow/.local/lib/python3.8/site-packages/airflow/utils/session.py", line 77, in wrapper return func(*args, session=session, **kwargs) File "/home/airflow/.local/lib/python3.8/site-packages/airflow/models/xcom_arg.py", line 431, in resolve raise XComNotFound(ti.dag_id, task_id, self.key) airflow.exceptions.XComNotFound: XComArg result from get_fullname at DAG_with_TaskFlow_API with key="firstname" is not found! [2023-09-10, 05:27:46 UTC] {taskinstance.py:1943} ERROR - Task failed with exception Traceback (most recent call last): File "/home/airflow/.local/lib/python3.8/site-packages/airflow/models/taskinstance.py", line 1518, in _run_raw_task self._execute_task_with_callbacks(context, test_mode, session=session) File "/home/airflow/.local/lib/python3.8/site-packages/airflow/models/taskinstance.py", line 1646, in _execute_task_with_callbacks task_orig = self.render_templates(context=context) File "/home/airflow/.local/lib/python3.8/site-packages/airflow/models/taskinstance.py", line 2291, in render_templates original_task.render_template_fields(context) File "/home/airflow/.local/lib/python3.8/site-packages/airflow/models/baseoperator.py", line 1244, in render_template_fields self._do_render_template_fields(self, self.template_fields, context, jinja_env, set()) File "/home/airflow/.local/lib/python3.8/site-packages/airflow/utils/session.py", line 77, in wrapper return func(*args, session=session, **kwargs) File "/home/airflow/.local/lib/python3.8/site-packages/airflow/models/abstractoperator.py", line 688, in _do_render_template_fields rendered_content = self.render_template( File "/home/airflow/.local/lib/python3.8/site-packages/airflow/template/templater.py", line 162, in render_template return tuple(self.render_template(element, context, jinja_env, oids) for element in value) File "/home/airflow/.local/lib/python3.8/site-packages/airflow/template/templater.py", line 162, in <genexpr> return tuple(self.render_template(element, context, jinja_env, oids) for element in value) File "/home/airflow/.local/lib/python3.8/site-packages/airflow/template/templater.py", line 158, in render_template return value.resolve(context) File "/home/airflow/.local/lib/python3.8/site-packages/airflow/utils/session.py", line 77, in wrapper return func(*args, session=session, **kwargs) File "/home/airflow/.local/lib/python3.8/site-packages/airflow/models/xcom_arg.py", line 431, in resolve raise XComNotFound(ti.dag_id, task_id, self.key) airflow.exceptions.XComNotFound: XComArg result from get_fullname at DAG_with_TaskFlow_API with key="firstname" is not found! [2023-09-10, 05:27:46 UTC] {taskinstance.py:1400} INFO - Marking task as FAILED. dag_id=DAG_with_TaskFlow_API, task_id=greet_multiple, execution_date=20230910T052333, start_date=20230910T052746, end_date=20230910T052746 [2023-09-10, 05:27:46 UTC] {standard_task_runner.py:104} ERROR - Failed to execute job 162 for task greet_multiple (XComArg result from get_fullname at DAG_with_TaskFlow_API with key="firstname" is not found!; 16082) [2023-09-10, 05:27:46 UTC] {local_task_job_runner.py:228} INFO - Task exited with return code 1 [2023-09-10, 05:27:46 UTC] {taskinstance.py:2784} INFO - 0 downstream tasks scheduled from follow-on schedule check
错误原因
你把multiple_outputs=True放在了get_fullname函数的参数列表里,而它应该是@task装饰器的参数。这个参数的作用是告诉Airflow:当任务返回字典时,要把字典的每个键值对作为单独的XCom条目存储,这样下游才能通过键名访问对应的值。
因为参数位置错误,Airflow没有启用多输出模式,get_fullname返回的字典被当作一个整体存进了XCom,下游尝试访问fullname_dict['firstname']时,Airflow会去寻找键为firstname的XCom条目,自然找不到,导致报错。
修正后的代码
from datetime import datetime, timedelta from airflow.decorators import dag, task default_args = { 'owner': 'ashish', 'retries': '2', 'retry_delay': timedelta(minutes=2) } @dag(dag_id='DAG_with_TaskFlow_API', default_args=default_args, description='Example for Task Flow API', start_date=datetime(2023, 8, 30, 2), schedule_interval='@daily') def hello(): @task() def get_name(): return 'Ashish' # 将multiple_outputs移到@task装饰器中 @task(multiple_outputs=True) def get_fullname(): return {'firstname': 'Ashish', 'lastname': 'Inamdar'} @task() def greet(name): print(f"Hello {name}") @task() def greet_multiple(firstname, lastname): print(f"Hi {firstname} {lastname}") name = get_name() greet(name=name) fullname_dict = get_fullname() greet_multiple(fullname_dict['firstname'], fullname_dict['lastname']) # 也可以用解构的方式传递参数 # greet_multiple(**fullname_dict) taskflow_api_dag = hello()
补充说明
- 启用
multiple_outputs=True后,Airflow会自动把返回字典的每个键作为XCom的键存储对应的值 - 下游任务除了用
fullname_dict['firstname']的方式访问,还可以通过解构参数greet_multiple(**fullname_dict)来传递,代码更简洁
内容的提问来源于stack exchange,提问作者Ashish Inamdar
相关产品推荐
相关产品推荐

