Airflow DAG 调度体系

本文系统讲解 Airflow 的调度体系与工程实践,回答 DAG 该怎么写、调度周期与数据区间如何对齐、回填为什么容易重跑出错、Executor 怎么选。覆盖 TaskFlow API、数据区间语义与 catchup、任务幂等设计、动态任务映射、传感器与 deferrable 模式、XCom 传递边界、池与并发控制、Asset 数据驱动调度、迁移管道编排与性能调优。

引言

Airflow 是数据领域的事实标准调度器。它的核心抽象非常朴素:一个 DAG 描述任务之间的依赖,调度器按时间周期触发 DagRun,Executor 把任务分发到 Worker 执行。这个模型解决了「每天凌晨要按顺序跑 200 个 SQL」这类需求,也让数据团队第一次有了统一的编排入口。

但 Airflow 的心智模型比看起来复杂。它最容易让人困惑的是时间语义:data_interval_start 到底指什么、为什么昨天写的 DAG 今天才跑、为什么回填会重跑出错。这些问题的根源在于 Airflow 的调度是「按数据区间驱动」而不是「按执行时刻驱动」,理解这一点之后大部分困惑会消失。

另一个常见误区是把 Airflow 当成通用工作流引擎。它没有「等人审批」的原语,也不适合做毫秒级的服务编排。它的强项是「周期性、批处理、可回填」的数据管道,判断标准是「这个任务是否有一个天然的数据区间」。本文按这个视角展开,先讲调度语义,再讲 DAG 的写法与工程实践,最后讲性能调优与测试。想先看整体选型框架的读者,可以从 工作流引擎全景与选型 开始。

目录

  1. Airflow 的核心概念与组件
  2. 调度语义:数据区间而非执行时刻
  3. 一个可运行的 DAG
  4. TaskFlow API 与参数传递
  5. catchup 与回填
  6. 幂等设计:让任务可以安全重跑
  7. 动态任务映射
  8. 传感器与延迟等待
  9. XCom 与数据传递的边界
  10. 连接、变量与密钥管理
  11. Executor 选型
  12. 池与并发控制
  13. Asset 与数据驱动调度
  14. 与数据仓库分层的协作
  15. 数据库迁移管道的编排
  16. 调度延迟与性能调优
  17. 可观测与告警
  18. 权衡取舍
  19. 常见坑清单
  20. 小结

1. Airflow 的核心概念与组件

Airflow 有六个核心概念,先对齐术语:

概念含义
DAG有向无环图,描述任务与依赖,是调度的单位
TaskDAG 中的一个节点,由 Operator 实例化
DagRunDAG 的一次执行,对应一个数据区间
TaskInstance任务的一次执行,有状态(queued / running / success / failed)
Operator任务类型的实现(Python、Bash、SQL、K8s Pod 等)
Executor任务分发机制(Local、Celery、Kubernetes)

运行时的四个组件:Scheduler(解析 DAG、创建 DagRun、分发任务)、Executor(决定任务在哪跑)、Worker(实际执行)、Metadata Database(存 DAG 结构、运行状态、变量、连接)。Airflow 3 增加了 API Server 与 DAG Processor 的独立部署形态,把「解析 DAG」从调度器里拆了出去。

理解「元数据库是唯一真相」很重要:任务的每次状态变化都写库,所以 Airflow 的吞吐上限往往由数据库写入决定,而不是计算资源。

2. 调度语义:数据区间而非执行时刻

这是 Airflow 最需要先理解的一点。当你写 schedule="@daily" 时,Airflow 创建的 DagRun 有一个 data_interval(数据区间),比如 2026-10-06 00:00 到 2026-10-07 00:00。这个 DagRun 会在区间结束后才被触发,也就是 2026-10-07 00:00 之后。

@dag(schedule="@daily", start_date=datetime(2026, 10, 1))
def daily_etl():
    ...
