一、场景:Airflow任务“睡死”没人喊醒的糟心事
做过定时任务调度的朋友,大概率碰过这种闹心事儿:早上打开任务看板,发现本该凌晨跑的报表任务、本该每小时同步的用户数据任务,全都卡在“待执行”或者“运行中”没进展的状态。更坑的是,前一天明明刚确认过调度规则没问题,也没改代码,怎么突然就集体“睡死”了?
这种情况在Airflow(目前最火的开源调度工具)里特别常见。你以为是任务本身报错了?翻日志看,任务连启动的机会都没有;你以为是调度器挂了?Airflow后台还能正常登录,其他无关的测试任务反而能正常跑。最烦的是,这种问题不是必现的,有时候熬一晚上自己好了,有时候得人工手动点“唤醒”(也就是触发任务)才动,完全摸不着规律。
我之前负责的用户数据同步系统,就踩过这个大坑:每天凌晨2点要跑全量用户数据同步,结果连续3天,任务都卡在“待执行”状态,直到运维同学早上8点手动触发才跑,导致下游的报表任务全乱了。后来排查发现,核心问题出在Airflow的调度逻辑死循环,以及后台数据库的锁冲突上。
二、核心问题拆解:调度器死循环和数据库锁等待
要解决任务睡死的问题,得先搞懂Airflow的调度逻辑,以及为什么会出现死循环和锁等待。
2.1 Airflow调度的基本逻辑
Airflow的调度逻辑其实不复杂:
- 调度器(Scheduler)会每隔一段时间(默认是1秒)去扫一遍所有的DAG(也就是任务流);
- 对每个DAG,检查它的下一次执行时间是不是到了;
- 如果到了,就把这个DAG对应的任务状态改成“待执行”,然后分配给执行器(Executor)去跑。
正常情况下,这个循环很顺:扫→检查→分配→扫→检查→分配。但就像人循环做事会出错一样,调度器的这个循环也可能出问题。
2.2 调度器死循环的成因
调度器死循环,简单说就是它卡在了某一步,反复做同一件事,没法往下走,导致新的任务分配不出去。
举个最常见的例子:假设你有一个DAG,设置了“一旦任务失败就自动重试”,而且重试次数设成了10次,每次重试间隔5分钟。如果这个DAG的某个任务因为依赖的接口挂了,连续失败,调度器会反复触发“检查任务状态→发现失败→触发重试”的循环,完全顾不上扫其他DAG。
我之前踩的坑就是这个:有个用户数据同步的DAG,因为第三方接口偶尔会返回500错误,我给它设了10次重试,结果某天接口连续挂了30分钟,调度器就卡在这个DAG的重试循环里,扫其他DAG的动作完全停了,导致所有依赖这个DAG的下游任务,还有其他独立的DAG,都没法被调度。
2.3 数据库锁等待的成因
Airflow的调度器需要把任务状态、DAG信息这些数据存在后台数据库里(默认是PostgreSQL或者MySQL),这就会涉及到数据库的锁。
数据库锁的作用很简单:防止两个操作同时改同一条数据,导致数据乱掉。比如调度器要把任务状态改成“待执行”,执行器要把任务状态改成“运行中”,如果两个操作同时来,数据库就会给其中一个加锁,等另一个操作做完再继续。
但如果锁的时间太长,就会出问题。比如调度器要更新一个DAG的下一次执行时间,这个操作需要先锁这个DAG对应的数据库行。如果某个慢查询(比如统计所有任务状态的报表查询)长时间占着这个锁,调度器就会一直等,等的时间久了,调度逻辑就会卡住,新任务自然没法被分配。
我之前还碰到过一次:运维同学为了排查任务延迟,跑了一个全量统计任务的SQL,这个SQL跑了15分钟才结束,结果在这15分钟里,所有DAG的调度都停了,因为那个SQL占着核心的DAG表锁,调度器没法更新DAG的执行状态。
三、检测手段:怎么快速找到“睡死”的根源
知道了问题的成因,接下来就是怎么检测。我总结了一套实用的检测流程,从易到难,一步步排查。
3.1 第一步:快速定位是不是调度器死循环
调度器死循环的最大特点是:部分任务正常,部分任务卡住。所以第一步要先做个简单的测试。
3.1.1 测试方法
新建一个最简单的测试DAG,只做一件事:打印“我是测试任务”,不设任何重试,不依赖任何外部资源。然后给它设成每分钟跑一次。
技术栈:Airflow 2.5.0(本次所有示例统一用这个版本)
# 测试DAG代码,存成test_dag.py
from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime, timedelta
# 定义默认参数
default_args = {
'owner': 'test',
'retries': 0, # 设为0,不重试,避免干扰
'retry_delay': timedelta(minutes=1),
}
# 定义测试任务的函数
def test_task():
print("我是测试任务,正常执行")
# 定义DAG
with DAG(
'test_scheduler', # DAG唯一ID
default_args=default_args,
description='测试调度器是否正常',
schedule_interval='* * * * *', # 每分钟跑一次
start_date=datetime(2024, 1, 1), # 启动时间,设成过去的时间
catchup=False, # 不补跑之前的任务
tags=['test'],
) as dag:
# 定义任务
task = PythonOperator(
task_id='test_task',
python_callable=test_task,
)
3.1.2 结果判断
部署这个DAG后,等2分钟,去Airflow后台看这个测试任务的状态:
- 如果每分钟都正常跑,说明调度器的核心逻辑是好的,死循环的概率不大;
- 如果测试任务也卡住,那调度器大概率是出问题了,比如卡在某个循环里。
3.2 第二步:抓调度器的循环日志
如果测试任务也卡住,接下来要抓调度器的日志,看它到底在反复做什么。
Airflow的调度器日志会记录它的每一步操作,我们可以用命令过滤出关键的日志内容。
技术栈:Airflow 2.5.0,Linux系统
# 查看调度器的日志,过滤出“检查DAG”和“触发任务”的关键词
# 假设调度器的日志文件路径是/var/log/airflow/scheduler.log
grep -E "Checking DAG|Triggering task" /var/log/airflow/scheduler.log | tail -n 20
3.2.1 结果判断
正常情况下,日志里会交替出现不同DAG的检查记录,比如:
[2024-05-20 09:00:00,123] {scheduler_job.py:XXX} INFO - Checking DAG: test_scheduler
[2024-05-20 09:00:00,124] {scheduler_job.py:XXX} INFO - Checking DAG: user_sync_dag
[2024-05-20 09:00:00,125] {scheduler_job.py:XXX} INFO - Triggering task: test_scheduler.test_task
如果日志里反复出现同一个DAG的检查和触发记录,比如连续10次都是:
[2024-05-20 09:00:00,123] {scheduler_job.py:XXX} INFO - Checking DAG: user_sync_dag
[2024-05-20 09:00:00,124] {scheduler_job.py:XXX} INFO - Triggering task: user_sync_dag.sync_task
[2024-05-20 09:00:00,125] {scheduler_job.py:XXX} INFO - Checking DAG: user_sync_dag
[2024-05-20 09:00:00,126] {scheduler_job.py:XXX} INFO - Triggering task: user_sync_dag.sync_task
那基本可以确定,调度器卡在了这个DAG的重试循环里,也就是死循环了。
3.3 第三步:检测数据库锁等待
如果测试任务正常跑,只是部分任务卡住,那大概率是数据库锁等待的问题。接下来要检测数据库的锁。
3.3.1 针对PostgreSQL的锁检测
Airflow默认用PostgreSQL,我们可以用PostgreSQL自带的命令查看锁等待情况。
技术栈:PostgreSQL 14.0(Airflow默认支持的版本)
-- 查看当前所有的锁等待情况,执行这个SQL(可以在psql里跑)
SELECT
lock.pid, -- 进程ID
lock.relation::regclass AS locked_table, -- 被锁的表名
lock.locktype, -- 锁的类型
lock.mode, -- 锁的模式
activity.query AS waiting_query -- 等待锁的SQL
FROM
pg_locks lock
JOIN
pg_stat_activity activity ON lock.pid = activity.pid
WHERE
lock.granted = false; -- 只看没拿到锁的等待状态
3.3.2 结果判断
如果这个SQL返回了结果,说明有进程在等锁。比如返回的结果里,locked_table是dag(Airflow核心的DAG表),waiting_query是类似“UPDATE dag SET next_execution_date = ...”的语句,那就是调度器在等这个表的锁,导致调度卡住。
3.3.3 针对MySQL的锁检测
如果Airflow用的是MySQL,检测方法类似:
技术栈:MySQL 8.0(Airflow支持的版本)
-- 查看当前的锁等待情况
SELECT
OBJECT_NAME, -- 被锁的表名
LOCK_TYPE, -- 锁的类型
LOCK_MODE, -- 锁的模式
PROCESSLIST_ID, -- 进程ID
PROCESSLIST_INFO -- 等待锁的SQL
FROM
performance_schema.data_locks
WHERE
LOCK_STATUS = 'WAITING'; -- 只看等待状态
结果判断和PostgreSQL一样,如果返回了dag表的锁等待,就是调度器在等锁。
3.4 第四步:定位占锁的“罪魁祸首”
找到锁等待后,还要知道是谁占着锁不放。
3.4.1 PostgreSQL的占锁进程查询
-- 查看占着锁的进程信息
SELECT
activity.pid, -- 进程ID
activity.query AS holding_query, -- 占锁的SQL
activity.state, -- 进程状态
activity.backend_start -- 进程启动时间
FROM
pg_locks lock
JOIN
pg_stat_activity activity ON lock.pid = activity.pid
WHERE
lock.granted = true -- 只看已经拿到锁的进程
AND lock.relation::regclass = 'dag'::regclass; -- 只看dag表的锁
3.4.2 MySQL的占锁进程查询
-- 查看占着锁的进程信息
SELECT
PROCESSLIST_ID, -- 进程ID
PROCESSLIST_INFO, -- 占锁的SQL
STATE, -- 进程状态
TIME -- 进程已经运行的时间(秒)
FROM
performance_schema.data_locks
JOIN
information_schema.PROCESSLIST ON data_locks.PROCESSLIST_ID = PROCESSLIST.ID
WHERE
LOCK_STATUS = 'GRANTED' -- 只看已经拿到锁的进程
AND OBJECT_NAME = 'dag'; -- 只看dag表的锁
3.4.3 结果判断
如果返回的占锁SQL是一个统计查询(比如“SELECT * FROM dag JOIN task_instance ON ...”),或者是一个长时间运行的业务任务的SQL,那这个就是导致锁等待的罪魁祸首。
四、应用场景、优缺点和注意事项
4.1 应用场景
这套检测方法适用于所有Airflow调度的场景,尤其是:
- 有大量DAG的生产环境:比如同时跑上百个DAG的大数据平台,很容易出现调度死循环和锁等待;
- 任务依赖复杂的场景:比如DAG之间有上下游依赖,一个DAG出问题会影响一片;
- 数据库压力大的场景:比如经常跑全量统计、备份的Airflow系统,容易出现锁冲突。
4.2 检测方法的优缺点
4.2.1 优点
- 简单易操作:不需要额外安装工具,用Airflow自带的日志和数据库自带的命令就能检测;
- 定位准确:从测试调度器到抓日志、查锁,一步步缩小范围,能快速找到问题根源;
- 成本低:不需要修改Airflow的核心配置,也不需要加额外的监控,适合中小团队快速排查问题。
4.2.2 缺点
- 被动检测:只能在问题出现后再排查,没法提前预警;
- 对复杂场景覆盖不全:如果是多个DAG同时出现死循环,或者锁等待的原因是多个进程共同导致的,检测起来会比较麻烦;
- 依赖人工操作:需要运维或开发同学手动执行命令,没法自动处理。
4.3 注意事项
- 测试DAG要尽量简单:不能设重试、不能依赖外部资源,避免干扰测试结果;
- 抓日志的时候要注意时间:如果调度器死循环的问题是偶现的,要在问题出现的时候马上抓日志,不然日志会被覆盖;
- 查数据库锁的时候要注意权限:需要用有超级用户权限的账号(比如PostgreSQL的postgres用户,MySQL的root用户),不然看不到所有的锁信息;
- 杀占锁进程要谨慎:如果占锁的进程是正常的业务任务,杀了可能会导致数据不一致,最好先等任务跑完,或者调整任务的执行时间(比如把全量统计任务改成凌晨低峰期跑)。
五、文章总结
Airflow任务睡死的问题,本质上是调度逻辑的异常循环和数据库锁冲突导致的。要解决这个问题,核心是先定位问题:先测试调度器是不是正常,再抓日志看是不是死循环,最后查数据库锁是不是有冲突。
我总结的这套检测方法,从易到难,不需要复杂的工具,适合不同基础的开发者操作。只要按照步骤一步步排查,就能快速找到问题的根源,避免任务睡死导致的业务故障。
最后,给大家提个小建议:如果你的Airflow系统已经稳定运行,最好提前加一些监控,比如监控调度器的循环时间、监控数据库的锁等待时间,这样能提前发现问题,避免出故障后再手忙脚乱。
评论
围绕“Airflow任务挂起无人唤醒的场景丛生,深挖调度器死循环与数据库锁等待的检测手段”参与讨论