能否基于Airflow动态映射任务输出构建任务树?
如何基于动态映射任务的输出构建连续的动态任务树?
我希望从动态任务的输出生成动态任务,每个映射任务返回一个列表,需要为列表中的每个元素创建单独的映射任务,让流程形成任务树结构。想知道能否对动态映射任务的输出直接扩展,实现连续的映射操作,而非先映射再聚合?
环境信息
本地使用的环境:
Astronomer Runtime 9.6.0 based on Airflow 2.7.3+astro.2 Git Version: .release:9fad9363bb0e7520a991b5efe2c192bb3405b675
为测试验证,我设置了三个任务,输入为单个字符串,输出为字符串列表。
1. 在已扩展的任务组上扩展(含映射任务的组再次映射)
尝试在已扩展的任务组内嵌套映射任务,再对任务组进行扩展:
import datetime import logging from airflow.decorators import dag, task, task_group @dag(schedule_interval=None, start_date=datetime.datetime(2023, 9, 27)) def try_dag3(): @task def first() -> list[str]: return ["0", "1"] first_task = first() @task_group def my_group(input: str) -> list[str]: @task def second(input: str) -> list[str]: logging.info(f"input: {input}") result = [] for i in range(3): result.append(f"{input}_{i}") # ['0_0', '0_1', '0_2'] # ['1_0', '1_1', '1_2'] return result second_task = second.expand(input=first_task) @task def third(input: str, input1: str = None): logging.info(f"input: {input}, input1: {input1}") return input third_task = third.expand(input=second_task) my_group.expand(input=first_task) try_dag3()
运行时触发错误:
NotImplementedError: operator expansion in an expanded task group is not yet supported
2. 对已扩展任务的结果直接扩展(映射任务上再次映射)
尝试直接对前一个映射任务的输出进行二次扩展:
import datetime import logging from airflow.decorators import dag, task @dag(start_date=datetime.datetime(2023, 9, 27)) def try_dag1(): @task def first() -> list[str]: return ["0", "1"] first_task = first() @task def second(input: str) -> list[str]: logging.info(f"source: {input}") result = [] for i in range(3): result.append(f"{input}_{i}") # ['0_0', '0_1', '0_2'] # ['1_0', '1_1', '1_2'] return result # 此处正常扩展为first_task返回列表对应的两个任务 second_task = second.expand(input=first_task) @task def third(input: str): logging.info(f"source: {input}") return input # 此处未扩展,生成两个映射任务,输入为列表而非字符串 third_task = third.expand(input=second_task) try_dag1()
实际运行后,second任务的结果未被展开,third任务的输入是完整的字符串列表而非单个元素:
third[0]任务日志:[2024-01-05, 11:40:30 UTC] {try_dag1.py:30} INFO - source: ['0_0', '0_1', '0_2']
(流程显示:third仅生成2个映射任务,每个任务接收second返回的完整列表)
3. 携带常量输入对已扩展任务进行扩展(测试结构可行性)
尝试为二次映射任务添加常量列表输入,验证扩展逻辑:
import datetime import logging from airflow.decorators import dag, task @dag(start_date=datetime.datetime(2023, 9, 27)) def try_dag0(): @task def first() -> list[str]: return ["0", "1"] first_task = first() @task def second(input: str) -> list[str]: logging.info(f"input: {input}") result = [] for i in range(3): result.append(f"{input}_{i}") # ['0_0', '0_1', '0_2'] # ['1_0', '1_1', '1_2'] return result second_task = second.expand(input=first_task) @task def third(input: str, input1: str = None): logging.info(f"input: {input}, input1: {input1}") return input third_task = third.expand(input=second_task, input1=["a", "b", "c"]) try_dag0()
运行后发现,映射任务可通过input1的常量列表完成扩展,但input仍为未展开的完整列表:
third[0]任务日志:[2024-01-05, 12:51:39 UTC] {try_dag0.py:33} INFO - input: ['0_0', '0_1', '0_2'], input1: a
(流程显示:third生成3个映射任务,每个任务的input是second返回的完整列表,input1依次取常量列表中的元素)
内容的提问来源于stack exchange,提问作者zacheusz
相关产品推荐
相关产品推荐

