
一个让数据工程师用 Python 写工作流、可视化调度、分布式执行的开源平台。从 Airbnb 内部工具到 Apache 顶级项目,11 年迭代、3.x 大重构让 ” 写 DAG” 从 JSON 配置文件变回纯 Python 函数。
写在前面
每个数据团队都经历过这种混乱:
# 凌晨 2 点
# ETL-3 失败了,但 ETL-4 还在跑,因为没人通知上一个脚本退出了
# 运营拉取的数据跟 BI 看板对不上,因为脚本跑的时间点漂了
# 新人入职第一周:学长留下的 47 个 cron + 19 个 shell 脚本 + 1 份没人维护的 README
Apache Airflow 就是为了终结这种 ”cron + 胶水脚本 ” 混乱而生的——用 Python 写 DAG、可视化调度 、 依赖管理 、 失败重试 、 告警,让工作流从 ” 脚本堆 ” 升级为 ” 可观测、可维护、可协作 ” 的工程产物。
11 年过去,它依然是数据 /ML 领域 事实标准的工作流编排器——Databricks、AWS、Microsoft、Google、Netflix、Slack、Spotify 等几千家公司在生产环境跑它。但很多人对 Airflow 的印象还停留在 “1.x/2.x 时代需要写一堆 JSON 配置 ”,Airflow 3.x 已经把这一切改写了。
一、它解决什么问题
一句话卖点:用 Python 定义、可视化调度、分布式执行的现代工作流编排平台,ETL / 数据管道 / ML pipeline 的事实标准。
| 维度 | 信息 |
|---|---|
| 项目名 | Apache Airflow |
| 出身 | Airbnb(2014 年 Maxime Beauchemin 创建) |
| 项目阶段 | Apache 顶级项目(2019 年毕业) |
| Stars | 46,460 |
| Forks | 17,577 |
| License | Apache-2.0 |
| 主语言 | Python ~91% + TypeScript ~7% + Go + Kotlin |
| 创建时间 | 2015-04-13(airbnb/airflow → apache/airflow 2016 年捐赠) |
| 最新稳定版 | 3.3.1(2026 年,PyPI/Docker 均已发布) |
| 最近 Push | 2026-08-13(今天, 持续活跃) |
| Open Issues | 1,892 |
| Subscribers | 781 |
| 核心依赖 | Python ≥ 3.10,SQLAlchemy,Jinja2,Flask,celeryExecutor/k8sExecutor |
| 官网 | https://airflow.apache.org/ |
二、核心能力 1:DAG + Python + 可视化
2.1 用 Python 写 DAG(不是 JSON/YAML)
Airflow 的 第一性原理:DAG 是 Python 代码,而不是配置文件。
# airflow/dags/etl_demo.py
from airflow.sdk import dag, task
from datetime import datetime
@dag(start_date=datetime(2026, 1, 1), schedule="@daily", catchup=False)
def etl_demo():
@task
def extract():
return fetch_from_api() # 你的数据源
@task
def transform(raw):
return clean(raw)
@task
def load(cleaned):
write_to_warehouse(cleaned)
load(transform(extract()))
etl_demo()
对比 YAML/JSON 配置:参数化、循环生成、子 DAG、动态分支、条件判断——Python 原生表达力碾压 YAML 模板。
2.2 三大核心抽象
| 抽象 | 角色 | 例子 |
|---|---|---|
| DAG | 工作流定义(有向无环图) | 一个 ETL pipeline |
| Operator/Task | 工作流节点 | BashOperator、PythonOperator、@task 装饰器 |
| Sensor | 等待外部事件 | S3KeySensor、ExternalTaskSensor |
2.3 调度与依赖
- Cron 表达式:
schedule="0 2 * * *"每天凌晨 2 点 - Dataset 触发:上游 task 产出 dataset → 自动触发下游 DAG(3.x 新增)
- 任务依赖:
task_a >> task_b >> task_c链式表达 - Branch / Shortcircuit:条件分支跳过任务
2.4 Web UI
┌─ Grid View ────────┬─ Graph View ─────┬─ Calendar ─────┐
│ 任务 N 状态色块 │ DAG 拓扑图 │ 调度日历 │
│ 时间轴 / 重试次数 │ 节点 / 边 / 状态 │ 跑批历史 │
│ 失败 task 红黄高亮 │ 鼠标悬停看日志 │ 季度对比 │
└────────────────────┴──────────────────┴───────────────┘
价值:新人不需要 SSH 上服务器 tail 日志就能看到所有 DAG 状态。
三、核心能力 2:Airflow 3.x 大重构
3.x 是 Airflow 自 1.x 以来 最大的一次架构革新, 主要变化:
3.1 TaskFlow API(2.0+ 引入,3.x 完善)
老写法(Operator 显式):
t1 = PythonOperator(task_id="extract", python_callable=extract_fn)
t2 = PythonOperator(task_id="transform", python_callable=transform_fn)
t1 >> t2
新写法(@task 装饰器, 自动 XCom):
@task
def extract(): return data
@task
def transform(data): return clean(data)
result = transform(extract()) # 自动 XCom 传递
好处:代码量减少 60%,XCom 自动化,IDE 自动补全。
3.2 Edge Executor / Edge Worker(3.x 新)
老架构:Scheduler / Worker 都要跑在中心集群,K8s 上部署繁琐。
3.x 新增:Worker 可以跑在边缘(3rd-party infra / 远程 GPU / 边缘节点), 通过 API 与中心通信。
适用场景:ML 团队想在 GPU 集群跑训练, 但调度仍由中心 Airflow 管。
3.3 DAG Bundles(3.x 新)
DAG 文件从 git/s3/GCS 加载, 支持 多 bundle——同一 Airflow 实例跑多个团队 / 多个 repo 的 DAG。
好处: 不必所有 DAG 都放一个 monorepo。
3.4 Task SDK(3.x 新)
老 Airflow 把 Scheduler / Executor / Webserver 都强耦合在 airflow 包里。3.x 拆出独立 apache-airflow-task-sdk,Worker 可以脱离 Airflow 主包运行。
3.5 Pandas 3 兼容(3.3.1 重要变更)
XCom DataFrame 序列化兼容 pandas 2 / pandas 3, 这是 3.3.1 的重要变更(2026-07 刚发)。
四、扩展生态:Providers + Executors + Hooks
4.1 Providers(80+ 官方)
每个第三方服务一个 provider 包:
| Provider | 内容 |
|---|---|
apache-airflow-providers-amazon |
AWS(S3 / EMR / Glue / SageMaker / Lambda) |
apache-airflow-providers-google |
GCP(BigQuery / GCS / Dataproc / Cloud Composer) |
apache-airflow-providers-cncf-kubernetes |
K8s Executor / K8sPodOperator |
apache-airflow-providers-microsoft-azure |
Azure 全家桶 |
apache-airflow-providers-apache-spark |
SparkSubmitOperator |
apache-airflow-providers-dbt-cloud |
dbt Cloud 集成 |
apache-airflow-providers-openai / -anthropic |
LLM 调用 |
apache-airflow-providers-slack / -discord |
告警通知 |
生态规模:80+ 官方 providers,400+ community hooks——任何常见集成都有现成 Operator。
4.2 Executors(5 种)
| Executor | 适用 |
|---|---|
SequentialExecutor |
开发测试 |
LocalExecutor |
中小规模单机 |
CeleryExecutor |
传统分布式(Redis/RabbitMQ 队列) |
KubernetesExecutor |
生产标准: 每个 task 一个 Pod |
CeleryKubernetesExecutor |
混合:critical task 走 K8s, 批量走 Celery |
4.3 Hooks(连接抽象)
Hook 是 Airflow 跟外部系统通信的统一抽象——S3Hook、BigQueryHook、SnowflakeHook 等。Task 通过 Hook 拿到连接, 不用关心认证 / 重试细节。
五、对比 Dagster / Prefect / Argo Workflows
| 维度 | Apache Airflow | Dagster | Prefect | Argo Workflows |
|---|---|---|---|---|
| Stars | 46K | 10K | 14K | 14K(K8s 原生) |
| License | Apache-2.0 | Apache-2.0 | Apache-2.0 | Apache-2.0 |
| 诞生 | 2015 Airbnb | 2018 | 2018 | 2017 Intuit |
| DAG 定义 | Python @dag | Python @job + ops/assets | Python @flow | YAML/CRD |
| 资产意识 | ❌ 1.x/2.x 弱,3.x 加 Dataset | ✅ 核心抽象 | ⚠️ 较新 | ❌ |
| 类型系统 | ❌ | ✅ type-safe | ⚠️ | ❌ |
| 本地开发体验 | ⭐⭐⭐ | ⭐⭐⭐⭐⭐ | ⭐⭐⭐⭐⭐ | ⭐⭐ |
| K8s 集成 | ⭐⭐⭐⭐⭐ | ⭐⭐⭐⭐ | ⭐⭐⭐⭐ | ⭐⭐⭐⭐⭐(原生) |
| 生态成熟度 | ⭐⭐⭐⭐⭐ | ⭐⭐⭐ | ⭐⭐⭐ | ⭐⭐⭐ |
| 学习曲线 | 中(概念多) | 中(资产模型新) | 低(更 Pythonic) | 中(K8s 知识) |
| 生产案例 | Airbnb/Netflix/Slack/Spotify | 多家中型公司 | 部分 SaaS | K8s 重用户 |
结论:
- 选 Airflow: 数据团队规模 ≥ 5 人、需要成熟生态、生产 K8s 部署
- 选 Dagster: 数据资产优先(配合 dbt)、想要 type-safe、愿意吃新概念
- 选 Prefect: 小团队、Pythonic 优先、本地开发体验优先
- 选 Argo:K8s 原生、Workflow 就是 K8s CRD、不需要 Python
六、实战:3 步跑通
6.1 Docker Compose 启动
# 1. 拉镜像
docker pull apache/airflow:3.3.1
# 2. 起服务(scheduler + webserver + postgres + redis)
docker run -d --name airflow-web -p 8080:8080 \
-e AIRFLOW__DATABASE__SQL_ALCHEMY_CONN=postgresql+psycopg2://airflow:airflow@postgres/airflow \
apache/airflow:3.3.1 webserver
docker run -d --name airflow-sched \
--link airflow-web \
apache/airflow:3.3.1 scheduler
# 3. 访问 http://localhost:8080(默认账号 airflow/airflow)
6.2 第一条 DAG(TaskFlow API)
把文件丢到 $AIRFLOW_HOME/dags/:
# dags/hello_world.py
from airflow.sdk import dag, task
from datetime import datetime
@dag(start_date=datetime(2026, 1, 1), schedule="@hourly", catchup=False)
def hello_world():
@task
def greet(name: str) -> str:
return f"Hello, {name}!"
@task
def shout(message: str) -> None:
print(message.upper())
shout(greet("Airflow 3.x"))
hello_world()
打开 Web UI → DAGs 列表 → 触发 → 看日志输出 HELLO, AIRFLOW 3.X!。
6.3 接 S3 + K8s Executor(生产配置)
# airflow.cfg
executor = KubernetesExecutor
kubernetes_namespace = airflow
# dags/etl_s3_k8s.py
from airflow.providers.amazon.aws.transfers.s3_to_redshift import S3ToRedshiftOperator
S3ToRedshiftOperator(
task_id="s3_to_redshift",
s3_bucket="my-bucket",
s3_key="data/{{ds}}/",
redshift_conn_id="redshift_default",
table="events",
)
部署 K8sExecutor 后, 每个 task 自动起一个 Pod, 任务结束 Pod 自动销毁。
七、风险与坑
| 风险 | 说明 | 缓解 |
|---|---|---|
| 概念多, 上手陡 | DAG/Operator/Sensor/Hook/Executor/Variables/Connections/Plugins 一堆 | TaskFlow API 大幅简化; 新项目直接上 3.x |
| 元数据库压力 | 所有 task 状态 / 日志 /XCom 全存 metadata DB(默认 SQLite) | 生产必须 PostgreSQL; 定期清理旧 task instance |
| 3.x 升级坑 | AIP-44/Dag Bundle 等破坏性变更, 需要重写部分 DAG | 跟 release notes; 考虑先在测试环境跑 1 个月 |
| 老 Operator 不维护 | 部分 community provider 几年没更新 | 优先用官方 providers |
| DAG 文件 import 时间长 | Top-level import 会拖慢 scheduler 解析 | 用 Dynamic DAG 或 TaskFlow 延迟加载 |
| Zombie/Orphan task | worker 异常退出留下僵尸 task | 启用 scheduler_health_check_threshold + min_file_process_interval |
八、总结
Airflow 最适合的三个场景:
- 数据团队的 ETL/ELT 编排 —— 5+ 个 DAG、跨 S3/BigQuery/Snowflake/Redshift, 需要统一调度
- ML pipeline 编排 —— 训练 / 评估 / 部署 / 监控全流程,3.x 的 Edge Executor 让 GPU 集群灵活接入
- 企业内部工作流平台 —— 跨团队共享 Operators/Hooks,Airflow 成为公司数据基础设施
一句话建议:
如果你团队 ≥ 3 个数据工程师、跨 ≥ 2 个数据源、需要 ≥ 5 个 DAG,直接上 Airflow 3.x——它仍是数据编排领域最稳的选择; 如果你只跑 1-2 个 cron 脚本, 先用 crontab + Python 脚本, 别过度工程。
先试一周:docker run apache/airflow:3.3.1 → 写一个 @dag + 3 个 @task → 接一个真实数据源。这周你会决定 Airflow 是不是你团队的 ” 标准 ETL 工具 ”。
参考
- apache/airflow — GitHub 仓库(46K stars)
- Apache Airflow 官网 — 官方主页
- Airflow 3.3.1 Release Notes — 最新版本
- Airflow 文档 — 完整文档
- Providers 列表 — 80+ 官方 providers
- Astronomer 文档 — 商业公司贡献的最佳实践
- TaskFlow API 教程 — 3.x 推荐写法
- Awesome Airflow — 生态资源汇总
研究文档(引用来源参考)
(no reference document available)