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

能否基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 01:32:33