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

如何在Airflow 2.3.x运行时动态设置任务的max_active_tis_per_dag值?

动态调整Airflow任务的max_active_tis_per_dag值

要实现根据batch_id动态调整sanitize_payment任务的max_active_tis_per_dag参数,核心是从DAG运行配置中获取batch_id并传递到任务组内部,再根据batch_id计算对应并发限制值,最终应用到任务装饰器上。以下是修改后的完整代码:

from airflow.decorators import dag, task, task_group
from datetime import timedelta
import pendulum
import logging

@dag(
    dag_id="data_transfer_dag",
    schedule_interval=None,
    tags=["Data transfer"],
    default_args={
        "owner": "Open",
        "start_date": pendulum.today("UTC").add(days=-1),
        "depends_on_past": False,
        "email_on_failure": False,
        "email_on_retry": False,
        "email": "",
        "retries": 1,
        "retry_delay": timedelta(seconds=10),
    },
    description="Data transfer dag"
)
def data_transfer():
    
    @task()
    def fetch_batch_user_to_process(**kwargs):
        batch_id = kwargs['dag_run'].conf.get('batch_id')
        # 返回包含batch_id和数据列表的字典,方便后续任务组获取所需信息
        return {"batch_id": batch_id, "data_list": [{"batch_id": batch_id}]*100}
    
    @task_group(group_id="sanitize_payment_group")
    def sanitize_payment_group(users_data):
        # 从传入数据中提取batch_id
        batch_id = users_data["batch_id"]
        # 根据batch_id动态设置并发限制值
        max_active = 10 if batch_id == 1 else 15 if batch_id == 2 else 16

        @task(max_active_tis_per_dag=max_active)
        def sanitize_payment(data):
            """Some operation"""
            correct_api_version_data = data  # Dummy operation
            return correct_api_version_data
        
        # 使用数据列表生成动态任务
        task_result = sanitize_payment.expand(data=users_data["data_list"])
        return task_result
    
    @task_group(group_id="process_payment_group")
    def process_payment_group(users_data):
        @task(max_active_tis_per_dag=1)
        def process_payment(payment_data):
            """Some operation"""
            data = payment_data  # Dummy operation
            return data
        
        task_result = process_payment.expand(payment_data=users_data["data_list"])
        return task_result
    
    @task_group(group_id="create_contact_group")
    def create_contact_group(user_data):
        @task(max_active_tis_per_dag=16)
        def create_contact(user_info):
            """Some Operation"""
            if_contact_present = user_info  # Dummy operation 
            return if_contact_present
        
        task_result = create_contact.expand(user_info=user_data["data_list"])
        return task_result
    
    @task()
    def end_processing():
        logging.info("ending the dag.")
    
    end = end_processing()
    batch_to_process = fetch_batch_user_to_process()
    
    process_payment_group(sanitize_payment_group(create_contact_group(batch_to_process))) >> end


DAG = data_transfer()

关键改动说明:

  • 调整fetch_batch_user_to_process返回结构:将原单元素列表改为包含batch_id和data_list的字典,既保留batch_id用于并发计算,也提供生成动态任务所需的数据集合。
  • 动态计算并发限制值:在sanitize_payment_group内提取batch_id,通过条件判断生成对应max_active值,直接传入任务装饰器的max_active_tis_per_dag参数。
  • 统一任务组数据处理:各任务组调用expand时使用字典中的data_list字段,保证动态任务正常生成。

如果需要扩展更多batch_id对应的并发规则,只需在max_active的条件判断中添加新分支即可。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 04:55:15