如何在Airflow中配置短于1分钟间隔的DAG 实现30秒调度周期
Airflow 30秒DAG调度间隔实现方案
Airflow原生默认DAG最小调度间隔为1分钟,但可以通过两种成熟方案实现30秒调度需求:
方案1:双DAG错峰触发(全版本通用,生产环境首选)
该方案不需要修改Airflow任何核心配置,稳定性最高,是生产环境最常用的实现方式:
- 原有DAG保持
schedule_interval = "*/1 * * * *"(每分钟调度一次)不变,正常执行业务逻辑 - 复制一份业务逻辑完全相同的新DAG,在所有业务任务前增加一个30秒等待的前置任务,示例如下:
# 第二个错峰DAG的前置等待任务 wait_30s = BashOperator( task_id="wait_30s", bash_command="sleep 30" ) - 两个DAG均为每分钟调度一次,错峰30秒执行,即可实现整体每30秒执行一次业务逻辑的效果
- 注意事项:如果业务逻辑非幂等,需要加分布式锁避免两个DAG同时执行业务逻辑时产生脏数据
方案2:秒级cron表达式(适用于Airflow 2.3及以上版本)
Airflow 2.3版本开始原生支持6位带秒的cron表达式,可以直接配置30秒调度间隔:
- 先确认
airflow.cfg中scheduler_heartbeat_sec配置项的值不大于10(默认值为5,无需修改即可满足需求),同时调整min_file_process_interval = 10避免DAG文件扫描过于频繁占用调度器资源 - 直接将DAG的
schedule_interval设置为6位秒级cron表达式*/30 * * * * *即可,示例代码如下:from datetime import datetime, timedelta from airflow import DAG from airflow.operators.bash import BashOperator default_args = { "owner": "airflow", "retries": 1, "retry_delay": timedelta(seconds=10) } with DAG( dag_id="30s_schedule_demo", default_args=default_args, start_date=datetime(2024, 1, 1), schedule_interval="*/30 * * * * *", # 6位cron,首位为秒,代表每30秒执行一次 catchup=False ) as dag: business_task = BashOperator( task_id="run_business_logic", bash_command="your business command" ) - 注意事项:需要确保单轮业务任务执行时长远小于30秒,避免出现任务堆叠;调度器负载较高时可能出现几秒的调度延迟,对调度实时性要求极高的场景谨慎使用。
内容的提问来源于stack exchange,提问作者Devendra Reddy
相关产品推荐
相关产品推荐

