Dagster 入门指南:现代数据编排框架

一、什么是 Dagster

Dagster 是一个用 Python 编写的开源数据编排平台,专为现代数据团队设计。它不只是一个"定时跑任务"的调度器,而是一个完整的数据开发环境,将管道定义、调度执行、监控告警和故障处理整合在一起。

核心哲学:从"任务优先"到"资产优先"

理解 Dagster 的关键,在于理解它和传统工具看世界的角度不同

传统工作流工具(以 Apache Airflow 为代表)关心的是任务——"先跑 A,再跑 B,然后跑 C"。至于这些任务产出了什么数据、数据之间有什么关系,它不太关心,需要你自己在脑子里维护。

Dagster 关心的是数据资产——"用户表是由订单表聚合而来的,报表又依赖用户表"。你声明资产之间的依赖关系,Dagster 自动推导执行顺序。整个数据血缘图,它自己就拼出来了,不用你手动连线。

打个比方:Airflow 像是一个"任务清单管家",它确保你按顺序把事情做完;Dagster 则像一个"数据地图",它不仅知道你要做什么,更知道每件事产出了什么、下游谁在用。

二、核心概念

2.1 资产(Assets)—— 最核心的概念

在 Dagster 中,资产就是你的数据产物——可以是数据库表、CSV 文件、机器学习模型,甚至是报表 PDF。用 @asset 装饰器声明一个资产,函数参数就是它的上游依赖:

from dagster import asset

@asset
def raw_orders():
    """从数据库拉取原始订单"""
    return fetch_from_db("orders")

@asset
def daily_summary(raw_orders):  # 参数名 = 上游资产名,依赖关系一目了然
    """基于原始订单生成每日汇总"""
    return raw_orders.groupby("date").agg({"amount": "sum"})

Dagster 看到 daily_summary(raw_orders) 就知道:要更新汇总表,得先有原始订单。这种声明式依赖让数据血缘追踪变得异常简单。

2.2 Ops 和 Graphs —— 传统任务编排方式

虽然资产是核心,Dagster 也兼容传统的任务编排模式:

  • Op:可复用的计算单元,类似 Airflow 的 Operator
  • Graph:将多个 Op 连接成完整的工作流(有向无环图)
from dagster import op, graph, Out

@op
def extract():
    return fetch_data()

@op
def transform(raw_data):
    return clean(raw_data)

@op
def load(clean_data):
    write_to_db(clean_data)

@graph
def etl_pipeline():
    load(transform(extract()))

2.3 软件定义资产(SDA)—— 创新的设计

Software-Defined Assets 是 Dagster 的一大创新:它将数据产物和生成这些产物的代码直接关联。每个资产不仅知道"自己是什么",还知道"自己是怎么来的"。这让数据血缘、版本追踪和依赖管理变得异常简单。

2.4 资产检查(Asset Checks)—— 内建数据质量

Dagster 允许给资产挂上质量断言,物化(执行)时自动校验:

from dagster import asset, AssetCheckResult

@asset
def user_table():
    return load_users()

@asset(check_specs=[
    AssetCheckSpec(name="no_empty_names", asset="user_table")
])
def check_user_table(user_table):
    empty_count = user_table["name"].isna().sum()
    return AssetCheckResult(
        passed=empty_count == 0,
        metadata={"empty_names": empty_count}
    )

2.5 Components(v1.11+)—— 低代码构建块

从 v1.11 开始,Dagster 引入了 Components 机制:用 YAML 或轻量 Python 定义可复用的管道构建块。比如引入 dbt 项目,只需几行配置:

type: dagster_dbt.DbtProjectComponent
attributes:
  project: "{{ project_root }}/dbt"
  select: "customers"

Dagster 会根据配置自动生成所有资产,大幅减少样板代码。内置的 Components 覆盖了 dbt、Fivetran、Airbyte、Sling、dlt、Power BI 等常见工具$TRAE_REF

三、Dagster 与其他框架的对比

3.1 Dagster vs Airflow

维度AirflowDagster
核心模型任务优先(DAG + Operator)资产优先(Asset + 声明式依赖)
数据血缘需手动维护,依赖图复杂时混乱自动生成,函数参数即依赖关系
重跑机制手动梳理下游,容易遗漏自动识别受影响资产,按依赖顺序重跑
类型系统弱类型,任务间传参靠 XCom强类型,支持类型注解和编译期检查
本地开发本地和线上行为有差异,调试困难本地完整运行,测试体验好
数据质量需额外集成 Great Expectations 等内建 Asset Checks
UI 可视化任务 DAG 图资产血缘图 + 执行历史 + 日志追踪

很多从 Airflow 迁移过来的团队反馈:最大的变化不是功能多寡,而是心智负担的减轻——依赖关系不再需要手动维护,重跑不再需要自己画图梳理,数据血缘可视化是白送的$TRAE_REF

3.2 Dagster vs Prefect

两者都是现代化的 Python 原生工作流框架,但定位有差异:

  • Prefect 更强调弹性(自动重试、断点续跑),适合"任务编排"场景
  • Dagster 更强调资产管理和可观测性,适合"数据工程"场景

选型建议:如果团队关注的是"数据产出了什么、谁在用",选 Dagster;如果关注的是"任务失败后怎么自动恢复",选 Prefect。

四、为什么选择 Dagster

4.1 开发体验

全 Python API,类型提示,本地开发环境和测试框架——你可以在提交到生产环境之前,在本地完整运行和调试管道。对于习惯了软件工程最佳实践的数据工程师来说,这种体验是降维打击。

4.2 资产感知

