Airflow 2 Taskflow API 如何使用自定义命名key推送XCom
Airflow 2 Taskflow API 自定义XCom键实现方案
方法1:任务内部主动调用xcom_push指定自定义键
该方法和传统XCom操作逻辑完全一致,只需在@task装饰的函数内获取当前任务上下文即可调用xcom_push自定义键名:
import json import requests from airflow.decorators import dag, task from airflow.operators.python import get_current_context from datetime import datetime @dag( schedule=None, start_date=datetime(2023, 1, 1), catchup=False ) def demo_dag(): @task(task_id="task_one") def get_height(): response = requests.get("https://swapi.dev/api/people/4") data = json.loads(response.text) height = int(data["height"]) # 获取当前任务实例对象 ti = get_current_context()["ti"] # 自定义XCom键名存储值 ti.xcom_push(key="character_height", value=height) @task(task_id="task_two") def check_height(): ti = get_current_context()["ti"] # 拉取指定任务、指定键的XCom值 val = ti.xcom_pull(task_ids="task_one", key="character_height") print(f"Value passed in is: {val}") # 定义任务依赖 get_height() >> check_height() demo_dag()
方法2:使用multiple_outputs自动生成自定义XCom键(Airflow 2.0+支持)
如果任务需要返回多个值存为不同的XCom键,可以开启multiple_outputs参数,函数返回字典的键会自动作为XCom的自定义键,无需手动调用push方法:
import json import requests from airflow.decorators import dag, task from datetime import datetime @dag( schedule=None, start_date=datetime(2023, 1, 1), catchup=False ) def demo_dag(): @task(task_id="task_one", multiple_outputs=True) def get_character_info(): response = requests.get("https://swapi.dev/api/people/4") data = json.loads(response.text) # 返回字典的key会自动成为XCom的独立键 return { "character_height": int(data["height"]), "character_name": data["name"] } @task(task_id="task_two") def check_height(height_val): print(f"Value passed in is: {height_val}") # 直接提取指定键的值传入下游任务 character_info = get_character_info() check_height(character_info["character_height"]) demo_dag()
说明
- 方法1适合仅需要自定义单个XCom键的场景,和传统写法兼容性更高
- 方法2适合多返回值的场景,代码更简洁,符合Taskflow的函数式编程风格
内容的提问来源于stack exchange,提问作者Indrid
相关产品推荐
相关产品推荐

