工作流引擎选型实战:Airflow/Dagster/Prefect 对比

当你手里的定时任务从几个 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 AirflowAirbnb 2015,生态最成熟DAG + Operator/TaskFlow通用数据管道、团队已熟悉
Dagster2018,资产优先范式Software-defined Asset数据资产治理、强类型血缘
Prefect2018,现代 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 开发习惯的体现。

六、五维对比,一眼选型

维度AirflowDagsterPrefect
调度模型静态 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 胜在开发体验。选型的本质是「用团队最舒服的方式,把任务依赖管起来」——先想清楚你的任务图谱长什么样,答案自然浮现。落地时记得一条铁律:先接监控再写业务,否则引擎跑得再稳,你也只是把「不知道为何失败」换成了「不知道何时失败」。

上一篇 Java 结构化并发实战:StructuredTaskScope 新范式
下一篇 CSS 容器查询实战:组件级响应式布局新范式