Airflow 2.3中能否实现DAG的多次动态expand操作?
Airflow 多次.expand() 嵌套场景的实现方案
你当前的代码逻辑里,transform.expand()会生成5个独立的transform任务,每个任务返回一个列表。但直接把这个结果传给load.expand()时,Airflow只会把每个transform任务的输出当作单个输入项,所以load只会生成5个任务(每个接收一个列表),不会自动展开嵌套的列表——这种原生嵌套expand的方式不可行。
要实现二次展开生成15个独立打印任务,用Airflow提供的.flatten()方法就能解决,它可以把多个任务返回的嵌套列表合并成一维列表。修改后的完整代码如下:
import typing from airflow.decorators import task, dag from datetime import datetime @dag(schedule=None, start_date=datetime(2024, 1, 1), catchup=False) def nested_expand_dag(): @task def extract(): return ["1", "2", "3", "4", "5"] @task def transform(value: str) -> typing.List[str]: return [value] * int(value) @task def load(value: str): print(value) order_data = extract() transformed_results = transform.expand(value=order_data) # 用.flatten()扁平化所有transform任务的输出结果 load.expand(value=transformed_results.flatten()) nested_expand_dag()
.flatten()会把5个transform任务返回的列表(["1"]、["2","2"]、["3","3","3"]等)合并成一个包含15个元素的一维列表,这样load.expand()就能基于这个列表生成15个独立的打印任务,每个任务接收单个字符串值。
内容的提问来源于stack exchange,提问作者Pwnergod
相关产品推荐
相关产品推荐