# 产生的 DagRun:
#   data_interval = [10-01, 10-02)  在 10-02 触发
#   data_interval = [10-02, 10-03)  在 10-03 触发

这个设计的理由是「处理完整的一天数据」:要处理 10 月 6 日的数据,必须等 10 月 6 日结束。所以任务里读数据的条件应该用 data_interval_start 与 data_interval_end,而不是 datetime.now()。

在任务里访问这个区间:

@task
def extract(data_interval_start=None, data_interval_end=None):
    # 用区间做查询条件,而不是用当前时间
    run_query(f"SELECT * FROM orders WHERE created_at >= '{data_interval_start}' "
              f"AND created_at < '{data_interval_end}'")

用 datetime.now() 的 DAG 在回填时会全部处理「今天」的数据,这是回填出错的头号原因。

3. 一个可运行的 DAG

from datetime import datetime, timedelta
from airflow.sdk import dag, task

@dag(
    dag_id="orders_daily_pipeline",
    schedule="0 2 * * *",
    start_date=datetime(2026, 10, 1),
    catchup=False,
    max_active_runs=1,
    default_args={
        "retries": 3,
        "retry_delay": timedelta(minutes=5),
        "retry_exponential_backoff": True,
        "execution_timeout": timedelta(hours=1),
    },
    tags=["orders", "daily"],
)
def orders_daily_pipeline():

    @task
    def extract(data_interval_start=None, data_interval_end=None):
        return run_query(
            "SELECT * FROM ods.orders "
            "WHERE created_at >= %s AND created_at < %s",
            (data_interval_start, data_interval_end),
        )

    @task
    def transform(rows):
        return [normalize(r) for r in rows]

    @task
    def load(rows):
        upsert_into("dws.orders_daily", rows)

    load(transform(extract()))

orders_daily_pipeline()

关键参数:catchup=False 避免一上线就补跑所有历史区间(start_date 是一年前的话会瞬间创建 365 个 DagRun);max_active_runs=1 保证同一 DAG 不并发跑多个区间,避免写同一张表冲突;execution_timeout 防止任务卡死占用槽位。

4. TaskFlow API 与参数传递

TaskFlow API(Airflow 2.0 引入)用装饰器把 Python 函数变成任务,返回值自动通过 XCom 传递。它的价值是让 DAG 的依赖关系由「函数调用」自然表达,而不是手工写 >>。

@task(retries=5, retry_delay=timedelta(seconds=30))
def call_api(endpoint: str) -> dict:
    return requests.get(endpoint, timeout=30).json()

@task
def summarize(payload: dict) -> str:
    return f"{len(payload['items'])} items"

summarize(call_api("https://internal/api/orders"))

要注意装饰器参数(retries、pool、trigger_rule)与函数参数的区别:前者是 Airflow 的 Task 属性,后者是 XCom 传递的数据。名字冲突时 Airflow 会把函数参数当成「任务输入」,容易踩坑,建议业务参数用明确的前缀。

trigger_rule 控制任务的触发条件,最常用的三个:all_success(默认,全部上游成功)、all_done(上游全部结束,不管成功失败,适合做清理)、one_failed(有上游失败就执行,适合做告警)。

5. catchup 与回填

catchup=True 时,Airflow 会为 start_date 到当前时间的每个区间都创建一个 DagRun,这是「回填历史数据」的机制。它有两种用法:一种是首次上线时自动补齐历史,另一种是用 airflow dags backfill 手工触发。

# 回填指定区间(Airflow 2.x 语法)
airflow dags backfill orders_daily_pipeline \
  --start-date 2026-09-01 --end-date 2026-09-30 \
  --reset-dagruns --yes

# 清空某天的任务状态,让它重跑
airflow tasks clear orders_daily_pipeline \
  --start-date 2026-10-05 --end-date 2026-10-05 --yes

回填能安全运行的前提是任务幂等。如果任务用 INSERT 而不是 INSERT OVERWRITE / MERGE,回填会重复插入数据。这是回填最常见的事故:回填一个月,表里数据翻倍。

