Dagster 与 Prefect 数据编排

本文对比 Dagster 与 Prefect 两代数据编排框架与 Airflow 的差异,回答资产为中心与任务为中心该怎么选、IO Manager 与资源如何组织、分区与回填怎么落地。覆盖软件定义资产、类型系统与数据契约、资源与 IO Manager、分区回填、Prefect 的 Flow 与 Deployment、Work Pool、重试缓存与抖动、动态工作流与异步,并给出从 Airflow 迁移的路径。

引言

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 调度体系 开始;想先看整体选型框架,可以从 工作流引擎全景与选型 开始。

目录

  1. 数据编排的第二次演进
  2. 任务为中心 vs 资产为中心
  3. Dagster 的软件定义资产
  4. 资产的依赖推导与血缘
  5. Dagster 的资源与 IO Manager
  6. Dagster 的类型系统与数据契约
  7. Dagster 的分区与回填
  8. Dagster 的调度与传感器
  9. Prefect 的 Flow 与 Task
  10. Prefect 的部署与 Work Pool
  11. Prefect 的重试、缓存与并发
  12. Prefect 的动态工作流与异步
  13. 三者的横向对比
  14. 从 Airflow 迁移的路径
  15. 元数据与可观测
  16. 落地路线图
  17. 权衡取舍
  18. 常见坑清单
  19. 小结

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. 三者的横向对比

维度AirflowDagsterPrefect
核心抽象任务与依赖数据资产Flow 与 Task
血缘任务级,需插件资产级,原生任务级,原生
数据传递XCom(小数据)IO Manager(任意)返回值 + Result
分区与回填数据区间一等公民需自己实现
动态任务2.3 起支持映射支持,偏静态原生,最灵活
异步支持弱弱原生 async
部署模型Worker 常驻常驻或 K8sWork Pool 拉取式
数据质量检查需自建Asset Check需自建
本地开发需跑调度器dagster dev 一条命令flow.serve() 一条命令
生态与 Operator最丰富中等中等
学习曲线中(时间语义绕)中高(概念多)低

三者的选择往往取决于团队构成:数据分析师为主、强调血缘与质量,选 Dagster;工程师为主、强调灵活与快速迭代,选 Prefect;需要最广的云服务集成与最大的人才池,选 Airflow。

14. 从 Airflow 迁移的路径

迁移不要一步到位。推荐的顺序是:

  1. 新管道直接用新框架,老管道保持不动。
  2. 把老管道里的「数据传递」从 XCom 改成对象存储引用(这一步与框架无关,是纯粹的改进)。
  3. 把「跨团队依赖」用 Asset / 血缘声明出来,先只在文档层面。
  4. 逐条管道迁移,优先迁移「最近改过、最容易出问题」的,而不是最老的。

迁移时的常见障碍是调度语义差异: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. 常见坑清单

  1. 把 Dagster 的资产函数写成「只做副作用不返回数据」,等于退回任务模型,失去血缘价值。
  2. IO Manager 配置为本地文件系统却跑在分布式 Worker 上,下游读不到上游的输出。
  3. 分区定义改了(比如从日分区改成小时分区),历史分区的物化记录无法对应,回填行为混乱。
  4. 在 Prefect 的 @task 里做阻塞 I/O 却用 async Flow,事件循环被阻塞,并发失效。
  5. Prefect 的缓存键只用任务名,不同参数的结果互相覆盖。
  6. Dagster 的传感器不用游标去重,重启后重复触发下游运行。
  7. 重试不设抖动,下游被打爆,重试放大成故障。
  8. 把资产的所有元数据(owners、SLA)都写在代码里却没人维护,一年后全是过期的假信息。
  9. 用 dagster dev 的本地 IO Manager 配置直接上线,数据全写在容器本地磁盘。
  10. Prefect 的 Work Pool 用 process 类型跑在容器里,Worker 退出后任务丢失。
  11. 迁移时只搬任务不搬分区语义,回填时全部处理当前区间,数据错乱。
  12. 把编排框架当成计算引擎,在资产函数里做重计算,Worker 资源成为瓶颈。

19. 小结

Dagster 与 Prefect 代表数据编排的两种改进方向:Dagster 往「数据资产」方向走,让编排理解数据的血缘、分区与质量;Prefect 往「开发者体验与动态性」方向走,让写管道像写普通 Python。Airflow 则守住了生态与人才池的优势。

选型的实用建议是:如果团队的核心痛点是「不知道数据从哪来、改了会怎样」,选 Dagster;如果痛点是「写管道太繁琐、动态场景做不了」,选 Prefect;如果痛点是「招不到人、云服务集成不够」,留在 Airflow。三者的能力边界正在收敛,真正的差异会越来越小,迁移成本反而成为主要考量。

下一步建议读 数据仓库与湖仓架构 ,理解编排之上的数据分层模型;如果管道里有跨小时的等待与补偿需求,则应该看 Temporal 与持久化执行 ,把「批处理编排」与「长事务编排」的边界划清。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「工作流引擎」更多文章

  1. 工作流成本优化
  2. 执行器与资源隔离
  3. 调度、回填与补数