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

在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)

关键特性说明

  1. 确保任务仅执行一次

    • catchup=False:Airflow不会回溯执行DAG启动前的历史调度任务
    • 固定时间Cron调度:每个自然日仅触发一次任务,成功后状态标记为「success」,不会重复运行
    • 失败自动重试:按照retries和retry_delay配置自动重试,直到成功或达到重试上限
  2. 多策略动态支持

    • 通过fetch_policies从数据库获取所有高峰时段策略
    • 循环为每个策略生成独立的降带宽/恢复带宽DAG,便于单独管理和监控
  3. 时区一致性

    • 所有时间配置统一使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 15:45:34