max_active_runs 在回填时尤其重要:不限制的话,回填 30 天会同时起 30 个 DagRun,把数据库和下游打爆。建议回填时用 --max-active-runs 控制并发。

6. 幂等设计:让任务可以安全重跑

Airflow 的任务会被重跑:重试、手工 clear、回填、调度器故障恢复。所以每个任务都必须是幂等的。四种常用模式:

-- 模式一:按区间覆盖(推荐用于数仓)
DELETE FROM dws.orders_daily WHERE biz_date = '{{ ds }}';
INSERT INTO dws.orders_daily SELECT ... WHERE created_at::date = '{{ ds }}';

-- 模式二:分区覆盖(推荐用于 Hive/Spark)
INSERT OVERWRITE TABLE dws.orders_daily PARTITION (dt='{{ ds }}') SELECT ...;

-- 模式三:按主键 upsert
MERGE INTO dws.orders_daily t USING staging s ON t.id = s.id
WHEN MATCHED THEN UPDATE SET ... WHEN NOT MATCHED THEN INSERT ...;
# 模式四:用执行标识去重
@task
def send_notification(run_id=None):
    if already_sent(run_id):
        return
    send(...)
    mark_sent(run_id)

{{ ds }} 是 Airflow 模板变量,取数据区间起始日期的 YYYY-MM-DD。模板变量只在支持模板化的字段里生效(bash_command、sql、op_kwargs 等),Python 函数体里要用参数注入。

最忌讳的模式是「先查有没有、没有则插入」而不加唯一约束:并发下两个实例会同时通过检查。正确做法是让数据库约束兜底(唯一索引 + ON CONFLICT DO NOTHING)。

7. 动态任务映射

Airflow 2.3 引入动态任务映射(Dynamic Task Mapping),解决了「任务数量在运行时才知道」的问题。传统做法是写一个 for 循环处理列表,但那样失败后无法只重跑失败的那一项。

@task
def list_files() -> list[str]:
    return ["part-001.csv", "part-002.csv", "part-003.csv"]

@task
def process_file(path: str):
    upload(transform(read(path)))

process_file.expand(path=list_files())

.expand() 会为列表里的每个元素创建一个任务实例,在 UI 上可以看到每个分片的独立状态。.partial() 用于固定部分参数:

process_file.partial(bucket="s3://data-lake").expand(path=list_files())

动态映射的限制是「展开数量在任务创建时就确定」,不能在执行过程中动态增加。如果分片数不确定(比如边处理边发现新文件),需要用 .expand_kwargs() 配合上游任务返回完整列表,或者改用循环任务。

8. 传感器与延迟等待

传感器(Sensor)是「等待某个条件成立」的任务。传统传感器会占用一个 Worker 槽位并轮询,代价高;Airflow 2.2 引入的可延迟传感器(Deferrable Sensor)把等待交给 Triggerer 进程,不占 Worker 槽位。

from airflow.providers.amazon.aws.sensors.s3 import S3KeySensor
from airflow.sensors.base import deferrable

# 传统写法:poke_interval 决定轮询频率,占槽位
wait_file = S3KeySensor(
    task_id="wait_input",
    bucket_name="data-lake",
    bucket_key="raw/{{ ds }}/orders.csv",
    poke_interval=300,
    timeout=60 * 60 * 6,
    mode="reschedule",   # 不占槽位,但调度开销大
)
# 可延迟写法:等待期间不占 Worker,由 Triggerer 托管
wait_file = S3KeySensor(
    task_id="wait_input",
    bucket_name="data-lake",
    bucket_key="raw/{{ ds }}/orders.csv",
    deferrable=True,
    timeout=60 * 60 * 6,
)

mode="reschedule" 与 deferrable=True 都能释放槽位,区别是前者仍需要调度器周期性唤醒(每次唤醒写一次数据库),后者由 Triggerer 在内存里等待。文件到达类等待优先用 deferrable,数据库轮询类等待优先用 reschedule。

