Apache Airflow 调度器深度实战:从 DAG 序列化、Critical Section 到 Deferrable Operator 与 Asset 驱动调度的工程全解
大多数人对 Airflow 的理解停留在「写个 Python 文件定义 DAG,然后 @daily 就能跑」。但当集群里躺着 800 个 DAG、3000 个 Task、高峰期每分钟要创建上百个 DagRun 时,你会发现真正的瓶颈从来不是你的 Python 业务逻辑,而是那个在后台默默打转的 Scheduler 进程。本文拆到源码级别,讲清楚 Airflow 2.x 调度器的每一处设计取舍,以及为什么这些取舍决定了你的任务延迟上限。
一、先认清一个事实:Airflow 不是调度执行引擎,是个状态机协调器
Airflow 的核心抽象是元数据驱动的状态机。数据库里有一张 task_instance 表,每一行是一个任务的某个执行实例(dag_id + task_id + run_id)。所谓「调度」,本质就是 Scheduler 进程不断扫描这张表,把符合条件的行从 scheduled 推到 queued,再由 Executor 推到 running,最后由 Worker 回写终态。
理解这一点很关键:Task 状态转换的驱动力来自数据库轮询,而不是内存里的定时器。这带来两个直接后果:
- 数据库是整个系统的吞吐上限。Scheduler 的每一次「心跳」都是一批 SQL。
- 任何优化本质上都是在做两件事:减少每次循环的 SQL 数量,或者缩短持有锁的时间。
Airflow 1.x 到 2.x 的全部改造,基本围绕这两条主线。
二、Serialized DAG:把 Python 对象砍出调度决策链路
2.1 为什么必须序列化
Airflow 1.x 的 Scheduler 是一个 DagBag 全量加载模型:主循环每 30 秒 collect_dags() 一次,把 dags_folder 里所有 .py 文件重新 import 成 Python DAG 对象,再逐个判断是否该触发。
这里有个致命问题:你的 DAG 文件里的任何顶层代码,都在 Scheduler 进程里以调度频率执行。一段顶层的 pandas.read_parquet("s3://...") 会让整个 Scheduler 卡住几十秒。更糟的是 pickle 出来的 DagRun 对象要在 Scheduler、Executor、Worker 之间传递,版本升级时 pickle 不兼容直接导致全站调度中断。
Airflow 2.0 引入 Serialized DAG:DAG 结构的真相来源从 Python 对象变成 serialized_dag 表里的 JSON 列。
2.2 序列化链路
DAG File ── DagFileProcessorProcess(独立进程)── 序列化 ── serialized_dag 表(JSONB)
│
▼
Scheduler 只读 JSON,永不 import 你的 py 文件
序列化后的 JSON 结构大致如下(简化):
{
"dag_id": "etl_user_profile",
"tasks": [{
"task_id": "extract",
"downstream_task_ids": ["transform"],
"operator": "airflow.operators.python.PythonOperator",
"start_date": 1700000000.0,
"retries": 2,
"retry_delay": 300.0,
"pool": "default_pool",
"priority_weight": 1,
"trigger_rule": "all_success",
"weight_rule": "downstream",
"executor_config": {"KubernetesExecutor": {"request_cpu": "500m"}}
}],
"dag_hash": "a3f9c2e1...",
"timetable": {"__type": "airflow.timetables.interval.CronDataIntervalTimetable",
"__var": {"expression": "0 2 * * *", "timezone": "UTC"}},
"tags": ["etl", "core"]
}
几个工程要点:
dag_hash是增量更新开关。DAG Processor 会计算序列化后的内容哈希,无变化就不写库,避免无谓的写放大和file_parse抖动。timetable用自定义反序列化协议(__type+__var)实现多态,这是 Airflow 2.2 之后把「cron 表达式」和「数据间隔」解耦的产物。- 序列化是有损的。你的
python_callable不会被存成逻辑,而是 pickle 成 bytes 塞进template_fields或SerializedBaseOperator的 blob。这就是为什么 Airflow 官方反复强调 top-level code 要轻——那部分代码每次 parse 都会跑。
# ❌ 反例:这段 SQL 查询在每个 file_parse 周期都会执行
df = pd.read_sql("SELECT * FROM dim_users", engine) # 顶层!
with DAG("etl_user_profile", ...) as dag:
@task
def transform():
return df.pipe(clean) # 闭包捕获,还会被 pickle 进 DB
# ✅ 正例:全部下沉到 task 内部,只在 Worker 执行时才跑
with DAG("etl_user_profile", ...) as dag:
@task
def transform():
df = pd.read_sql("SELECT * FROM dim_users", engine)
return df.pipe(clean)
这条规则看着简单,但它是生产环境 Scheduler CPU 打满的头号原因。
三、Scheduler 主循环与 Critical Section
3.1 三段式循环
Scheduler 的每个 loop 分成三块:
while True:
self.adopt_or_reset_orphaned_tasks() # ① 清理僵尸
self._run_scheduler_loop() # ② 主调度
self._emit_pool_metrics() # ③ 指标上报
其中 ② 的核心是:
# airflow/jobs/scheduler_job_runner.py(示意逻辑)
for callback in callbacks: # DagRun / TaskInstance 成功失败回调
callback.on_loop()
self._start_queued_dagruns() # 把 queued 的 DagRun 置 running
self._schedule_dag_run(dag_run) # 核心:从 TaskInstance 里挑可执行项
→ self._critical_section_enqueue_task_instances(...)
self._send_dag_callbacks_to_processor(...)
3.2 Critical Section 到底在保护什么
critical_section 装饰器内部是一个基于数据库行锁的互斥机制:
@provide_session
def _critical_section_enqueue_task_instances(self, session):
# 1. 尝试获取 scheduler_job 行的排他行锁(SELECT ... FOR UPDATE NOWAIT)
locked = self.job.executor.... _try_to_lock_scheduler_slot(session)
if not locked:
return 0 # 别的 scheduler 拿到了,本轮直接跳过
...
它保护的资源是:「判断并发额度 → 扣减 Pool slot → 写入 queued 状态」这一组操作的原子性。
如果不加锁,两个 Scheduler 实例同时读到 pool 还剩 1 个 slot,各自扣掉 1,结果超发到 -1,集群直接过载。所以 Airflow 选择了「悲观锁 + 快速失败」:拿不到锁的那个实例本轮什么都不做,让给同伴。这是典型的 throughput-vs-latency 取舍——宁可牺牲单次循环的效率,也要保证正确性。
实战含义:Scheduler 横向扩容并不是线性收益。3 个 Scheduler 的吞吐通常只有单实例的 1.8~2.2 倍,因为它们在争抢同一把行锁。超过 4 个实例,收益基本停滞,反而会因为 NOWAIT 抢锁失败产生大量空转。实测建议:Scheduler 实例数 ≤ 3,优先级 WRITE DAG > 加机器。
3.3 TaskInstance 状态机与四个并发阀门
一个 Task 从定义到执行,要连过四道闸:
| 阀门 | 作用范围 | 元数据字段 | 表现 |
|---|---|---|---|
max_active_tasks | 单个 DAG Run | dag.max_active_tasks | 同一 DagRun 内并行上限 |
max_active_runs | 单个 DAG | dag.max_active_runs | 未跑完的 DagRun 上限 |
| Pool slots | 全局,跨 DAG | slot_pool 表 | 资源池,比如 database_pool: 32 |
| Executor slots | 物理 worker 容量 | executor 配置 | K8s/Celery worker 并发数 |
with DAG(
dag_id="feature_store_backfill",
max_active_runs=3, # 同时最多 3 个 DagRun
max_active_tasks=16, # 每个 DagRun 内最多 16 个 task 并行
catchup=True,
schedule="0 * * * *",
) as dag:
backfill = PythonOperator(
task_id="backfill",
pool="spark_pool", # 走专用的 spark_pool,避免打爆 default_pool
pool_slots=4, # 单个 task 占 4 个 slot
priority_weight=10,
weight_rule="absolute", # 不随下游膨胀,固定优先级
execution_timeout=timedelta(hours=2),
retries=2,
retry_delay=timedelta(minutes=5),
)
priority_weight 默认用 weight_rule="downstream",即下游所有 task 的权重之和。这意味着 DAG 末端的叶子节点权重最低——如果你的关键路径在下游,反而会被降级。涉及 SLA 的核心链路,显式设 weight_rule="absolute" 更安全。
四、Deferrable Operator:把等待从 Worker 手里抢出来
4.1 传统 Sensor 的资源陷阱
一个 ExternalTaskSensor 默认每 60 秒 poll 一次,等外部 DAG 完成。假设下游任务要等 8 小时:
Worker slot 占用:8 小时
实际有效工作时间:~15 秒(几十次 SQL 查询)
资源利用率:0.05%
这是 Airflow 里最典型的浪费。100 个这样的 Sensor 就能把 Celery worker 队列彻底堵死。
4.2 Triggerer + asyncio 的解法
DeferrableOperator 的执行模型是两段式:
deferrable=True ──► Task 执行到 self.defer(...) ──► 释放 Worker slot,状态转 deferred
│
Trigger 注册到 Triggerer 进程
(asyncio event loop,单进程可挂数千协程)
│ 事件就绪
TaskInstance 状态转 scheduled ──► 重新排队执行
自定义一个 Deferrable Operator:
from airflow.sensors.base import BaseSensorOperator
from airflow.triggers.base import BaseTrigger, TriggerEvent
from typing import Any
class MLJobDoneTrigger(BaseTrigger):
"""轮询某个 ML 训练服务,直到 job 完成。"""
def __init__(self, job_id: str, poll_interval: float = 30.0):
super().__init__()
self.job_id = job_id
self.poll_interval = poll_interval
def serialize(self) -> tuple[str, dict[str, Any]]:
# Trigger 必须可被 JSON 序列化才能在进程重启后恢复
return ("my_provider.triggers.ml.MLJobDoneTrigger",
{"job_id": self.job_id, "poll_interval": self.poll_interval})
async def run(self):
try:
while True:
status = await self._async_check(self.job_id) # 非阻塞 HTTP
if status in ("SUCCEEDED", "FAILED"):
yield TriggerEvent({"status": status, "job_id": self.job_id})
return
await asyncio.sleep(self.poll_interval)
except asyncio.CancelledError:
# Task 被 clear / DagRun 超时,这里做清理
await self._async_cancel(self.job_id)
raise
class WaitMLJobSensor(BaseSensorOperator):
def execute(self, context):
status = self._sync_check(context["params"]["job_id"])
if status not in ("SUCCEEDED", "FAILED"):
# 关键:defer 之后本 task 会释放 slot,由 Triggerer 接管
self.defer(
trigger=MLJobDoneTrigger(
job_id=context["params"]["job_id"], poll_interval=30.0),
method_name="execute_complete",
timeout=timedelta(hours=12),
)
return status
def execute_complete(self, context, event=None):
# 第二段:Triggerer 触发后的回调,此时才重新占 Worker slot
if event["status"] != "SUCCEEDED":
raise AirflowException(f"ML job failed: {event}")
self.log.info("ML job %s succeeded", event["job_id"])
注意 serialize() 必须返回可 JSON 化的 dict:Triggerer 进程可能在 Task 还 deferred 时被重启,Trigger 需要能从 DB 里的字段重建。把连接池、HTTP client 塞进 Trigger 会导致序列化失败——所有 IO 都必须在协程内部现开。
生产收益是立竿见影的:把集群里 60% 的 Sensor 换成 deferrable 版本后,同样的 Celery worker 规模,任务排队延迟通常能下降一个数量级。这是 Airflow 性价比最高的单项优化。
五、Asset:从时间驱动到数据驱动
Airflow 2.4 引入 Dataset,2.5 定名为 Asset,解决的是一个长期痛点:DAG A 产出某张表,DAG B 要消费它,但两者的 cron 周期天然对不齐。
传统做法是用 ExternalTaskSensor 跨 DAG 等待——又回到 worker slot 被占用的老路上。Asset 的做法是把依赖声明提升到数据契约层:
# 生产者:不关心谁消费
users_asset = Asset(
name="s3://lake/dim_users/dt=*",
group="warehouse",
extra={"owner": "data-platform", "freshness_sla_minutes": 60},
)
with DAG(dag_id="build_dim_users", schedule="@daily") as producer:
@task(outlets=[users_asset]) # outlets 声明产出
def build():
write_partition("dim_users")
# 消费者:不关心谁生产、什么时候生产
with DAG(
dag_id="build_user_features",
schedule=[users_asset], # ← 注意:schedule 接受 Asset 列表
) as consumer:
@task
def features(): ...
底层实现是:任何带 outlets 的 task 成功时,往 asset_event 表插一条记录;Scheduler 侧的 AssetTriggeredTimetable 检查每个 Asset 是否有新的 asset_event,若有则为所有订阅它的 DAG 创建 DagRun。
两个重要细节:
- 多 Asset 默认是 AND 还是 OR? 2.9 之前是 OR(任一 asset 更新就触发),之后引入了
AssetAny/AssetAll显式表达。混用会导致反直觉的触发频率,务必显式写清楚。 AssetAlias是给 CI/CD 用的逃生舱。生产/测试环境 Asset URI 不同(比如 bucket 前缀不同),用 Alias 可以让 DAG 代码跨环境复用,运行时再解析到真实 Asset。
Asset 的真正价值在于:它让「数据新鲜度」成为调度的一等公民,而不是用 cron + 传感器去猜。
六、生产调优清单(按收益排序)
| 问题 | 根因 | 动作 |
|---|---|---|
| Scheduler CPU 长期 100% | DAG 文件顶层重 IO | 顶层代码下沉;[scheduler] parsing_processes 与 CPU 核数对齐 |
| 任务触发延迟不稳定 | 文件解析抖动 | [scheduler] min_file_process_interval 调大;file_parsing_sort_mode 设为 modified_time |
| UI 页面加载慢 | serialized_dag 表膨胀 | 定期 airflow dags reserialize + 清理孤儿行 |
| 任务创建延迟高 | Critical Section 锁竞争 | Scheduler 实例 ≤ 3;拆分超大 DAG |
| Worker 队列常年满 | Sensor 占用 slot | 迁移到 Deferrable Operator |
| 数据库 IOPS 打满 | polling SQL 太频繁 | [scheduler] scheduler_idle_sleep_time 适当放大(默认 1s,可到 5s) |
| 任务重复执行 | Zombie 检测延迟 | 配好 Worker 心跳与 [scheduler] zombie_detection 间隔 |
最后一条建议:不要把 Airflow 当作流处理引擎用。catchup=True + @hourly 跑两年历史数据会产生 17000 个 DagRun,任何优化都救不了这种用法。批处理编排器和流式引擎的边界要划清楚——前者管依赖与重试,后者管毫秒级延迟。混淆这两者,是 Airflow 集群崩溃最常见的人为原因。
总结:Airflow 2.x 的调度器是一个围绕数据库构建的状态机协调器,其全部设计都可以归结为一句话——在保证多实例并发正确的前提下,尽可能减少每次循环对数据库的冲击。Serialized DAG 解决了 CPU 侧的开销,Critical Section 保证了扣减额度的原子性,Deferrable Operator 消灭了无意义的 slot 占用,Asset 则把调度从时间维度解放到数据维度。理解这四者的因果,你就掌握了 Airflow 调优的全部主干。

发表评论 取消回复