引言
Airflow 定义了数据编排的第一代范式:任务为中心(task-centric),一个 DAG 由若干任务组成,任务之间靠依赖连线。这个模型很直观,但它有一个根本局限:它回答的是「跑过什么」,而不是「产出了什么」。当数据团队想知道「这张表是谁产的、上游改了要不要重跑、下游报表受影响吗」,Airflow 的 DAG 图给不出答案。
Dagster 用「软件定义资产」(Software-Defined Asset)重新组织了编排的抽象:编排的单位不是任务,而是数据资产;任务只是产出资产的实现细节。这样「表 A 由表 B 加工而来」这个事实成为模型的一等公民,血缘、分区、回填、 freshness 检查都从它推导出来。
Prefect 走的是另一条路:保留任务为中心的模型,但把开发者体验与动态性做到极致。它的 Flow 是普通 Python 函数,支持运行时动态生成任务、原生 async、以及更轻量的部署模型(Work Pool 而非长期运行的 Worker 集群)。
本文先讲资产为中心与任务为中心的心智差异,再分别深入 Dagster 的资产、资源、分区、类型系统,与 Prefect 的 Flow、Deployment、Work Pool、Block,然后做三方对比并给出迁移路径。想先看 Airflow 的基础语义,可以从 Airflow DAG 调度体系 开始;想先看整体选型框架,可以从 工作流引擎全景与选型 开始。
目录
- 数据编排的第二次演进
- 任务为中心 vs 资产为中心
- Dagster 的软件定义资产
- 资产的依赖推导与血缘
- Dagster 的资源与 IO Manager
- Dagster 的类型系统与数据契约
- Dagster 的分区与回填
- Dagster 的调度与传感器
- Prefect 的 Flow 与 Task
- Prefect 的部署与 Work Pool
- Prefect 的重试、缓存与并发
- Prefect 的动态工作流与异步
- 三者的横向对比
- 从 Airflow 迁移的路径
- 元数据与可观测
- 落地路线图
- 权衡取舍
- 常见坑清单
- 小结
1. 数据编排的第二次演进
第一代编排框架(Airflow、Oozie、Azkaban)解决的是「按时间顺序把任务串起来」。第二代(Dagster、Prefect、Flyte、Metaflow)解决的是「让编排理解数据本身」:数据是什么、在哪、谁产的、变了会怎样。
这个演进的驱动力是数据规模与团队规模的增长。当管道只有 10 个任务时,任务视图足够;当有 5000 张表、跨 8 个团队时,你需要的是资产视图与血缘,否则每次上游改 schema 都要靠人肉通知。
具体到能力上,第二代框架通常提供:资产级别的血缘与分区、数据质量检查(freshness、schema 断言)、本地可运行的开发体验、以及更强的类型与契约表达。这些能力在 Airflow 里要靠插件与约定拼出来。
判断自己需不需要第二代框架的标准很简单:如果你现在维护着一张 Excel 记录「哪张表依赖哪张表」,那就是需要的信号。
2. 任务为中心 vs 资产为中心
两种模型的差异可以用同一个需求对照。需求:把 raw_orders 清洗成 clean_orders,再汇总成 daily_revenue。
任务为中心的写法描述「做什么」:
with DAG("orders", schedule="@daily") as dag:
clean = PythonOperator(task_id="clean_orders", python_callable=clean_fn)
agg = PythonOperator(task_id="daily_revenue", python_callable=agg_fn)
clean >> agg
资产为中心的写法描述「产出什么」:
@asset
def raw_orders(): ...
@asset
def clean_orders(raw_orders): ...
@asset
def daily_revenue(clean_orders): ...
差别在于依赖的来源:前者由人手工连线,后者由「函数参数引用哪个资产」自动推导。这意味着资产模型里不存在「漏连线」这个错误类别——你不可能引用一个资产而不建立依赖。
另一个差别是重跑语义。资产模型能回答「daily_revenue 需要重算,它的上游哪些是过期的」,因为每个资产有「最新一次物化的时间」。任务模型只能回答「这个任务上次成功是什么时候」。
3. Dagster 的软件定义资产
Dagster 的核心 API 是 @asset 装饰器。一个资产函数返回它产出的数据,参数声明它依赖的资产:
from dagster import asset, AssetExecutionContext
@asset(group_name="orders", compute_kind="python")
def raw_orders(context: AssetExecutionContext):
return extract_from_source(context.partition_key)
@asset(group_name="orders", deps=[raw_orders], compute_kind="spark")
def clean_orders(context: AssetExecutionContext):
df = load_parquet("s3://raw/orders")
return df.dropna(subset=["order_id", "amount"])
compute_kind 是给 UI 用的标签(标明这个资产是 Python 算的、Spark 算的、还是 dbt 建的),group_name 用于在 UI 里分组。这两个属性看起来是装饰性的,但在几百个资产的项目里,它们决定了 UI 是否可用。
资产函数不一定要返回值。如果资产由外部系统产出(比如 dbt 模型、Spark 作业),可以用 deps 声明依赖而不返回值,或者用 AssetSpec 只做元数据声明。这让「纳入编排但不重写」成为可能。
4. 资产的依赖推导与血缘
Dagster 通过两种方式建立依赖:函数参数(数据依赖)与 deps(顺序依赖)。前者会自动传递上游的返回值(如果配置了 IO Manager),后者只保证顺序。
@asset(deps=[clean_orders]) # 只依赖顺序,不接收数据
def publish_report():
run_dbt("publish_report")
血缘的价值在变更影响分析。当 raw_orders 的 schema 变了,Dagster 能立刻列出所有下游资产及其负责人(从代码所有权推导)。这是 Airflow 做不到的,因为 Airflow 里「哪张表对应哪个任务」这个映射只存在于工程师脑子里。
Dagster 还支持「资产检查」(Asset Check),它把数据质量断言挂到资产上:
from dagster import asset_check, AssetCheckResult
@asset_check(asset=clean_orders)
def no_null_order_id(clean_orders) -> AssetCheckResult:
return AssetCheckResult(
passed=clean_orders["order_id"].notna().all(),
metadata={"null_count": int(clean_orders["order_id"].isna().sum())},
)
检查失败会让资产标记为「物化但有问题」,而不是让整个管道失败。这个区分很重要:数据产出了但有质量问题,与数据没产出,是两类不同的告警。
5. Dagster 的资源与 IO Manager
资源(Resource)是外部依赖的抽象:数据库连接、Spark 会话、S3 客户端。它们通过 Definitions 注入,可以在测试时替换成假实现。
from dagster import Definitions, resource, ConfigurableResource
class WarehouseResource(ConfigurableResource):
host: str
database: str
def query(self, sql: str): ...
@asset
def clean_orders(context, warehouse: WarehouseResource):
return warehouse.query("SELECT * FROM ods.orders WHERE dt = %s",
context.partition_key)
defs = Definitions(
assets=[raw_orders, clean_orders],
resources={"warehouse": WarehouseResource(host="db.internal", database="dw")},
)
IO Manager 决定资产返回值如何持久化与如何被下游读取。默认的 IO Manager 把数据存本地文件系统(开发用),生产环境通常用 S3 IO Manager 或数据库 IO Manager。
from dagster_aws.s3 import S3PickleIOManager, S3Resource
defs = Definitions(
assets=[raw_orders, clean_orders],
resources={
"io_manager": S3PickleIOManager(
s3_resource=S3Resource(),
s3_bucket="dagster-io",
)
},
)
IO Manager 的存在让「资产之间怎么传数据」变成可配置的横切关注点,而不是每个任务自己决定。这是 Dagster 相对 Airflow 的 XCom 的一个显著优势:XCom 只能传小数据,IO Manager 可以传任意大小的 DataFrame。
6. Dagster 的类型系统与数据契约
Dagster 的类型(Dagster Type)在资产边界做校验与文档化:
from dagster import DagsterType, TypeCheck
def is_valid_orders_df(_context, value) -> bool:
return {"order_id", "amount"} <= set(value.columns)
OrdersDataFrame = DagsterType(
name="OrdersDataFrame",
type_check_fn=is_valid_orders_df,
description="订单明细表,必须包含 order_id 与 amount 列",
)
@asset
def clean_orders(raw_orders) -> OrdersDataFrame:
return raw_orders.dropna()
类型检查在生产环境是可选的(默认在物化时不跑,只在 UI 里展示),但它的文档价值很高:读代码的人能立刻知道资产的数据形状。
更实用的是与 Pandera、Pydantic 的集成,可以写出声明式的 schema 校验:
import pandera as pa
schema = pa.DataFrameSchema({
"order_id": pa.Column(str, nullable=False, unique=True),
"amount": pa.Column(float, pa.Check.gt(0)),
})
@asset
def clean_orders(raw_orders):
return schema.validate(raw_orders)
7. Dagster 的分区与回填
分区(Partition)是 Dagster 相对 Airflow 最实用的增强之一。分区把「数据的时间维度」变成一等公民,回填、增量、freshness 都从它推导。
from dagster import DailyPartitionsDefinition, asset
daily = DailyPartitionsDefinition(start_date="2026-01-01")
@asset(partitions_def=daily)
def raw_orders(context):
return extract(context.partition_key) # partition_key 是 '2026-10-07'
@asset(partitions_def=daily)
def daily_revenue(context, raw_orders):
return aggregate(raw_orders)
回填是指定分区范围的一次物化:
dagster asset materialize --select daily_revenue \
--partition 2026-09-01...2026-09-30
分区模型让「重跑 9 月的数据」变成显式的、受控的操作,而不是 Airflow 里那种「clear 一批任务状态」的隐式操作。分区还能表达「上游分区到下游分区的映射」(PartitionMapping),比如「月度资产依赖 30 个日分区」。
8. Dagster 的调度与传感器
Dagster 的调度挂在「作业」(Job)上,而作业可以从资产图自动生成:
from dagster import define_asset_job, ScheduleDefinition, Definitions
daily_job = define_asset_job("daily_orders", selection=[clean_orders, daily_revenue])
daily_schedule = ScheduleDefinition(
job=daily_job,
cron_schedule="0 2 * * *",
default_status=DefaultScheduleStatus.RUNNING,
)
defs = Definitions(assets=[...], schedules=[daily_schedule])
传感器(Sensor)用于事件驱动的触发:
from dagster import sensor, RunRequest, SensorEvaluationContext
@sensor(job=daily_job, minimum_interval_seconds=60)
def new_file_sensor(context: SensorEvaluationContext):
for key in list_new_files("s3://raw/orders/"):
if not context.has_cursor_seen(key):
yield RunRequest(run_key=key, run_config={"key": key})
sensor 的游标(cursor)是它相对 Airflow 传感器的优势:游标持久化,重启后不会重复触发。而 Airflow 的传感器重启后会重新轮询所有条件,容易重复触发下游。
9. Prefect 的 Flow 与 Task
Prefect 的核心是 @flow 与 @task。Flow 是编排单位,Task 是其中的步骤。与 Dagster 的资产模型不同,Prefect 保留任务为中心,但用装饰器把重试、缓存、日志、状态追踪都内建进函数。
from prefect import flow, task
from prefect.tasks import exponential_backoff
@task(retries=3, retry_delay_seconds=exponential_backoff(backoff_factor=2))
def fetch_orders(day: str) -> list[dict]:
return requests.get(f"{API}/orders?day={day}", timeout=30).json()
@task(cache_key_fn=task_input_hash, cache_expiration=timedelta(hours=1))
def transform(rows: list[dict]) -> list[dict]:
return [normalize(r) for r in rows]
@flow(log_prints=True)
def orders_flow(day: str):
rows = fetch_orders(day)
print(f"fetched {len(rows)} rows")
return transform(rows)
log_prints=True 把 print 自动转成结构化日志,cache_key_fn 基于输入哈希做缓存(相同输入在一小时内不重复执行)。这些在 Airflow 里都要自己实现。
10. Prefect 的部署与 Work Pool
Prefect 的部署模型是它相对 Airflow 的最大差异。Airflow 需要常驻的 Worker 集群;Prefect 的 Work Pool 是「拉取式」的,Worker 可以按需启停,甚至可以完全跑在无服务器环境(比如用 Prefect 的托管执行)。
from prefect import flow
if __name__ == "__main__":
orders_flow.serve(
name="orders-daily",
cron="0 2 * * *",
parameters={"day": "2026-10-07"},
tags=["orders", "daily"],
)
.serve() 是最轻量的部署方式:一个进程既跑调度又跑执行,适合小规模。生产环境通常用 prefect deploy 生成部署定义,推到 Work Pool,由 Worker 拉取:
# prefect.yaml 片段
deployments:
- name: orders-daily
entrypoint: flows/orders.py:orders_flow
work_pool:
name: k8s-pool
job_variables:
image: registry.internal/prefect-flows:1.2.0
cpu_request: "2"
schedules:
- cron: "0 2 * * *"
timezone: Asia/Shanghai
Work Pool 的类型决定执行环境:process(本机进程)、docker(容器)、kubernetes(Pod)。这与 Airflow 的 Executor 选型对应,但 Prefect 的 Worker 是纯拉取式的无状态进程,部署与升级更简单。
11. Prefect 的重试、缓存与并发
Prefect 的重试语义比 Airflow 更细:可以按异常类型决定是否重试,可以设置重试的时间窗口与抖动。
from prefect import task
from prefect.exceptions import RetryFailed
@task(
retries=5,
retry_delay_seconds=[1, 5, 30, 120, 600], # 每次重试的间隔
retry_jitter_factor=0.5, # 抖动,避免惊群
retry_condition_fn=lambda task, run, state: isinstance(state.result(), TimeoutError),
)
def call_flaky_api():
...
retry_jitter_factor 是容易被忽略但很重要的参数:50 个任务同时失败重试,没有抖动会同时打向下游,造成二次故障。
并发控制用标签与限制:
from prefect import flow
@flow(flow_run_name="orders-{day}")
def orders_flow(day: str): ...
# 通过 concurrency limit 限制同名任务的最大并发
Prefect 3 引入「全局并发限制」(Global Concurrency Limits),可以限制跨部署的资源使用,比如「所有访问数仓的任务最多 10 个并发」。这与 Airflow 的池概念对应。
12. Prefect 的动态工作流与异步
Prefect 的杀手级特性是「运行时动态生成任务」。这解决了一类 Airflow 难以处理的需求:任务数量取决于运行时发现的数据。
from prefect import flow, task
@task
def process_partition(part: str):
...
@flow
def dynamic_flow():
parts = discover_partitions() # 运行时才知道有多少个
futures = process_partition.map(parts) # 动态生成任务
return [f.result() for f in futures]
.map() 会为每个元素创建独立的任务实例,在 UI 上能看到每个实例的状态。这比 Airflow 的动态任务映射更灵活,因为它在 Flow 运行过程中随时可以再次 map。
Prefect 原生支持 async:
@flow
async def concurrent_fetch(urls: list[str]):
results = await asyncio.gather(*[fetch(u) for u in urls])
return results
async 让「等 100 个 HTTP 请求」变成一个协程而不是 100 个线程,资源占用低得多。这类场景在数据管道里很常见(调 100 个 API 拉数据),Airflow 处理起来要麻烦得多。
13. 三者的横向对比
| 维度 | Airflow | Dagster | Prefect |
|---|---|---|---|
| 核心抽象 | 任务与依赖 | 数据资产 | Flow 与 Task |
| 血缘 | 任务级,需插件 | 资产级,原生 | 任务级,原生 |
| 数据传递 | XCom(小数据) | IO Manager(任意) | 返回值 + Result |
| 分区与回填 | 数据区间 | 一等公民 | 需自己实现 |
| 动态任务 | 2.3 起支持映射 | 支持,偏静态 | 原生,最灵活 |
| 异步支持 | 弱 | 弱 | 原生 async |
| 部署模型 | Worker 常驻 | 常驻或 K8s | Work Pool 拉取式 |
| 数据质量检查 | 需自建 | Asset Check | 需自建 |
| 本地开发 | 需跑调度器 | dagster dev 一条命令 | flow.serve() 一条命令 |
| 生态与 Operator | 最丰富 | 中等 | 中等 |
| 学习曲线 | 中(时间语义绕) | 中高(概念多) | 低 |
三者的选择往往取决于团队构成:数据分析师为主、强调血缘与质量,选 Dagster;工程师为主、强调灵活与快速迭代,选 Prefect;需要最广的云服务集成与最大的人才池,选 Airflow。
14. 从 Airflow 迁移的路径
迁移不要一步到位。推荐的顺序是:
- 新管道直接用新框架,老管道保持不动。
- 把老管道里的「数据传递」从 XCom 改成对象存储引用(这一步与框架无关,是纯粹的改进)。
- 把「跨团队依赖」用 Asset / 血缘声明出来,先只在文档层面。
- 逐条管道迁移,优先迁移「最近改过、最容易出问题」的,而不是最老的。
迁移时的常见障碍是调度语义差异:Airflow 的 data_interval 概念在 Prefect 里没有直接对应(Prefect 用参数传递),在 Dagster 里对应分区。建议先在 Dagster 里为每个 Airflow DAG 定义一个分区定义,把区间语义显式化。
另一条路径是「共存 + 单向依赖」:用 Dagster 做上层资产编排,把 Airflow 的 DAG 用 dagster-airflow 包装成一个资产。这样血缘图里能看到 Airflow 的部分,而不用重写。缺点是调试要跨两个系统。
15. 元数据与可观测
两个框架都提供运行历史与日志,但观测的侧重点不同。
Dagster 的 UI 以「资产」为主视角:每个资产有物化历史、分区状态、最近一次的质量检查结果、上游新鲜度。它还能做「新鲜度策略」(Freshness Policy),声明「这个资产必须每 6 小时更新一次」,超时自动告警。
Prefect 的 UI 以「运行」为主视角:Flow Run 与 Task Run 的列表、状态、日志、重试次数。它的事件系统(Prefect Events)可以把运行状态推给外部系统做告警。
两者的指标都可以导出到 Prometheus,与现有的告警体系集成。告警设计的通用原则见 监控与告警设计 。核心指标是「预期存在的运行是否发生」与「运行是否在 SLA 内完成」,而不只是「运行是否失败」。
16. 落地路线图
- 第 1 周:选一条「3 到 5 个步骤、有明确输入输出表」的管道,用新框架重写一遍,对比开发体验。
- 第 2 周:把资源与凭证管理接进去(Dagster 的 Resource 或 Prefect 的 Block),确保没有硬编码凭证。
- 第 3 周:加上数据质量检查与告警,验证「产出但有问题」与「没产出」两类告警都能触发。
- 第 4 周:做一次回填演练,验证分区回填(Dagster)或参数化重跑(Prefect)的行为符合预期。
评估时重点看两个数字:一是「从改代码到看到结果」的时间(本地开发体验),二是「出问题时定位原因」的时间(可观测性)。这两个数字决定了框架在长期维护中的真实成本。
17. 权衡取舍
| 选择 | 收益 | 代价 |
|---|---|---|
| 资产为中心(Dagster) | 血缘、分区、质量检查原生 | 概念多,团队需要学习成本 |
| 任务为中心(Prefect) | 上手快,灵活 | 血缘要自己维护 |
| IO Manager 传数据 | 任意大小、可配置 | 引入一层间接,调试要理解配置 |
| Work Pool 拉取式 | 空闲成本低、部署简单 | 冷启动延迟,调试不如常驻直观 |
| 原生 async | 高并发 I/O 场景效率高 | 需要理解协程与阻塞调用的区别 |
| 动态任务生成 | 运行时决定任务数 | 静态分析难,资源预估难 |
| 自建 Prefect Server | 数据自主 | 需维护 API 与数据库 |
| 用 Prefect Cloud | 零运维 | 数据出境与成本考量 |
18. 常见坑清单
- 把 Dagster 的资产函数写成「只做副作用不返回数据」,等于退回任务模型,失去血缘价值。
- IO Manager 配置为本地文件系统却跑在分布式 Worker 上,下游读不到上游的输出。
- 分区定义改了(比如从日分区改成小时分区),历史分区的物化记录无法对应,回填行为混乱。
- 在 Prefect 的
@task里做阻塞 I/O 却用 async Flow,事件循环被阻塞,并发失效。 - Prefect 的缓存键只用任务名,不同参数的结果互相覆盖。
- Dagster 的传感器不用游标去重,重启后重复触发下游运行。
- 重试不设抖动,下游被打爆,重试放大成故障。
- 把资产的所有元数据(owners、SLA)都写在代码里却没人维护,一年后全是过期的假信息。
- 用
dagster dev的本地 IO Manager 配置直接上线,数据全写在容器本地磁盘。 - Prefect 的 Work Pool 用
process类型跑在容器里,Worker 退出后任务丢失。 - 迁移时只搬任务不搬分区语义,回填时全部处理当前区间,数据错乱。
- 把编排框架当成计算引擎,在资产函数里做重计算,Worker 资源成为瓶颈。
19. 小结
Dagster 与 Prefect 代表数据编排的两种改进方向:Dagster 往「数据资产」方向走,让编排理解数据的血缘、分区与质量;Prefect 往「开发者体验与动态性」方向走,让写管道像写普通 Python。Airflow 则守住了生态与人才池的优势。
选型的实用建议是:如果团队的核心痛点是「不知道数据从哪来、改了会怎样」,选 Dagster;如果痛点是「写管道太繁琐、动态场景做不了」,选 Prefect;如果痛点是「招不到人、云服务集成不够」,留在 Airflow。三者的能力边界正在收敛,真正的差异会越来越小,迁移成本反而成为主要考量。
下一步建议读 数据仓库与湖仓架构 ,理解编排之上的数据分层模型;如果管道里有跨小时的等待与补偿需求,则应该看 Temporal 与持久化执行 ,把「批处理编排」与「长事务编排」的边界划清。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。