这是 Dagster 与传统工具最根本的区别。传统工具关心"任务有没有跑完",Dagster 关心"数据是不是最新的"。这种模式更符合数据工程师的思维方式——我们最终关心的是数据产物,而不仅仅是任务执行

4.3 强大的调试和可观测性

Dagster 的 Web UI(Dagit)提供了:

  • 资产依赖图:一眼看清整条数据链路
  • 数据血缘追踪:从源表到报表,上下游一目了然
  • 执行历史:每次运行的详细日志和状态
  • 错误定位:精准定位失败步骤,支持单步重跑

4.4 灵活的部署选项

  • 本地开发:pip install dagster 即可
  • 自托管:Docker、Kubernetes
  • 云服务:Dagster Cloud(含 Serverless 和 Hybrid 模式)

五、实战:构建一个完整的 ETL 流程

下面用 Dagster 实现一个从 CSV 提取 → 清洗转换 → 写入数据库的完整 ETL 管道:

import pandas as pd
from dagster import asset, job, Definitions, ScheduleDefinition

# ========== 定义资产 ==========

@asset
def raw_sales_data():
    """资产1:从 CSV 提取原始销售数据"""
    df = pd.read_csv("data/sales_2024.csv")
    return df

@asset
def cleaned_sales(raw_sales_data):
    """资产2:清洗数据(依赖 raw_sales_data)"""
    df = raw_sales_data.copy()
    # 去除空值
    df = df.dropna(subset=["order_id", "amount"])
    # 金额转正数
    df["amount"] = df["amount"].abs()
    return df

@asset
def daily_report(cleaned_sales):
    """资产3:生成日报(依赖 cleaned_sales)"""
    report = (
        cleaned_sales
        .groupby("date")
        .agg(
            total_orders=("order_id", "count"),
            total_amount=("amount", "sum"),
            avg_amount=("amount", "mean")
        )
        .reset_index()
    )
    report.to_csv("output/daily_report.csv", index=False)
    return report

# ========== 组装定义 ==========

defs = Definitions(
    assets=[raw_sales_data, cleaned_sales, daily_report],
    jobs=[define_asset_job("etl_job", selection="*")],
    schedules=[
        ScheduleDefinition(
            job=define_asset_job("daily_etl", selection="*"),
            cron_schedule="0 6 * * *",  # 每天早上6点
        )
    ],
)

启动:

pip install dagster dagit
dagster dev -f etl_pipeline.py

然后访问 http://localhost:3000,就能看到资产依赖图和执行面板。

六、高级特性

6.1 分区与回填

数据工程中常需要按时间分区处理数据,Dagster 对此有原生支持:

from dagster import asset, DailyPartitionsDefinition

@asset(
    partitions_def=DailyPartitionsDefinition(start_date="2024-01-01")
)
def daily_sales(context):
    partition_date = context.partition_key  # 如 "2024-03-15"
    return fetch_sales_for_date(partition_date)

在 UI 中可以选择任意日期范围进行回填,Dagster 会自动按分区逐个执行。

6.2 资源管理

数据库连接、API 客户端等外部依赖通过 Resources 统一管理,测试时可注入模拟资源:

from dagster import asset, resource, Definitions

@resource
def database_client():
    return create_db_connection()

@asset(required_resource_keys={"db"})
def user_table(context):
    return context.resources.db.query("SELECT * FROM users")

defs = Definitions(
    assets=[user_table],
    resources={"db": database_client},
)

6.3 声明式自动化(Declarative Automation)

这是 Dagster 近年推出的重要特性。传统编排是"我告诉你什么时候跑什么",声明式自动化是"我告诉你数据应该是什么状态,Dagster 自己判断什么时候该跑、跑什么"。你定义期望的资产新鲜度,Dagster 持续监控并在需要时自动触发物化$TRAE_REF

七、实际应用场景

数据仓库 ETL/ELT

Dagster 的资产模型与 dbt 天然契合——dbt 的每个 model 对应一个资产,Dagster 自动追溯表之间的血缘关系。Vanta 从 Airflow 迁移到 Dagster 后,关键业务数据的新鲜度从 7 小时缩短到 30 分钟以内$TRAE_REF

机器学习管道

从数据准备、特征工程到模型训练和评估,Dagster 可以管理整个 ML 生命周期,并追踪每个步骤产出的数据和模型版本。

报表生成系统

当报表之间有复杂依赖关系(比如月报依赖周报,周报依赖日报),Dagster 的资产依赖图能让整个系统一目了然,新增报表只需声明它依赖哪些上游数据即可。

八、入门建议

  1. 安装pip install dagster dagit 即可开始
  2. 先理解核心概念:资产(Asset)、Op、Job 是三个最重要的基础概念
  3. 从小项目开始:不要一上来就构建复杂系统,先用一个简单的 ETL 流程跑通
  4. 善用本地开发dagster dev 一键启动本地环境,边写边测
  5. 渐进式迁移:如果已有 Airflow,不必一次性全搬。Dagster 可以集成现有 Airflow 实例,新管道用 Dagster、老旧管道逐步迁移
  6. 利用社区资源:Dagster 有活跃的 Slack 社区和详细的官方文档,遇到问题很容易找到答案

总结

Dagster 的根本创新在于将数据编排的关注点从"任务"转向"资产"。这种视角转换带来的不仅是技术上的便利(自动血缘、智能重跑、内建数据质量),更是一种思维方式的改变——数据工程师真正关心的,从来都是数据本身,而不是搬运数据的管道。

如果你正在开始一个新项目,或者被 Airflow 的复杂依赖和重跑搞得焦头烂额,Dagster 值得认真考虑。

标签: none

添加新评论