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

如何将Airflow映射任务组中的装饰器任务替换为BashOperator

问题

现有一个简化的Airflow DAG,通过读取文本文件获取若干ID,展开任务组为每个ID执行对应操作,当前使用任务装饰器可正常运行。但生产环境需使用BashOperator替代自定义任务装饰器,现需将其中say_hello和say_bye替换为BashOperator,实现相同功能,该如何修改?

原DAG代码:

from airflow import DAG
from airflow.operators.empty import EmptyOperator
from airflow.decorators import task, task_group
from datetime import days_ago

# 假设FILE_PATH已提前定义
FILE_PATH = "/path/to/your/file.txt"

@dag(dag_id="dynamic_example_with_mapped_task_group", schedule=None, start_date=days_ago(1), catchup=False, tags=["test"])
def dynamic_example_with_mapped_task_group():
    start = EmptyOperator(task_id="start")
    
    @task
    def read_names():
        with open(FILE_PATH, "r") as file:
            entries = [line.strip() for line in file.readlines() if line.strip()]
        return entries

    # Read names
    names = read_names()

    @task
    def say_hello(entry):
        return f"Hello {entry}"

    @task
    def say_bye(entry):
        return f"Bye {entry}"

    @task_group
    def hello_bye_task_group(entry):
        say_hello(entry) >> say_bye(entry)

    end = EmptyOperator(task_id="end")

    hbtg = hello_bye_task_group.expand(entry=names)

    start >> names >> hbtg >> end
解决方案

修改核心要点

  1. 导入BashOperator替代原@task装饰的函数
  2. 用echo命令模拟原任务的输出逻辑,通过Airflow模板变量引用传入的ID参数
  3. 保留原有任务组的动态展开逻辑,确保每个ID对应一组独立的hello/bye任务

修改后完整代码

from airflow import DAG
from airflow.operators.empty import EmptyOperator
from airflow.operators.bash import BashOperator
from airflow.decorators import task, task_group
from datetime import days_ago

FILE_PATH = "/path/to/your/file.txt"

@dag(dag_id="dynamic_example_with_mapped_task_group", schedule=None, start_date=days_ago(1), catchup=False, tags=["test"])
def dynamic_example_with_mapped_task_group():
    start = EmptyOperator(task_id="start")
    
    @task
    def read_names():
        with open(FILE_PATH, "r") as file:
            entries = [line.strip() for line in file.readlines() if line.strip()]
        return entries

    names = read_names()

    @task_group
    def hello_bye_task_group(entry):
        # 替换为BashOperator,通过echo输出对应内容
        say_hello = BashOperator(
            task_id="say_hello",
            # 注意:Python字符串中需用双层大括号转义Airflow模板变量
            bash_command=f'echo "Hello {{{{ entry }}}}"'
        )
        
        say_bye = BashOperator(
            task_id="say_bye",
            bash_command=f'echo "Bye {{{{ entry }}}}"'
        )
        
        say_hello >> say_bye

    end = EmptyOperator(task_id="end")

    hbtg = hello_bye_task_group.expand(entry=names)

    start >> names >> hbtg >> end

关键细节说明

  • 模板变量转义:在Python f-string中定义bash_command时,Airflow的模板变量{{ entry }}需要写成{{{{ entry }}}},外层Python会解析一层大括号,最终传递给Airflow的是正确的模板语法。
  • 功能一致性:echo命令的输出会记录在任务日志中,和原Python任务的返回值效果一致,都能体现对每个ID的处理结果。
  • 动态映射逻辑:原有的hello_bye_task_group.expand(entry=names)完全保留,依然会为read_names返回的每个ID生成独立的任务组,保持原DAG的结构和执行逻辑不变。

内容的提问来源于stack exchange,提问作者Nélia Fonseca

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 03:55:10