在Airflow中创建条件任务:ISP带宽调度需求实现
解决方案:基于Airflow实现ISP带宽调度策略
核心思路
要实现「到点仅执行一次」的带宽调整任务,关键利用Airflow的精确调度触发和任务状态持久化特性:
- 降带宽任务:设置为每天高峰开始时间(如09:00)触发,执行成功后标记为完成,不再重复运行
- 恢复带宽任务:设置为每天高峰结束时间(如13:00)触发,同理仅执行一次
- 支持从数据库动态加载多时段策略,自动生成对应调度任务
改造后的完整代码
from __future__ import annotations from datetime import datetime, timedelta from airflow import DAG from airflow.operators.python import PythonOperator import pendulum # 从数据库获取带宽高峰策略 def fetch_policies(): # 实际环境替换为数据库查询逻辑 policies = [ {'id': 1, 'peak_start': '09:00', 'peak_end': '13:00'}, {'id': 2, 'peak_start': '17:00', 'peak_end': '21:00'} ] return policies # 降带宽执行函数 def decrease_bandwidth(policy_id): # 替换为实际带宽调整逻辑 print(f"[策略{policy_id}] 执行降带宽操作,当前时间: {datetime.now()}") # 恢复带宽执行函数 def return_to_normal_bandwidth(policy_id): # 替换为实际带宽恢复逻辑 print(f"[策略{policy_id}] 执行带宽恢复操作,当前时间: {datetime.now()}") # 通用默认参数 default_args = { 'owner': 'isp_admin', 'depends_on_past': False, 'start_date': pendulum.datetime(2024, 1, 1, tz="Europe/Istanbul"), 'email_on_failure': False, 'email_on_retry': False, 'retries': 2, # 失败重试次数 'retry_delay': timedelta(minutes=5), } # 生成降带宽DAG:每天高峰开始时间触发 def create_decrease_dag(policy): dag_id = f"decrease_bandwidth_policy_{policy['id']}" # 转换为Cron表达式:分 时 日 月 周 schedule = f"0 {policy['peak_start'].split(':')[0]} * * *" dag = DAG( dag_id=dag_id, default_args=default_args, catchup=False, # 禁止回溯执行历史日期任务 tags=["isp_bandwidth", "decrease"], schedule_interval=schedule, ) with dag: PythonOperator( task_id=f"decrease_bw_policy_{policy['id']}", python_callable=decrease_bandwidth, op_kwargs={'policy_id': policy['id']}, ) return dag # 生成恢复带宽DAG:每天高峰结束时间触发 def create_restore_dag(policy): dag_id = f"restore_bandwidth_policy_{policy['id']}" schedule = f"0 {policy['peak_end'].split(':')[0]} * * *" dag = DAG( dag_id=dag_id, default_args=default_args, catchup=False, tags=["isp_bandwidth", "restore"], schedule_interval=schedule, ) with dag: PythonOperator( task_id=f"restore_bw_policy_{policy['id']}", python_callable=return_to_normal_bandwidth, op_kwargs={'policy_id': policy['id']}, ) return dag # 动态加载所有策略对应的DAG policies = fetch_policies() for policy in policies: globals()[f"dag_decrease_{policy['id']}"] = create_decrease_dag(policy) globals()[f"dag_restore_{policy['id']}"] = create_restore_dag(policy)
关键特性说明
确保任务仅执行一次
catchup=False:Airflow不会回溯执行DAG启动前的历史调度任务- 固定时间Cron调度:每个自然日仅触发一次任务,成功后状态标记为「success」,不会重复运行
- 失败自动重试:按照
retries和retry_delay配置自动重试,直到成功或达到重试上限
多策略动态支持
- 通过
fetch_policies从数据库获取所有高峰时段策略 - 循环为每个策略生成独立的降带宽/恢复带宽DAG,便于单独管理和监控
- 通过
时区一致性
- 所有时间配置统一使用
Europe/Istanbul时区,避免时区偏差导致触发时间错误
- 所有时间配置统一使用
替代方案:单DAG分支处理(适合简单场景)
如果不需要多策略动态生成,单DAG也可实现需求:
from __future__ import annotations from datetime import datetime, timedelta from airflow import DAG from airflow.operators.python import PythonOperator from airflow.operators.datetime import BranchDateTimeOperator from airflow.operators.empty import EmptyOperator import pendulum default_args = { 'owner': 'isp_admin', 'depends_on_past': False, 'start_date': pendulum.datetime(2024, 1, 1, tz="Europe/Istanbul"), 'retries': 2, 'retry_delay': timedelta(minutes=5), } dag = DAG( dag_id="isp_bandwidth_scheduler", default_args=default_args, catchup=False, tags=["isp_bandwidth"], schedule_interval="@hourly", # 每小时检查一次时间窗口 ) def decrease_bandwidth(): print("执行降带宽操作") def restore_bandwidth(): print("执行带宽恢复操作") with dag: # 判断是否处于高峰时段 branch = BranchDateTimeOperator( task_id="check_peak_window", follow_task_ids_if_true=["decrease_bw"], follow_task_ids_if_false=["check_restore_time"], target_lower=pendulum.time(9, 0, 0), target_upper=pendulum.time(12, 59, 59), dag=dag, ) # 判断是否到达恢复时间点 restore_branch = BranchDateTimeOperator( task_id="check_restore_time", follow_task_ids_if_true=["restore_bw"], follow_task_ids_if_false=["do_nothing"], target_lower=pendulum.time(13, 0, 0), target_upper=pendulum.time(13, 0, 0), # 仅13:00触发恢复 dag=dag, ) decrease_task = PythonOperator( task_id="decrease_bw", python_callable=decrease_bandwidth, execution_timeout=timedelta(minutes=10), ) restore_task = PythonOperator( task_id="restore_bw", python_callable=restore_bandwidth, execution_timeout=timedelta(minutes=10), ) do_nothing = EmptyOperator(task_id="do_nothing") branch >> [decrease_task, restore_branch] restore_branch >> [restore_task, do_nothing]
注意事项
- 生产环境中,建议为带宽操作函数添加幂等性校验:执行前检查当前带宽状态,避免重复调整
- 可结合Airflow的
Variable或数据库存储任务执行状态,进一步确保不会重复执行 - 配置任务失败告警,避免高峰时段带宽调整失败导致网络过载
内容的提问来源于stack exchange,提问作者Muhammad Khaled
相关产品推荐
相关产品推荐