9. XCom 与数据传递的边界

XCom 是任务之间传递数据的机制,底层存在元数据库的 xcom 表里。它的设计用途是传递小数据(路径、配置、计数),不是传数据集。

@task
def extract() -> str:
    # 返回路径而不是数据本身
    write_parquet("s3://staging/orders/{{ ds }}.parquet")
    return "s3://staging/orders/{{ ds }}.parquet"

@task
def load(path: str):
    read_parquet(path)

把 100 MB 的 DataFrame 塞进 XCom 会把元数据库拖垮,这是 Airflow 最经典的性能事故。Airflow 2.6 之后有 XCom 后端可配置(存 S3/GCS),但仍然建议只传引用。

另一个细节是 XCom 的大小限制与序列化:默认用 JSON(enable_xcom_pickling=False),复杂对象要自己序列化。返回 Pandas DataFrame 时会因为 JSON 序列化失败而报错,这是新手常踩的坑。

10. 连接、变量与密钥管理

连接(Connection)存外部系统的地址与凭证,变量(Variable)存配置值。两者都存在元数据库,通过 UI 或 CLI 管理:

airflow connections add 'warehouse' \
  --conn-type 'postgres' \
  --conn-host 'db.internal' \
  --conn-login 'etl' \
  --conn-password "$(cat /run/secrets/db_pass)"

airflow variables set etl_batch_size 5000

元数据库里存明文凭证是安全风险,生产环境应该用 Secret Backend(AWS Secrets Manager、Vault、GCP Secret Manager),Airflow 会按需从后端拉取。配置方式:

[secrets]
backend = airflow.providers.amazon.aws.secrets.secrets_manager.SecretsManagerBackend
backend_kwargs = {"connections_prefix": "airflow/connections", "variables_prefix": "airflow/variables"}

DAG 里通过 BaseHook.get_connection("warehouse") 或 Variable.get("etl_batch_size") 访问。注意 Variable.get 在 DAG 顶层作用域调用会导致每次解析都查数据库,应该用 Variable.get(..., deserialize_json=True) 配合模板,或者放到任务函数里。

11. Executor 选型

Executor适用场景隔离性扩展方式
LocalExecutor单机、开发、小规模无(同进程)垂直扩容
CeleryExecutor中等规模、任务以 Python 为主进程级加 Worker 节点
KubernetesExecutor任务资源需求差异大Pod 级每任务一个 Pod
CeleryKubernetes混合:常规任务 Celery,重任务 K8s混合混合

选择标准是「任务之间是否需要不同的依赖或资源」。如果所有任务共用同一套 Python 依赖,Celery 最简单;如果任务需要不同镜像、不同 CPU/内存配额,K8s Executor 更合适,代价是每个任务启动一个 Pod(几秒的启动延迟)。

KubernetesExecutor 的一个实战细节:Pod 的启动延迟会显著影响短任务的总耗时。一个跑 5 秒的任务,加上拉镜像与调度可能要 30 秒。解决办法是用 pod_override 配置镜像拉取策略为 IfNotPresent,或者把短任务合并成一个任务。

12. 池与并发控制

池(Pool)是限制「同时运行的任务数」的机制,按资源维度划分:

airflow pools set warehouse_pool 5 "限制同时访问数仓的任务数"
airflow pools set spark_pool 20 "Spark 集群并发上限"
@task(pool="warehouse_pool", pool_slots=2)
def heavy_query():
    ...

pool_slots 让一个任务占用多个槽位,用于表达「这个任务消耗 2 份资源」。池的价值是保护下游系统:数仓连接数有限时,用池把并发压住比调大重试次数有效得多。

除了池,还有几层并发控制:max_active_tasks(DAG 级)、max_active_runs(DAG 级区间并发)、parallelism(整个 Airflow 的并行上限)、dag_concurrency。排查「任务排队不动」时按这个顺序逐层检查。

13. Asset 与数据驱动调度

