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

Apache Airflow on Composer2.0任务报'Not yet started'解决方法

问题背景
  • 异常现象:单DAG包含600余个任务,运行过程中偶发部分任务执行失败,任务状态异常显示为Not yet started
  • 运行环境配置:
    • Composer版本:2.0.10
    • Airflow版本:2.2.3
    • 执行器:CeleryExecutor
    • Worker节点:2~6节点自动扩缩容
    • Scheduler节点:2节点高可用部署
根因排查路径

1. 双Scheduler竞态锁问题

Airflow 2.2.3版本存在已确认的调度器竞态缺陷:双Scheduler同时运行时,两个进程可能同时扫描到同一个待调度任务实例,触发数据库行锁争抢,抢锁失败的调度器不会重试派发任务,也不会更新任务状态,最终任务卡在Not yet started状态直到DAG运行超时。
排查操作:

  • 拉取两个Scheduler节点的服务日志,检索关键词Task instance lock not acquired、Failed to dispatch task,匹配异常出现时间点的日志条目,确认是否存在锁争抢记录
  • 核对配置项scheduler_heartbeat_sec、task_queued_timeout:若心跳间隔大于5s、任务排队超时阈值小于任务平均排队时长,调度器会误判未派发的任务为丢失状态,跳过后续处理流程

2. Worker自动扩缩容丢任务

CeleryExecutor模式下,Worker节点缩容时如果未等已接收、未ACK的任务执行完成就被销毁,会导致已派发的任务直接丢失,且Worker无法向元数据库回传任务状态,任务会一直停留在Not yet started。
排查操作:

  • 查看Composer底层集群的扩缩容事件,确认任务异常时间点是否刚好触发Worker节点缩容
  • 核对Celery配置worker_prefetch_multiplier:若参数值大于1,单个Worker会预取多个待执行任务,缩容时预取但未启动的任务会全部丢失
  • 核对Worker配置termination_grace_period_seconds:若优雅终止等待时间小于300s,长任务未执行完就会被强制杀死

3. 元数据库连接压力过载

单DAG 600+任务的场景下,双Scheduler并发扫描、更新任务状态会给元数据库带来较高负载,若出现连接池耗尽、SQL执行超时,会出现任务已派发但状态未更新为queued的情况,后续调度逻辑会直接跳过该任务,最终状态卡在Not yet started。
排查操作:

  • 查看元数据库监控指标,确认异常时间点是否存在连接数打满、慢查询占比过高、CPU使用率超过80%的情况
  • 核对数据库连接池配置sql_alchemy_pool_size、sql_alchemy_max_overflow:若连接池总大小小于20,高负载场景下极易出现连接耗尽问题

4. Celery Broker消息异常

Composer 2.0.x默认使用Redis作为Celery Broker,若Broker内存不足、消息过期时间配置过短,排队中的任务消息会被自动淘汰,调度器侧记录任务已入队,但Broker侧无对应消息,Worker永远无法接收任务,状态会停留在初始的Not yet started。
排查操作:

  • 查看Redis实例的内存使用率、消息淘汰计数指标,确认异常时间点是否存在内存占满触发淘汰、消息过期被清理的记录
  • 核对Celery配置broker_transport_options中的visibility_timeout参数:若该值小于任务最长执行时长,会触发任务重复派发、原任务状态锁失效的问题
解决方案
  • 针对Scheduler竞态问题:
    1. 优先将Airflow版本升级至2.2.5及以上的2.2.x稳定版本,该版本已修复双Scheduler场景下任务锁争抢导致的状态异常缺陷;若暂时无法升级,可临时将Scheduler节点数调整为1规避竞态,该操作会牺牲调度层高可用能力
    2. 调整调度器配置:将scheduler_heartbeat_sec设为5,task_queued_timeout设为600,给任务排队预留足够超时窗口;将scheduler.max_tis_per_query从默认值512调整为128,降低单次批量更新任务给元数据库带来的压力
  • 针对Worker扩缩容丢任务问题:
    1. 调整Worker扩缩容策略,配置缩容时优先驱逐空闲节点;将worker_prefetch_multiplier设为1,关闭Worker预取多任务特性,避免单节点绑定过多待执行任务
    2. 将Worker的termination_grace_period_seconds调整为600以上,配置Celery参数worker_shutdown_timeout = 300,确保节点缩容时能等待已接收任务执行完成再退出
    3. 配置Celery任务ACK机制为task_acks_late = True,任务执行完成后再向Broker返回ACK,避免Worker异常退出时任务直接丢失
  • 针对元数据库压力问题:
    1. 扩容元数据库规格,将sql_alchemy_pool_size设为10,sql_alchemy_max_overflow设为20,确保调度器、Worker有足够的数据库连接可用
    2. 给大DAG配置合理的并发上限,单DAG的max_active_tasks不要超过200,避免单次DAG运行同时触发数百个任务,给数据库和Broker造成流量冲击
  • 针对Broker异常问题:
    1. 扩容Redis实例内存规格,确保峰值场景下内存使用率不超过70%;将Redis淘汰策略调整为noeviction,避免任务消息被意外清理
    2. 将visibility_timeout调整为比单任务最长执行时间多300s,避免任务还在执行就被Broker判定为超时重新入队,导致状态错乱

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 01:48:15