一、场景:Airflow任务“睡死”没人喊醒的糟心事

做过定时任务调度的朋友,大概率碰过这种闹心事儿:早上打开任务看板,发现本该凌晨跑的报表任务、本该每小时同步的用户数据任务,全都卡在“待执行”或者“运行中”没进展的状态。更坑的是,前一天明明刚确认过调度规则没问题,也没改代码,怎么突然就集体“睡死”了?

这种情况在Airflow(目前最火的开源调度工具)里特别常见。你以为是任务本身报错了?翻日志看,任务连启动的机会都没有;你以为是调度器挂了?Airflow后台还能正常登录,其他无关的测试任务反而能正常跑。最烦的是,这种问题不是必现的,有时候熬一晚上自己好了,有时候得人工手动点“唤醒”(也就是触发任务)才动,完全摸不着规律。

我之前负责的用户数据同步系统,就踩过这个大坑:每天凌晨2点要跑全量用户数据同步,结果连续3天,任务都卡在“待执行”状态,直到运维同学早上8点手动触发才跑,导致下游的报表任务全乱了。后来排查发现,核心问题出在Airflow的调度逻辑死循环,以及后台数据库的锁冲突上。

二、核心问题拆解:调度器死循环和数据库锁等待

要解决任务睡死的问题,得先搞懂Airflow的调度逻辑,以及为什么会出现死循环和锁等待。

2.1 Airflow调度的基本逻辑

Airflow的调度逻辑其实不复杂:

  1. 调度器(Scheduler)会每隔一段时间(默认是1秒)去扫一遍所有的DAG(也就是任务流);
  2. 对每个DAG,检查它的下一次执行时间是不是到了;
  3. 如果到了,就把这个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调度的场景,尤其是:

  1. 有大量DAG的生产环境:比如同时跑上百个DAG的大数据平台,很容易出现调度死循环和锁等待;
  2. 任务依赖复杂的场景:比如DAG之间有上下游依赖,一个DAG出问题会影响一片;
  3. 数据库压力大的场景:比如经常跑全量统计、备份的Airflow系统,容易出现锁冲突。

4.2 检测方法的优缺点

4.2.1 优点

  1. 简单易操作:不需要额外安装工具,用Airflow自带的日志和数据库自带的命令就能检测;
  2. 定位准确:从测试调度器到抓日志、查锁,一步步缩小范围,能快速找到问题根源;
  3. 成本低:不需要修改Airflow的核心配置,也不需要加额外的监控,适合中小团队快速排查问题。

4.2.2 缺点

  1. 被动检测:只能在问题出现后再排查,没法提前预警;
  2. 对复杂场景覆盖不全:如果是多个DAG同时出现死循环,或者锁等待的原因是多个进程共同导致的,检测起来会比较麻烦;
  3. 依赖人工操作:需要运维或开发同学手动执行命令,没法自动处理。

4.3 注意事项

  1. 测试DAG要尽量简单:不能设重试、不能依赖外部资源,避免干扰测试结果;
  2. 抓日志的时候要注意时间:如果调度器死循环的问题是偶现的,要在问题出现的时候马上抓日志,不然日志会被覆盖;
  3. 查数据库锁的时候要注意权限:需要用有超级用户权限的账号(比如PostgreSQL的postgres用户,MySQL的root用户),不然看不到所有的锁信息;
  4. 杀占锁进程要谨慎:如果占锁的进程是正常的业务任务,杀了可能会导致数据不一致,最好先等任务跑完,或者调整任务的执行时间(比如把全量统计任务改成凌晨低峰期跑)。

五、文章总结

Airflow任务睡死的问题,本质上是调度逻辑的异常循环和数据库锁冲突导致的。要解决这个问题,核心是先定位问题:先测试调度器是不是正常,再抓日志看是不是死循环,最后查数据库锁是不是有冲突。

我总结的这套检测方法,从易到难,不需要复杂的工具,适合不同基础的开发者操作。只要按照步骤一步步排查,就能快速找到问题的根源,避免任务睡死导致的业务故障。

最后,给大家提个小建议:如果你的Airflow系统已经稳定运行,最好提前加一些监控,比如监控调度器的循环时间、监控数据库的锁等待时间,这样能提前发现问题,避免出故障后再手忙脚乱。