Airflow 2.4 引入 Dataset(3.0 改名为 Asset),让 DAG 可以由「数据更新」而不是「时间」触发:

from airflow.sdk import asset

@asset(schedule="@daily")
def raw_orders():
    ...

@asset
def dwd_orders(raw_orders):   # 声明依赖,raw_orders 更新后自动触发
    ...

Asset 的价值是解耦:生产方不需要知道有哪些下游,只需声明「我产出了 asset X」;消费方声明「我依赖 X」,调度关系自动建立。这比在同一个 DAG 里画依赖更松耦合,也让跨团队的管道可以拼接。

限制是 Asset 的触发是「整个 DAG 级别」的,粒度不如任务级依赖;且 Asset 本身不携带数据,只是信号。用 Asset 时要配合数据区间语义,避免「上游更新一次、下游跑十次」的重复触发。

14. 与数据仓库分层的协作

Airflow 与数仓分层的标准对应关系是「每层一个 DAG 或一组任务」:

数仓层Airflow 中的体现调度频率
ODS抽取任务(从源库/日志抽取)每小时或每天
DWD清洗与明细加工每天,ODS 完成后
DWS轻度汇总每天,DWD 完成后
ADS应用层宽表与指标每天,DWS 完成后

分层的架构与建模细节见 数据仓库与湖仓架构 。Airflow 侧要注意的是「层间依赖用 Asset 还是用 DAG 内依赖」:同一团队用 DAG 内依赖更直观,跨团队用 Asset 更松耦合。

一个实战建议是把「调度时间」与「数据就绪」分开:不要假设 ODS 在 2 点一定跑完,而是让 DWD 用 Asset 或传感器等待 ODS 产出。用时间硬编码会在上游延迟时产生错误结果(读到不完整数据)而不是失败。

15. 数据库迁移管道的编排

数据库 schema 变更(DDL)是数据管道里最危险的部分,因为它不可回滚。Airflow 编排迁移的标准做法是「先兼容、再切换、后清理」的三阶段:

@dag(schedule=None, tags=["migration"])
def schema_migration_v14():
    @task
    def pre_check():
        assert_no_long_running_tx()
        assert_replication_lag_below(seconds=5)

    @task
    def add_column():
        execute_ddl("ALTER TABLE orders ADD COLUMN channel VARCHAR(32)")

    @task
    def backfill():
        execute_sql("UPDATE orders SET channel = 'unknown' WHERE channel IS NULL")

    @task
    def verify():
        assert_null_ratio("orders", "channel", below=0.001)

    pre_check() >> add_column() >> backfill() >> verify()

注意 schedule=None 表示这个 DAG 只手工触发,不参与周期调度——DDL 迁移必须是显式的、有人审批的动作。更完整的迁移策略(expand-contract 模式、双写切换、回滚预案)见 数据库迁移管道 。

16. 调度延迟与性能调优

Airflow 的默认调度间隔是 30 秒(scheduler_heartbeat_sec),加上 DAG 解析时间与任务队列延迟,端到端延迟通常在 1 到 3 分钟。要缩短延迟:

  • 减少 DAG 文件数量与顶层代码执行时间(顶层代码每次解析都会跑)。
  • 把 Variable.get 从顶层移到任务里。
  • 提高 parsing_processes(Airflow 2.x 默认 2,可调到 CPU 核数)。
  • 用 min_file_process_interval 控制重新解析频率(默认 30 秒,DAG 多时调大)。
  • 3.0 用 DAG Processor 独立进程,解析不再影响调度循环。
[scheduler]
parsing_processes = 8
min_file_process_interval = 60
scheduler_heartbeat_sec = 10

数据库层面要关注 task_instance 与 dag_run 表的增长,历史运行记录要定期清理(airflow db clean)。一个每天 1000 个任务实例的集群,一年后 task_instance 表会有 3.6 亿行,不清理会显著拖慢 UI 与调度。

17. 可观测与告警

