如何访问XComArg中存储的字典?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

