You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何访问XComArg中存储的字典?Airflow任务传参问题

如何通过XComArg从JSON中提取特定键值并传递给Airflow任务

问题描述

我希望从XComArg传递的JSON中获取特定键对应的值,并将其作为参数传递给另一个任务。以下是相关代码:

from airflow.decorators import dag, task

@dag(schedule=None)
def multiple_outputs():

    @task(multiple_outputs=True)
    def extract_employee():
        return {
            "employee": {
                "name": "John",
                "surname": "Doe",
                "age": 34 
            }
        }
    
    @task
    def transform_name(name):
        print(name)

    result = extract_employee()
    # transform_name(result["employee"]["name"])
    result >> transform_name("{{ ti.xcom_pull(task_ids='extract_employee')['employee']['name'] }}")
        
multiple_outputs()

使用Jinja模板的方式可以正常工作,但我想要通过XComArg实现相同效果。代码中注释的行向transform任务传递的是None而非员工姓名,请问是否有解决该问题的方法?

解决方案

出现这个问题的核心原因是你给extract_employee任务设置了multiple_outputs=True,此时Airflow会将任务返回的字典的顶级键拆分为独立的XCom条目。也就是说,extract_employee推送的XCom中,只有一个键为employee的条目,对应的值是{"name": "John", ...}这个子字典,但在DAG解析阶段,result["employee"]["name"]无法直接解析出实际运行时的具体值,导致传递给transform_name的是None。

你可以通过以下几种方式解决:

方法1:移除multiple_outputs=True

如果不需要将顶级键拆分为独立XCom,直接移除该参数,让整个字典作为单个XCom推送,这样就可以直接通过XComArg的键索引获取嵌套值:

from airflow.decorators import dag, task

@dag(schedule=None)
def multiple_outputs():

    @task  # 移除multiple_outputs=True
    def extract_employee():
        return {
            "employee": {
                "name": "John",
                "surname": "Doe",
                "age": 34 
            }
        }
    
    @task
    def transform_name(name):
        print(name)

    result = extract_employee()
    transform_name(result["employee"]["name"])  # 现在可以正常获取name值
        
multiple_outputs()

方法2:使用.map()提取嵌套值(保留multiple_outputs=True)

如果需要保留multiple_outputs=True,可以通过XComArg的.map()方法对返回的子字典进行处理,提取出name字段:

from airflow.decorators import dag, task

@dag(schedule=None)
def multiple_outputs():

    @task(multiple_outputs=True)
    def extract_employee():
        return {
            "employee": {
                "name": "John",
                "surname": "Doe",
                "age": 34 
            }
        }
    
    @task
    def transform_name(name):
        print(name)

    result = extract_employee()
    # 用map从employee字典中提取name
    transform_name(result["employee"].map(lambda emp: emp["name"]))
        
multiple_outputs()

方法3:拆分更细的输出(可选)

如果后续还需要用到surname、age等字段,可以直接在extract_employee中将这些字段作为顶级键返回,这样multiple_outputs=True会自动将它们拆分为独立XCom,直接通过result["name"]获取:

from airflow.decorators import dag, task

@dag(schedule=None)
def multiple_outputs():

    @task(multiple_outputs=True)
    def extract_employee():
        employee_data = {
            "name": "John",
            "surname": "Doe",
            "age": 34 
        }
        # 返回包含所有需要字段的顶级字典
        return {"name": employee_data["name"], "surname": employee_data["surname"], "age": employee_data["age"]}
    
    @task
    def transform_name(name):
        print(name)

    result = extract_employee()
    transform_name(result["name"])  # 直接获取name
        
multiple_outputs()

内容的提问来源于stack exchange,提问作者Przemek Krysztofiak

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.22 22:09:57