Apache Airflow 调研:46K Stars 的 Python 工作流编排事实标准,3.x 革命让 DAG 不再是”配置地狱”

60次阅读
Apache Airflow 调研:46K Stars 的 Python 工作流编排事实标准,3.x 革命让 DAG 不再是

一个让数据工程师用 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 工作流节点 BashOperatorPythonOperator@task 装饰器
Sensor 等待外部事件 S3KeySensorExternalTaskSensor

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 跟外部系统通信的统一抽象——S3HookBigQueryHookSnowflakeHook 等。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 最适合的三个场景:

  1. 数据团队的 ETL/ELT 编排 —— 5+ 个 DAG、跨 S3/BigQuery/Snowflake/Redshift, 需要统一调度
  2. ML pipeline 编排 —— 训练 / 评估 / 部署 / 监控全流程,3.x 的 Edge Executor 让 GPU 集群灵活接入
  3. 企业内部工作流平台 —— 跨团队共享 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 工具 ”。

参考

研究文档(引用来源参考)

(no reference document available)

正文完