如何在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
相关产品推荐
相关产品推荐