Airflow 自带的 UI 能看 DAG 结构、运行历史、任务日志。生产环境还需要三类外部观测:

  • 指标:airflow_scheduler_heartbeat、airflow_dag_processing_total_parse_time、executor_queued_tasks、scheduler.tasks.running。
  • 告警:DAG 失败、任务超时、任务排队超过阈值、DAG 未按预期时间启动(missed SLA)。
  • 日志聚合:任务日志写到远端(S3/ES),UI 通过 remote_logging 读取。
default_args = {
    "on_failure_callback": notify_oncall,   # 任务失败告警
    "sla": timedelta(hours=3),              # SLA 未达成告警
    "email_on_failure": False,              # 用回调而不是邮件
}

「DAG 未按预期启动」是最容易被忽略的告警类型。如果调度器挂了 4 小时,任务不会失败,只是没跑,而下游可能已经用旧数据出了报表。监控手段是「检查预期存在的 DagRun 是否真的存在」,比如每天早上检查昨天的 DagRun 状态。相关实践见 监控与告警设计 。

18. 权衡取舍

选择收益代价
catchup=True 自动补齐历史上线即补齐首次上线可能瞬间创建大量 DagRun
catchup=False上线平稳历史数据要手工回填
CeleryExecutor部署简单、生态成熟任务共用依赖,隔离性差
KubernetesExecutorPod 级隔离与资源配额每任务启动延迟数秒
传感器轮询实现简单占槽位或增加调度开销
deferrable 传感器不占 Worker需要 Triggerer 组件
XCom 传数据使用方便大数据拖垮元数据库
任务内做重计算DAG 简单Python 依赖膨胀,环境难维护
Asset 驱动跨团队解耦触发粒度粗,调试链路长

19. 常见坑清单

  1. 任务里用 datetime.now() 而不是 data_interval_start,回填时全部处理当天数据。
  2. 回填时任务用 INSERT 而非覆盖或 upsert,表里数据成倍增长。
  3. catchup=True 配上一年前的 start_date,上线瞬间创建 365 个 DagRun 打爆数据库。
  4. 把 DataFrame 或大 JSON 塞进 XCom,元数据库膨胀到几百 GB。
  5. 在 DAG 顶层作用域调 Variable.get,每次解析都查一次数据库,解析时间飙升。
  6. 传感器用默认 poke_interval=60 且 mode="poke",几十个传感器占满 Worker 槽位。
  7. 不设 max_active_runs,同一 DAG 的多个区间并发跑,写同一张表互相覆盖。
  8. execution_timeout 不设,任务挂死占用槽位直到运维发现。
  9. 所有任务共用一个池,重查询把并发占满,轻任务全部排队。
  10. DAG 数量多但 parsing_processes 保持默认 2,DAG 更新要几分钟才生效。
  11. 只在任务失败时告警,不监控「DAG 未按预期启动」,调度器故障静默数小时。
  12. 用 airflow tasks clear 重跑但没考虑下游依赖,下游读到中间状态。

20. 小结

Airflow 的复杂度集中在时间语义上:它按数据区间调度而不是按执行时刻调度,理解这一点之后,回填、幂等、依赖等待这些设计都会变得自然。工程上的三条底线是「任务幂等」「数据用引用传递」「并发用池控制」。

Airflow 的定位应该保持克制:它是编排器,不是计算引擎。把重计算推给 Spark、Trino、数仓,Airflow 只负责「什么时候跑、什么顺序跑、失败怎么办」,这样 DAG 的依赖能保持极简,环境也容易维护。

如果团队更关注数据血缘与资产视图,可以对比 Dagster 与 Prefect 数据编排 ;如果流程里有人工审批与长等待,Airflow 不是合适的工具,应该看 BPMN 2.0 与 Camunda 实战 或 Temporal 与持久化执行 。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「工作流引擎」更多文章

  1. 工作流成本优化
  2. 执行器与资源隔离
  3. 调度、回填与补数