引言
Airflow 是数据领域的事实标准调度器。它的核心抽象非常朴素:一个 DAG 描述任务之间的依赖,调度器按时间周期触发 DagRun,Executor 把任务分发到 Worker 执行。这个模型解决了「每天凌晨要按顺序跑 200 个 SQL」这类需求,也让数据团队第一次有了统一的编排入口。
但 Airflow 的心智模型比看起来复杂。它最容易让人困惑的是时间语义:data_interval_start 到底指什么、为什么昨天写的 DAG 今天才跑、为什么回填会重跑出错。这些问题的根源在于 Airflow 的调度是「按数据区间驱动」而不是「按执行时刻驱动」,理解这一点之后大部分困惑会消失。
另一个常见误区是把 Airflow 当成通用工作流引擎。它没有「等人审批」的原语,也不适合做毫秒级的服务编排。它的强项是「周期性、批处理、可回填」的数据管道,判断标准是「这个任务是否有一个天然的数据区间」。本文按这个视角展开,先讲调度语义,再讲 DAG 的写法与工程实践,最后讲性能调优与测试。想先看整体选型框架的读者,可以从 工作流引擎全景与选型 开始。
目录
- Airflow 的核心概念与组件
- 调度语义:数据区间而非执行时刻
- 一个可运行的 DAG
- TaskFlow API 与参数传递
- catchup 与回填
- 幂等设计:让任务可以安全重跑
- 动态任务映射
- 传感器与延迟等待
- XCom 与数据传递的边界
- 连接、变量与密钥管理
- Executor 选型
- 池与并发控制
- Asset 与数据驱动调度
- 与数据仓库分层的协作
- 数据库迁移管道的编排
- 调度延迟与性能调优
- 可观测与告警
- 权衡取舍
- 常见坑清单
- 小结
1. Airflow 的核心概念与组件
Airflow 有六个核心概念,先对齐术语:
| 概念 | 含义 |
|---|---|
| DAG | 有向无环图,描述任务与依赖,是调度的单位 |
| Task | DAG 中的一个节点,由 Operator 实例化 |
| DagRun | DAG 的一次执行,对应一个数据区间 |
| 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 | 部署简单、生态成熟 | 任务共用依赖,隔离性差 |
| KubernetesExecutor | Pod 级隔离与资源配额 | 每任务启动延迟数秒 |
| 传感器轮询 | 实现简单 | 占槽位或增加调度开销 |
| deferrable 传感器 | 不占 Worker | 需要 Triggerer 组件 |
| XCom 传数据 | 使用方便 | 大数据拖垮元数据库 |
| 任务内做重计算 | DAG 简单 | Python 依赖膨胀,环境难维护 |
| Asset 驱动 | 跨团队解耦 | 触发粒度粗,调试链路长 |
19. 常见坑清单
- 任务里用
datetime.now()而不是data_interval_start,回填时全部处理当天数据。 - 回填时任务用
INSERT而非覆盖或 upsert,表里数据成倍增长。 catchup=True配上一年前的start_date,上线瞬间创建 365 个 DagRun 打爆数据库。- 把 DataFrame 或大 JSON 塞进 XCom,元数据库膨胀到几百 GB。
- 在 DAG 顶层作用域调
Variable.get,每次解析都查一次数据库,解析时间飙升。 - 传感器用默认
poke_interval=60且mode="poke",几十个传感器占满 Worker 槽位。 - 不设
max_active_runs,同一 DAG 的多个区间并发跑,写同一张表互相覆盖。 execution_timeout不设,任务挂死占用槽位直到运维发现。- 所有任务共用一个池,重查询把并发占满,轻任务全部排队。
- DAG 数量多但
parsing_processes保持默认 2,DAG 更新要几分钟才生效。 - 只在任务失败时告警,不监控「DAG 未按预期启动」,调度器故障静默数小时。
- 用
airflow tasks clear重跑但没考虑下游依赖,下游读到中间状态。
20. 小结
Airflow 的复杂度集中在时间语义上:它按数据区间调度而不是按执行时刻调度,理解这一点之后,回填、幂等、依赖等待这些设计都会变得自然。工程上的三条底线是「任务幂等」「数据用引用传递」「并发用池控制」。
Airflow 的定位应该保持克制:它是编排器,不是计算引擎。把重计算推给 Spark、Trino、数仓,Airflow 只负责「什么时候跑、什么顺序跑、失败怎么办」,这样 DAG 的依赖能保持极简,环境也容易维护。
如果团队更关注数据血缘与资产视图,可以对比 Dagster 与 Prefect 数据编排 ;如果流程里有人工审批与长等待,Airflow 不是合适的工具,应该看 BPMN 2.0 与 Camunda 实战 或 Temporal 与持久化执行 。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。