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

Airflow动态映射任务组:移除子任务依赖与解析映射输入

解决Airflow动态Task Group映射中的两个技术问题

问题1:移除get_files到process_file_step2的多余依赖连线

当前DAG出现多余依赖的核心原因是process_file_step2直接引用了Task Group外部传入的file_path参数,Airflow会自动为参数的来源创建依赖关系。要消除这个依赖,需让process_file_step2的输入仅来自process_file_step1,而非直接使用外部参数。

修改思路:

  • 让process_file_step1返回包含file_path的结构化结果(比如字典)
  • 让process_file_step2接收process_file_step1的输出,从中提取file_path使用

问题2:在Task Group内解析file_path获取文件名并供下游使用

可以在Task Group内部新增专门的解析任务,或者在现有任务中完成解析,再将文件名作为Task Group的返回值,下游任务即可直接获取该值。

修改后的完整代码

from airflow.decorators import dag, task_group, task
from airflow.operators.empty import EmptyOperator
from airflow.utils.dates import days_ago
from pathlib import Path

@dag(
    start_date=days_ago(1),
    schedule=None,
    catchup=False
)
def dynamic_task_group_mapping():
    @task_group(group_id="process_file")
    def tg1(file_path):
        @task
        def process_file_step1(file_path):
            # 处理文件并返回包含file_path的结构化结果
            processed_content = f"Step 1 processed {file_path}"
            return {"processed_result": processed_content, "source_path": file_path}

        @task
        def process_file_step2(step1_output):
            # 从step1的输出中提取file_path,不再直接依赖外部参数
            target_path = step1_output["source_path"]
            return f"Step 2 processed {target_path}"

        @task
        def parse_file_name(file_path):
            # 解析file_path获取文件名
            return Path(file_path).name

        # 定义任务依赖关系
        step1_task = process_file_step1(file_path)
        step2_task = process_file_step2(step1_task)
        file_name_task = parse_file_name(file_path)
        
        # 将文件名作为Task Group的输出,供下游任务调用
        return {"step2_result": step2_task, "file_name": file_name_task}

    @task
    def get_files():
        return ['file1.txt', 'file2.csv', 'file3.log']

    # 展开Task Group并获取返回结果
    tg_expanded_result = tg1.expand(file_path=get_files())
    # 下游任务可通过tg_expanded_result访问解析后的文件名
    tg_expanded_result >> EmptyOperator(task_id='Notify', trigger_rule="none_failed_min_one_success")

dynamic_task_group_mapping()

关键说明

  1. 消除多余依赖:
    • process_file_step2不再直接接收外部file_path,仅依赖process_file_step1的输出,此时DAG图中只会保留process_file_step1 -> process_file_step2的依赖连线,get_files与process_file_step2的多余依赖会自动消失。
  2. 解析并传递文件名:
    • 新增parse_file_name任务专门处理文件名解析,也可整合到process_file_step1中完成解析并返回。
    • 通过Task Group的return语句将文件名作为输出,下游任务可通过Task Group的展开结果直接获取该值。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 08:43:14