Dagster入门
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
| 维度 | Airflow | Dagster |
|---|---|---|
| 核心模型 | 任务优先(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 的资产依赖图能让整个系统一目了然,新增报表只需声明它依赖哪些上游数据即可。
八、入门建议
- 安装:
pip install dagster dagit即可开始 - 先理解核心概念:资产(Asset)、Op、Job 是三个最重要的基础概念
- 从小项目开始:不要一上来就构建复杂系统,先用一个简单的 ETL 流程跑通
- 善用本地开发:
dagster dev一键启动本地环境,边写边测 - 渐进式迁移:如果已有 Airflow,不必一次性全搬。Dagster 可以集成现有 Airflow 实例,新管道用 Dagster、老旧管道逐步迁移
- 利用社区资源:Dagster 有活跃的 Slack 社区和详细的官方文档,遇到问题很容易找到答案
总结
Dagster 的根本创新在于将数据编排的关注点从"任务"转向"资产"。这种视角转换带来的不仅是技术上的便利(自动血缘、智能重跑、内建数据质量),更是一种思维方式的改变——数据工程师真正关心的,从来都是数据本身,而不是搬运数据的管道。
如果你正在开始一个新项目,或者被 Airflow 的复杂依赖和重跑搞得焦头烂额,Dagster 值得认真考虑。