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 状态转换的驱动力来自数据库轮询,而不是内存里的定时器。这带来两个直接后果:

  1. 数据库是整个系统的吞吐上限。Scheduler 的每一次「心跳」都是一批 SQL。
  2. 任何优化本质上都是在做两件事:减少每次循环的 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 Rundag.max_active_tasks同一 DagRun 内并行上限
max_active_runs单个 DAGdag.max_active_runs未跑完的 DagRun 上限
Pool slots全局,跨 DAGslot_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。

两个重要细节:

  1. 多 Asset 默认是 AND 还是 OR? 2.9 之前是 OR(任一 asset 更新就触发),之后引入了 AssetAny / AssetAll 显式表达。混用会导致反直觉的触发频率,务必显式写清楚。
  2. 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 调优的全部主干。

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部