MWAA Airflow 2.2.2传感器任务被误标记失败的问题排查
问题定位
你遇到的是Airflow 2.2.2多调度器架构搭配Celery执行器的典型状态竞争bug:
- 重调度型Sensor每5分钟会释放Worker槽位并重新进入队列,多个调度器实例可能同时读取到任务的
queued状态 - 其中一个调度器会错误判定该任务已完成(实际是重调度的状态流转异常),直接将任务标记为
failed,但Worker上的Sensor进程并未收到终止信号,仍会持续轮询API直到拿到预期响应
临时缓解方案(无需重启MWAA)
- 调大Sensor的
reschedule_timeout参数:设为7分钟(比轮询周期长2分钟),降低调度器误判任务超时的概率 - 临时切为单调度器:在MWAA控制台把调度器数量改成1,切换过程只有短暂调度延迟,不用等20-30分钟的停机
- 优化
on_failure_handler:在Slack通知里加个校验逻辑——调用Airflow API查任务的latest_heartbeat,如果心跳在最近5分钟内,就标注「疑似状态误报,任务仍在运行」
根因排查要点
- 拉取CloudWatch中所有调度器实例的日志,对比同一任务的状态更新时间戳,确认是否有多个调度器同时修改
task_instance状态 - 查MWAA关联的RDS元数据库锁日志,看是否存在分布式锁超时或未正确释放的情况(多调度器依赖元数据库锁避免状态冲突)
- 检查Sensor的
poke方法:如果poke里有未捕获的异常,会导致调度器误判状态,给poke加详细日志,记录每次轮询的返回值和异常信息
长期修复建议
- 升级MWAA版本:Airflow 2.3及以上版本修复了多调度器重调度Sensor的状态同步bug,AWS后续的MWAA版本已整合该修复
- 替换重调度Sensor为
AsyncSensor:异步Sensor不用释放Worker槽位,从根源避免重调度带来的状态流转问题,更适配多调度器架构 - 配置调度器
task_queued_timeout:设为大于Sensor轮询周期的值,防止调度器过早把queued状态的任务标记为失败
内容的提问来源于stack exchange,提问作者evzdfx
相关产品推荐
相关产品推荐

