当你手里的定时任务从几个 crontab 涨到几十个、彼此之间还有依赖时,纯脚本很快就扛不住了。工作流引擎(Workflow Engine) 正是为解决「任务编排、依赖调度、失败重试、可视化监控」而生的基础设施。本文横向对比 Apache Airflow、Dagster、Prefect 三款主流开源引擎,从调度模型、开发体验、数据血缘到可观测性,帮你按团队规模与场景选出最合适的那一个。
一、为什么需要工作流引擎
1.1 裸 crontab 的痛点
最原始的方案是把任务塞进 crontab,简单场景够用,但一旦任务变多就会出现三个典型问题:任务之间有先后依赖却只能靠「错峰时间」硬撑;某一步失败不会自动重试,也不会告警;整条链路无法可视化,出问题只能逐个翻日志。更糟的是,当任务 A 产出要给任务 B、C 共用,而 B、C 又各有下游时,时间错峰会迅速变成「等最慢的那个再加缓冲」的玄学调参,半夜告警了你都不敢动。
# 裸 crontab:任务 B 依赖任务 A,只能靠时间错峰,A 失败 B 照样跑
0 2 * * * /opt/scripts/etl_extract.sh
30 2 * * * /opt/scripts/etl_load.sh # 假设 extract 已就绪,实则无保障
当编排逻辑复杂到这种程度,就轮到工作流引擎登场了。它把「任务」抽象成图里的节点,把「依赖」抽象成边,由调度器统一推进——这正是它与 GitHub Actions 这类 CI/CD 流水线 的本质区别:后者面向代码构建,前者面向数据与业务任务的长期编排。
二、三大引擎速览
| 引擎 | 诞生背景 | 核心抽象 | 最佳场景 |
|---|---|---|---|
| Apache Airflow | Airbnb 2015,生态最成熟 | DAG + Operator/TaskFlow | 通用数据管道、团队已熟悉 |
| Dagster | 2018,资产优先范式 | Software-defined Asset | 数据资产治理、强类型血缘 |
| Prefect | 2018,现代 Python 体验 | Flow + Task(动态图) | Python 团队、混合云调度 |
三、Apache Airflow:老牌王者
Airflow 用 DAG(有向无环图)描述任务依赖,调度器按依赖顺序执行。它的最大优势是生态:数百个内置 Operator 覆盖几乎所有数据源,社区资料也最丰富,踩坑几乎都能搜到答案。
3.1 用 TaskFlow API 定义 DAG
from datetime import datetime
from airflow.decorators import dag, task
@dag(schedule="0 2 * * *", start_date=datetime(2024, 1, 1), catchup=False)
def etl_pipeline():
@task
def extract():
return [1, 2, 3] # 模拟抽取数据
@task
def transform(data):
return [x * 2 for x in data]
@task
def load(result):
print("loaded:", result)
load(transform(extract()))
etl_pipeline()
3.2 适用与不足
Airflow 适合任务拓扑固定、需要丰富生态的团队。不足在于:DAG 是「静态」的,运行时才能决定的分支要写复杂逻辑;Web UI 偏运维视角,数据血缘展示弱;大规模部署对元数据库压力不小。如果你的任务每天重复同一张图、只是参数不同,Airflow 依然是稳妥选择;但若你希望「数据从哪来、算完落到哪」成为一等公民,就该看看 Dagster。
四、Dagster:资产优先的新范式
Dagster 把关注点从「任务」转向「资产(Asset)」——你声明的是「一张表由哪些 upstream 资产算出来」,框架自动推导依赖图。这种模式天然带来强类型与可观测的数据血缘,特别适合数据平台团队。
4.1 用 @asset 定义数据资产
from dagster import asset
@asset
def raw_orders():
return load_from_db("orders") # 上游资产
@asset
def daily_revenue(raw_orders): # 依赖自动推导
return raw_orders.groupby("day").sum()
五、Prefect:最 Pythonic 的现代体验
Prefect 用普通 Python 函数 + 装饰器即可定义流程,支持动态图(运行时才确定分支),开发体验最接近写原生 Python。它的混合执行模型(Hybrid)把调度留在云端、代码跑在你自己环境,兼顾便利与安全。
5.1 用 @flow / @task 定义流程
from prefect import flow, task
@task(retries=3, retry_delay_seconds=30)
def extract():
return [1, 2, 3]
@task
def transform(data):
return [x * 2 for x in data]
@flow(name="etl")
def etl_flow():
transform(extract())
if __name__ == "__main__":
etl_flow()
注意上面 retries=3 的重试声明——失败自动重试是工作流引擎的标配,而裸脚本要自己写这套逻辑。Prefect 的另一大优势是「结果持久化」:被 @task 装饰的函数返回值会自动缓存,重跑流程时未变的上游无需重新计算,调试长链路时尤其省时间。这也是它比 Airflow 更贴近日常 Python 开发习惯的体现。
六、五维对比,一眼选型
| 维度 | Airflow | Dagster | Prefect |
|---|---|---|---|
| 调度模型 | 静态 DAG | 资产图 | 动态 Flow |
| 开发体验 | 中等,需学概念 | 较好,类型友好 | 最佳,原生 Python |
| 数据血缘 | 弱 | 强(核心卖点) | 中等 |
| 可观测性 | 成熟(含 Prometheus 指标) | 强(Asset 视角) | 良好 |
| 学习曲线 | 较陡 | 中等 | 平缓 |
七、怎么选:一份决策清单
- 团队已用 Airflow、任务拓扑固定:继续用,别为了换而换。
- 做数据平台、重视资产血缘与质量:选 Dagster,长期治理收益大。
- Python 为主、要动态分支与低门槛:选 Prefect,上手最快。
- 只是简单 ETL + 少量任务:先评估是否真的需要重引擎,轻量编排可能更省心。
八、落地:容器化部署与监控
无论选哪款,都建议用 Docker 容器化 部署,配合 Kubernetes 做弹性伸缩,并把调度器的指标接入 Prometheus + Grafana 做任务成功率、耗时、失败告警的看板。下面是一段 Airflow 的最小 docker-compose 片段:
version: "3.8"
services:
postgres:
image: postgres:15
environment:
POSTGRES_DB: airflow
POSTGRES_PASSWORD: airflow
webserver:
image: apache/airflow:2.9
command: webserver
ports:
- "8080:8080"
depends_on: [postgres]
部署后第一步不是写 DAG,而是把监控和告警接上——否则引擎跑得再好,你也不知道哪天某个任务悄悄失败了。
九、三个真实踩坑与规避
引擎选型只是开始,真正掉坑都在运行期。下面三类是团队最常踩的:
| 踩坑 | 现象 | 规避 |
|---|---|---|
| 任务雪崩重试 | 上游抖动导致下游全量重试,数据库被打垮 | 加 retry_delay 与指数退避,限制并发 |
| DAG 导入即执行副作用 | 模块顶层写查询,调度器解析即报错 | 逻辑只放 @task 内部,顶层只声明 |
| 时区混乱 | 本地时间 vs UTC,任务半夜乱跑 | 统一设 timezone,调度计划显式标注 |
其中最隐蔽的是「DAG 导入即执行副作用」:Airflow 调度器会反复 import 你的 DAG 文件来解析依赖,如果顶层直接连数据库或发请求,解析阶段就会触发,既拖慢调度又制造噪声。正确做法是一切副作用都收进被装饰的函数里,文件顶层只做声明与导入。
结语
工作流引擎没有绝对胜负:Airflow 胜在生态与沉淀,Dagster 胜在资产治理,Prefect 胜在开发体验。选型的本质是「用团队最舒服的方式,把任务依赖管起来」——先想清楚你的任务图谱长什么样,答案自然浮现。落地时记得一条铁律:先接监控再写业务,否则引擎跑得再稳,你也只是把「不知道为何失败」换成了「不知道何时失败」。



