Article
第一章:Dagster 简介与核心概念
1.1 什么是 Dagster?
| 概念 | 说明 | 注意事项 |
|---|---|---|
| Dagster | 一个现代数据编排(data orchestration)框架,用于构建、测试、监控和部署数据管道。它强调”资产驱动”(Asset-centric)的开发模式,支持从简单 ETL 到复杂数据平台的构建。 | - 不是传统 DAG 工具的简单替代,而是面向数据资产的编程模型。 - 适用于数据工程、机器学习、分析流水线等场景。 |
| 核心设计理念 | 以数据”资产”为核心组织逻辑,强调可观察性、可测试性与可维护性。 | - 强调开发者体验(DX),提供强大可视化工具 Dagit。 - 支持增量构建与依赖自动推导。 |
| Dagit | Dagster 自带的 Web UI,用于可视化 Job 结构、执行历史、资产依赖、日志等。 | - 开发阶段必备工具,通过 dagit 命令启动。- 不用于生产部署,仅用于开发与调试。 |
1.2 DAG、Pipeline、Job 的区别
| 术语 | 说明 | 注意事项 |
|---|---|---|
| DAG(有向无环图) | 通用术语,指由节点(任务)和边(依赖)构成的图结构,不能有循环依赖。 | - 是 Airflow、Prefect 等工具的核心概念。 - 在 Dagster 中,Job 是 DAG 的具体实现。 |
| Pipeline(已弃用) | 旧版 Dagster 中表示执行流程的术语,已被 Job 取代。 | - 从 Dagster 0.14+ 开始逐步弃用。 - 新项目应使用 @job 而非 @pipeline。 |
| Job | 当前版本中表示可执行流程的实体,由 Op 或 Graph 组成,定义任务执行顺序。 | - 使用 @job 装饰器定义。- 可以被 Schedule 或 Sensor 触发执行。 - 在 Dagit 中可视化为 DAG 图。 |
1.3 Asset(资产)驱动编程模型
| 概念 | 说明 | 注意事项 |
|---|---|---|
| Asset(资产) | 表示一个可识别、可观察、可重构建的数据单元,如数据库表、文件、DataFrame。 | - 是 Dagster 推荐的现代开发范式。 - 强调”数据是什么”,而非”任务怎么做”。 |
| 资产依赖 | 当一个资产的生成依赖于另一个资产时,Dagster 自动建立依赖关系图。 | - 依赖通过函数参数自动推导,无需手动声明。 - 支持跨文件资产依赖。 |
| 资产材料化(Materialization) | 指资产被实际计算并持久化的过程,如写入数据库或保存为 Parquet 文件。 | - 每次 Job 运行会记录资产的材料化事件。 - 可在 Dagit 中查看资产历史与元数据。 |
| AssetKey | 资产的唯一标识符,通常由名称和可选命名空间组成。 | - 自动生成,也可手动指定。 - 用于构建资产依赖图和调度。 |
1.4 Op(操作单元)与 Graph(图)
| 概念 | 说明 | 注意事项 |
|---|---|---|
| Op(Operation) | 最小的执行单元,类似于函数,封装一段业务逻辑。 | - 使用 @op 装饰器定义。- 可接收输入、输出数据,并记录日志。 - 是 Job 和 Asset 的基础构建块。 |
| Graph | 一组 Op 的组合,定义它们之间的执行顺序和数据流。 | - 使用 @graph 装饰器定义。- 可作为单个节点嵌入更高层级的 Job 中。 - 不包含业务逻辑,仅定义结构。 |
| Op 输入/输出 | Op 可通过参数接收输入,通过 yield 或 return 输出数据。 | - 输入类型需与上游输出匹配。 - 支持结构化数据传递。 |
| Op 执行顺序 | 通过数据依赖或显式 .after() 控制执行顺序。 | - 默认按数据流顺序执行。 - 无依赖的 Op 可并行执行。 |
1.5 Run、Execution、Schedule 与 Sensor 概念解析
| 概念 | 说明 | 注意事项 |
|---|---|---|
| Run | 一次 Job 或 Asset Job 的具体执行实例,包含开始时间、状态、日志等信息。 | - 每次手动或自动触发都会生成一个 Run。 - 可在 Dagit 中查看每个 Run 的详细信息。 |
| Execution | 指 Run 的实际执行过程,包括 Op 的调度、资源分配、错误处理等。 | - 由 Dagster 的执行引擎管理。 - 支持本地、分布式(如 Dask)执行模式。 |
| Schedule | 基于时间规则(如 cron)自动触发 Job 执行的机制。 | - 使用 @schedule 装饰器定义。- 依赖 dagster-daemon 进程运行。- 适用于周期性任务(如每日 ETL)。 |
| Sensor | 监听外部事件(如文件到达、API 回调)并决定是否触发 Job 的机制。 | - 使用 @sensor 装饰器定义。- 更灵活,适用于事件驱动场景。 - 可结合 RunRequest 动态配置执行参数。 |
第二章:开发环境搭建与项目初始化
2.1 安装 Dagster 与依赖管理
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| pip 安装 | pip install dagster | 安装核心库 | pip install dagster | - 建议使用虚拟环境(venv 或 conda)。 - pip install dagster[dev] 包含开发依赖。 |
| 安装 CLI 工具 | pip install dagster-project | 使用现代项目模板 | pip install dagster-project | - 用于 dagster project 命令创建项目。- 替代旧版 dagster init。 |
| 检查版本 | dagster --version | 验证安装成功 | dagster --version | - 确保版本 >= 1.0 以使用现代 API。 |
2.2 使用 dagster project CLI 创建项目
| 命令 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| 创建项目 | dagster project init | 初始化新 Dagster 项目 | dagster project init my_dagster_project | - 交互式选择项目模板。 - 生成标准结构(pyproject.toml, src/, tests/)。 |
| 从模板创建 | dagster project scaffold | 基于模板生成项目 | dagster project scaffold --name my_etl_job | - 适用于快速生成 Job、Asset 模板。 - 可指定模板类型(如 asset job)。 |
| 安装项目包 | pip install -e . | 将项目安装为可导入模块 | 在项目根目录执行:pip install -e . | - 必须执行,否则 Dagit 找不到模块。 - -e 表示可编辑安装。 |
2.3 项目结构详解(assets, jobs, ops, schedules, sensors)
| 目录/文件 | 说明 | 注意事项 |
|---|---|---|
src/ | 源码目录,存放 Python 模块 | - 推荐结构:src/my_project/assets/, jobs/, ops/ 等。- 包名在 pyproject.toml 中定义。 |
src/my_project/ops/ | 存放 @op 函数 | - 每个文件可包含多个 Op。 - 建议按功能模块划分。 |
src/my_project/jobs/ | 存放 @job 定义 | - Job 组合多个 Op 或 Asset。 - 可定义多个 Job。 |
src/my_project/assets/ | 存放 @asset 定义 | - 推荐现代开发方式。 - 每个文件可定义多个相关资产。 |
src/my_project/schedules/ | 存放 @schedule 定义 | - 关联特定 Job。 - 使用 cron 表达式配置频率。 |
src/my_project/sensors/ | 存放 @sensor 定义 | - 监听外部系统状态。 - 返回 RunRequest 触发 Job。 |
pyproject.toml | 项目配置文件 | - 定义包名、依赖、入口点等。 - 必须包含 [tool.dagster] 配置。 |
2.4 启动 Dagit UI 进行可视化调试
| 命令 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| 启动 Dagit | dagit -f <file.py> | 加载单个文件中的定义 | dagit -f src/my_project/jobs/my_job.py | - 适用于简单测试。 - -f 指定入口文件。 |
| 启动 Dagit(模块模式) | dagit -m <module> | 加载 Python 模块 | dagit -m my_project | - 推荐方式,自动发现所有 assets/jobs。 - 模块必须可导入(已 pip install -e .)。 |
| 指定端口 | dagit -p <port> | 自定义 Web 服务端口 | dagit -m my_project -p 3001 | - 默认端口为 3000。 - 端口冲突时使用。 |
| 查看帮助 | dagit --help | 查看所有命令选项 | dagit --help | - 了解 -w(工作区文件)等高级用法。 |
提示:启动后访问
http://localhost:3000即可查看 Dagit UI,浏览 Job、Asset、Run 历史等。
第三章:构建第一个 Pipeline / Job
3.1 定义一个简单的 Op
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
@op 装饰器 | @opdef my_op(context):... | 定义一个最小执行单元(操作) | @opdef hello_op(context):context.log.info("Hello from Op!")return "hello" | - 必须使用 @op 装饰函数。- 参数 context 是可选的,用于日志、资源等。- 返回值会作为输出传递给下游。 |
| Op 名称 | 默认使用函数名,或指定 name 参数 | 自定义 Op 在 Dagit 中显示的名称 | @op(name="custom_hello")def hello_op(context):return "hello" | - 若不指定 name,则使用函数名。 - 名称在 Job 图中可见。 |
| 无逻辑 Op | 可定义空函数用于占位或结构测试 | 用于测试组合逻辑 | @opdef placeholder_op():pass | - 不推荐在生产中使用。 - 可用于流程原型设计。 |
3.2 组合多个 Op 构建 Job
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| 数据流连接 | 在 Job 中调用 Op 函数 | 建立 Op 之间的依赖关系 | @jobdef my_job():result = op_a()op_b(result) | - 调用顺序决定执行顺序。 - 上游 Op 的返回值自动传入下游 Op 的第一个参数。 |
| 并行执行 | 多个 Op 无依赖时可并行 | 提高执行效率 | @jobdef parallel_job():op_a()op_b()op_c() | - Dagster 自动并行调度无依赖 Op。 - 需执行器支持(如多进程)。 |
| 显式依赖 | 使用 .after() 强制顺序 | 控制无数据依赖的 Op 执行顺序 | @jobdef ordered_job():op_b(after=op_a()) | - 适用于需要顺序执行但无数据传递的场景。 - 不如数据流直观,建议优先使用数据依赖。 |
3.3 使用 @job 装饰器定义执行流程
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
@job 装饰器 | @jobdef my_job():... | 将 Op 组合为可执行的 Job | @jobdef simple_etl_job():data = extract_op()processed = transform_op(data)load_op(processed) | - 是现代 Dagster 推荐方式。 - 替代已弃用的 @pipeline。 |
| Job 名称 | 使用 name 参数自定义 | 控制 Job 在 Dagit 中的显示名 | @job(name="daily_etl")def etl_job():... | - 若不指定,使用函数名。 - 名称需唯一,避免冲突。 |
| Job 描述 | 使用 description 参数 | 提供 Job 的说明信息 | @job(description="Daily ETL for sales data")def etl_job():... | - 在 Dagit 中显示,提升可读性。 - 建议为每个 Job 添加描述。 |
3.4 在 Dagit 中运行和调试 Job
| 操作 | 语法 / 步骤 | 用途 | 示例 | 注意事项 |
|---|---|---|---|---|
| 启动 Dagit | dagit -m my_project | 加载项目并启动 Web UI | 在项目根目录执行:dagit -m my_project | - 确保已 pip install -e .。- 默认监听 localhost:3000。 |
| 查看 Job 图 | 进入 Dagit → Playground → 选择 Job | 可视化 Job 的 DAG 结构 | 无代码,UI 操作 | - 节点为 Op,箭头为数据流。 - 可查看每个 Op 的输入输出类型。 |
| 配置运行参数 | 在 Playground 中填写配置 YAML | 传入 Op 的 config_schema 参数 | config:ops:my_op:config:filename: "data.csv" | - 支持 JSON 或 YAML 格式。 - 配置需符合 config_schema 定义。 |
| 执行 Job | 点击 “Launch Execution” | 触发 Job 运行 | 无代码,UI 操作 | - 成功后生成 Run 记录。 - 可实时查看日志。 |
| 查看日志 | 在 Run 详情页点击 Op 节点 | 调试执行过程 | 无代码,UI 操作 | - 日志包含 context.log.info() 输出。- 错误信息会高亮显示。 |
第四章:Op 与数据流编程
4.1 Op 输入与输出定义(Input, Output)
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| 输入参数 | Op 函数定义参数 | 接收上游 Op 的输出 | @opdef process_data(data):return data.upper() | - 参数名需与上游 Op 返回值变量名匹配(在 Job 中赋值)。 - 类型提示可增强可读性。 |
| 输出返回 | 使用 return 或 yield Output() | 向下游传递数据 | @opdef generate_data():return ["a", "b", "c"]# 或使用 Outputfrom dagster import Output@opdef emit_value():yield Output("value", "result") | - return 更简洁,适用于单输出。- yield Output(...) 支持命名输出(见多输出场景)。 |
| 多输出定义 | @op(out=...) + yield Output() | 一个 Op 产生多个输出 | from dagster import Out@op(out={"out1": Out(), "out2": Out()})def split_data():yield Output("part1", "out1")yield Output("part2", "out2") | - 必须使用 yield 和 Output(output_name=...)。- 在 Job 中通过 .out("out1") 引用特定输出。 |
4.2 数据在 Op 之间的传递
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| 直接返回传递 | 上游 return,下游接收参数 | 最简单的数据流 | @opdef op_a():return "hello"@opdef op_b(greeting):return f"{greeting} world"@jobdef my_job():result = op_a()op_b(result) | - 参数名 greeting 匹配变量名 result。- 类型自动传递,无需声明。 |
| 命名输出传递 | 使用 .out("name") 引用特定输出 | 多输出场景下的精确传递 | @jobdef multi_out_job():out1, out2 = split_data()consumer_op(out1=out1, out2=out2.out("out2")) | - 必须使用 .out("output_name") 访问命名输出。- 混合使用位置和命名参数时注意语法。 |
| 结构化数据传递 | 传递 dict、list、pandas DataFrame 等 | 处理复杂数据 | @opdef load_csv():import pandas as pdreturn pd.read_csv("data.csv")@opdef analyze(df):return df.describe() | - 支持任意 Python 对象。 - 建议使用 IO Manager 进行序列化控制。 |
4.3 Op 配置参数(config_schema)
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
config_schema 参数 | @op(config_schema=...) | 定义 Op 可接收的外部配置 | from dagster import String@op(config_schema={"filename": String})def read_file(context):filename = context.op_config["filename"]with open(filename) as f:return f.read() | - 配置在 Dagit Playground 中传入。 - 使用 context.op_config 访问。 |
| 支持类型 | String, Int, Float, Bool, Dict, Array 等 | 类型安全的配置定义 | @op(config_schema={"batch_size": Int, "shuffle": Bool})def train_model(context):size = context.op_config["batch_size"] | - 导入自 dagster。- 提供配置验证和文档生成。 |
| 默认值 | 使用 Field() 设置默认值 | 降低配置复杂度 | from dagster import Field@op(config_schema={"delay": Field(Int, default_value=5)})def wait_op(context):import timetime.sleep(context.op_config["delay"]) | - default_value 避免必填。- 在 Dagit 中仍可覆盖。 |
4.4 Op 上下文(context)的使用
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
context.log | context.log.info()/error()/debug() | 记录日志信息 | @opdef my_op(context):context.log.info("Starting process...")context.log.debug("Debug info") | - 日志在 Dagit 中可查看。 - debug 级别需启用调试模式。 |
context.op_config | context.op_config["key"] | 访问配置参数 | 见 4.3 示例 | - 仅在定义了 config_schema 时可用。- 键不存在会报错,建议用 .get()。 |
context.resources | context.resources.db_conn | 访问注入的资源 | @opdef query_db(context):conn = context.resources.databasereturn conn.execute("SELECT * FROM table") | - 资源需在 Job 中定义并配置。 - 是依赖注入的核心机制。 |
context.instance | context.instance | 访问 Dagster 实例(高级) | context.instance.run_storage | - 用于访问运行存储、资产存储等。 - 一般用于 Sensor、IO Manager 等场景。 |
4.5 Op 生命周期与日志记录
| 概念 | 说明 | 注意事项 |
|---|---|---|
| Op 执行生命周期 | 初始化 → 配置加载 → 执行(setup → op logic → cleanup)→ 结束 | - 支持 setup/cleanup 钩子(通过资源或 IO Manager)。 - 异常会中断流程并记录失败。 |
| 日志级别 | debug, info, warning, error, critical | - info 是默认级别。- debug 需在配置中启用。 |
| 日志结构化 | 支持添加元数据 | context.log.info("Processed rows", metadata={"rows": 100}) |
| 错误处理 | 抛出异常即中断 Op 执行 | try:risky_operation()except Exception as e:context.log.error(f"Failed: {e}")raise |
第五章:Asset 驱动开发(现代 Dagster 核心)
5.1 使用 @asset 定义数据资产
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
@asset 装饰器 | @assetdef my_asset():... | 定义一个可观察、可重构建的数据资产 | @assetdef raw_users():import pandas as pdreturn pd.read_csv("data/users.csv") | - 函数名即资产名(AssetKey)。 - 返回值为资产内容,用于下游依赖。 |
| 资产描述 | @asset(description="...") | 提供资产语义说明 | @asset(description="Raw user data from source system")def raw_users():... | - 在 Dagit 中显示,提升可读性。 - 建议每个资产都添加描述。 |
| 资产键(AssetKey) | 自动生成,或通过 key 参数指定 | 控制资产的唯一标识和命名空间 | @asset(key=["staging", "users"])def staged_users():... | - 默认为 [function_name]。- 支持层级结构(如 ["sales", "monthly_report"])。 |
| 输入依赖 | 参数名匹配上游资产名 | 自动建立资产依赖关系 | @assetdef cleaned_users(raw_users):return raw_users.dropna() | - 不需显式声明依赖。 - 参数名必须与上游资产函数名一致。 |
5.2 多资产(@multi_asset)与资产组(AssetGroup)
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
@multi_asset | @multi_asset(outs={...}) | 一个函数生成多个资产 | from dagster import AssetOut, multi_asset@multi_asset(outs={"table_a": AssetOut(),"table_b": AssetOut()})def extract_tables():yield Output(df_a, "table_a")yield Output(df_b, "table_b") | - 适用于批量提取或转换场景。 - 必须使用 yield Output(..., output_name)。 |
AssetOut | AssetOut(description=..., io_manager_key=...) | 配置单个资产输出属性 | @multi_asset(outs={"report": AssetOut(description="Monthly summary")})def generate_report():... | - 可设置描述、IO Manager、资产键等。 - 替代 Out 用于资产场景。 |
AssetGroup | AssetGroup(assets=[...]) | 将多个资产组织为可调度单元 | from dagster import AssetGroupmy_group = AssetGroup(assets=[raw_users, cleaned_users]) | - 用于模块化组织资产。 - 可在 workspace.yaml 中注册。 |
| 资产组名称 | AssetGroup(name="...") | 自定义组名 | AssetGroup(name="etl_group", assets=[...]) | - 在 Dagit 中用于分类显示。 - 避免名称冲突。 |
5.3 资产依赖自动推导与手动指定
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| 自动依赖推导 | 参数名匹配资产函数名 | 系统自动建立依赖图 | @assetdef model_features(cleaned_users):return cleaned_users[["age", "income"]] | - 最常用方式,简洁直观。 - 依赖关系在 Dagit 资产图中可视化。 |
| 手动指定依赖 | @asset(non_argument_deps={...}) | 声明无参数输入的依赖 | @asset(non_argument_deps={"raw_logs"})def processed_logs():# 依赖 raw_logs 但不接收其数据return transform_files() | - 适用于监听文件到达等场景。 - 依赖资产必须先材料化。 |
| 跨模块依赖 | 导入资产函数作为参数 | 在不同文件中建立依赖 | # in assets/users.pydef cleaned_users(): ...# in assets/analytics.py@assetdef user_stats(cleaned_users): ... | - 需确保模块可导入。 - 推荐使用 AssetSelection 进行 Job 构建。 |
5.4 资产元数据与观察者模式
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
MaterializeResult | yield MaterializeResult(metadata={...}) | 返回资产材料化结果及元数据 | @assetdef summary_report():count = 100yield MaterializeResult(metadata={"record_count": count,"preview": MetadataValue.md("...")}) | - 替代 return。- 支持 text, md, url, table 等元数据类型。 |
MetadataValue | MetadataValue.text(), md(), table() 等 | 构造结构化元数据 | from dagster import MetadataValueMetadataValue.table(records=[{"a": 1}, {"a": 2}]) | - 在 Dagit 中渲染为富文本。 - 提升数据可观测性。 |
| 观察者模式(Observation) | @observable_source_asset | 定义外部系统状态资产 | from dagster import observable_source_asset@observable_source_assetdef api_status(context):return DataVersion("timestamp") | - 用于监听外部系统(如 API、文件)。 - 不执行计算,仅报告版本。 |
5.5 使用 define_asset_job 构建资产同步任务
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
define_asset_job | define_asset_job(name="...", selection=...) | 创建基于资产依赖的 Job | from dagster import define_asset_joball_assets_job = define_asset_job(name="all_assets_job",selection="*") | - 自动生成执行顺序。 - selection 控制包含哪些资产。 |
AssetSelection | AssetSelection.assets(asset1, asset2) | 精确选择资产子集 | job_daily = define_asset_job(name="daily_job",selection=AssetSelection.assets(raw_users, cleaned_users)) | - 支持 upstream(), downstream(), all() 等链式操作。- 用于构建增量同步任务。 |
| 资产 Job 名称 | name 参数 | 控制 Job 显示名 | define_asset_job(name="staging_sync") | - 必须唯一。 - 可用于 Schedule/Sensor 触发。 |
| 执行参数 | 在 Dagit 中配置 config | 为资产 Job 传入配置 | ops:raw_users:config:filename: "data_v2.csv" | - 配置作用于底层 Op。 - 需资产函数支持 config_schema。 |
第六章:配置管理与资源(Resources)
6.1 使用 Config Schema 进行参数化配置
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
config_schema in @asset | @asset(config_schema={...}) | 为资产定义可配置参数 | @asset(config_schema={"year": int})def annual_report(context):year = context.op_config["year"]return load_data(year) | - 与 Op 的 config_schema 用法一致。- 通过 context.op_config 访问。 |
| 嵌套配置 | Dict 类型支持层级结构 | 组织复杂配置 | config_schema=Dict({"db": Dict({"host": str, "port": int}),"batch": Field(int, default_value=100)}) | - 在 Dagit 中以树形结构展示。 - 提升配置可维护性。 |
| 配置验证 | 类型系统自动验证 | 防止非法输入 | 若传入 "port": "abc",会报错 | - Dagster 在运行前验证配置。 - 减少运行时错误。 |
6.2 定义与注入 Resources(如数据库连接、API 客户端)
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
@resource 装饰器 | @resourcedef my_resource(context):yield ... | 定义可重用的外部服务连接 | @resourcedef database(context):conn = create_connection()yield connconn.close() | - 使用 yield 支持 cleanup。- context 可用于配置。 |
| 注入到 Job | 在 Job 定义中通过 resources_def 配置 | 将资源绑定到执行环境 | from dagster import ResourceDefinition@job(resource_defs={"database": database})def my_job():use_db_op() | - resource_defs 是字典,键为资源名。- 资源名在 Op 中通过 context.resources 访问。 |
| 在 Op 中使用 | context.resources.<name> | 获取注入的资源实例 | @opdef use_db_op(context):db = context.resources.databasedb.execute("SELECT * FROM users") | - 资源名必须与 resource_defs 中的键一致。- 资源未定义会报错。 |
6.3 资源生命周期管理(setup/cleanup)
| 概念 | 说明 | 注意事项 |
|---|---|---|
| setup 阶段 | yield 之前的代码 | - 用于初始化连接、认证、文件打开等。 - 若失败,Job 不会启动。 |
| cleanup 阶段 | yield 之后的代码 | - 无论成功或失败都会执行。 - 用于关闭连接、删除临时文件、释放锁。 |
| 异常处理 | setup 失败不触发 cleanup | try:conn = connect()yield connfinally:if 'conn' in locals():conn.close() |
| 资源作用域 | Job 级别:每个 Job 执行一次 setup/cleanup | - 适用于数据库连接、HTTP 会话等。 - 不支持跨 Job 共享状态。 |
6.4 使用 @resource 和 ResourceDefinition
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
@resource | @resource(config_schema=...) | 定义带配置的资源 | @resource(config_schema={"token": str})def api_client(context):token = context.resource_config["token"]return APIClient(token) | - context.resource_config 访问配置。- 可在不同环境传入不同 token。 |
ResourceDefinition | ResourceDefinition.hardcoded_resource(...) 或自定义类 | 创建资源的轻量方式或复用 | dev_db = ResourceDefinition.hardcoded_resource("sqlite://dev.db")prod_db = database.configured({"host": "prod.db.com"}) | - hardcoded_resource 用于固定值。- configured() 用于预配置资源。 |
| 资源覆盖 | 在 Job 或 Schedule 中替换资源 | 实现环境隔离 | @job(resource_defs={"database": dev_db})def dev_job(): ... | - 开发、测试、生产使用不同资源。 - 是实现多环境部署的关键。 |
第七章:调度(Schedules)与触发(Sensors)
7.1 基于时间的 Schedule 配置(Cron 表达式)
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
@schedule 装饰器 | @schedule(job=..., cron_schedule=...) | 定义基于时间的自动调度任务 | from dagster import schedule@schedule(job=daily_etl_job,cron_schedule="0 0 * * *")def daily_schedule(context):return {} | - 必须关联一个 Job。 - cron_schedule 支持标准 cron 语法或预设字符串(如 "daily")。 |
| Cron 表达式格式 | "min hour day month weekday" | 精确控制执行时间 | "30 8 * * 1-5" → 工作日 8:30 执行 | - 推荐使用 crontab.guru 验证。 - 注意时区,默认为系统时区。 |
| 指定时区 | execution_timezone="America/New_York" | 控制调度使用的时区 | @schedule(job=daily_job,cron_schedule="0 0 * * *",execution_timezone="Asia/Shanghai")def shanghai_daily(): ... | - 避免跨时区部署的歧义。 - 必须为 IANA 时区名称。 |
返回 run_config | 在函数中返回配置字典 | 动态传入每次运行的参数 | def daily_schedule(context):date = context.scheduled_execution_time.strftime("%Y-%m-%d")return {"ops": {"extract_op": {"config": {"date": date}}}} | - context.scheduled_execution_time 提供调度时间。- 用于日期分区等场景。 |
7.2 使用 Sensor 监听外部事件触发 Job
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
@sensor 装饰器 | @sensor(job=..., minimum_interval_seconds=...) | 定义事件驱动的触发器 | from dagster import sensor@sensor(job=daily_etl_job, minimum_interval_seconds=300)def file_arrival_sensor(context):... | - 不依赖时间,而是轮询或监听外部系统。 - minimum_interval_seconds 控制检查频率。 |
| 外部系统检查 | 自定义逻辑判断是否触发 | 如检查文件、API、数据库状态 | import osif os.path.exists("/data/latest.csv"):yield RunRequest(run_key="file_v1") | - 在函数体中实现检查逻辑。 - 可结合 context.cursor 记录上次状态。 |
run_key | RunRequest(run_key="unique_id") | 防止重复触发相同事件 | yield RunRequest(run_key="file_20250930") | - 相同 run_key 的请求不会重复执行。- 建议基于事件标识生成(如文件名、时间戳)。 |
| 多 Job 支持 | @sensor(jobs=[job1, job2]) | 一个 Sensor 触发多个 Job | @sensor(jobs=[etl_job, report_job])def combined_sensor(): ... | - 更灵活,避免多个 Sensor 冗余轮询。 - 适用于复杂事件响应。 |
7.3 Sensor 与 Run Request 编程模型
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
RunRequest | RunRequest(run_key, run_config, tags) | 请求执行一次 Job | yield RunRequest(run_key="file_123",run_config={"ops": {"load": {"config": {"path": "/data/123.csv"}}}},tags={"source": "s3"}) | - 必须 yield 返回。- run_config 覆盖 Job 默认配置。 |
SkipReason | yield SkipReason("message") | 表示本次检查无触发必要 | if not new_file_detected():yield SkipReason("No new file arrived") | - 提升可读性,记录跳过原因。 - 在 Dagit Sensor 日志中可见。 |
context.cursor | context.cursor 和 context.update_cursor() | 持久化 Sensor 状态 | last_seen = context.cursornew_file = get_latest_file_after(last_seen)if new_file:yield RunRequest(run_key=new_file)context.update_cursor(new_file.timestamp) | - 跨执行保持状态。 - 类型为字符串,需自行序列化(如 JSON)。 |
RunRequest 参数校验 | 确保 run_config 合法 | 避免触发失败 | 使用 Job.validate_config(run_config) 提前验证 | - 不做强制校验,错误会在 Run 中体现。 - 建议在复杂动态配置时添加校验。 |
7.4 动态调度与条件判断
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| 条件触发 Schedule | 在 @schedule 函数中返回 SkipReason | 满足条件才调度 | def conditional_schedule(context):if is_holiday(context.scheduled_execution_time):return SkipReason("Today is holiday")return {} | - 调度仍按 cron 触发,但可跳过执行。 - 用于节假日、维护期等场景。 |
| 动态 Sensor 决策 | 基于外部状态生成多个 RunRequest | 批量触发或参数化运行 | files = scan_new_files()for f in files:yield RunRequest(run_key=f.name,run_config={"ops": {"process": {"config": {"file": f.path}}}}) | - 一个 Sensor 检查可触发多次 Run。 - 适用于文件批量处理。 |
| 结合 Assets 状态 | 使用 context.instance 查询资产材料化历史 | 仅当依赖资产更新时触发 | latest = context.instance.get_latest_materialization_record(asset_key)if latest and (now - latest.timestamp).seconds < 3600:yield RunRequest(...) | - 实现”变更驱动”而非”时间驱动”。 - 需启用 dagster-daemon。 |
| 环境感知调度 | 根据环境变量或配置决定行为 | 多环境差异化调度 | if os.getenv("ENV") == "prod":yield RunRequest(...)else:yield SkipReason("Not in prod") | - 避免测试环境误触发生产任务。 - 建议通过资源或配置中心管理。 |
第八章:测试与调试
8.1 使用 materialize() 测试单个 Asset
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
materialize() | materialize([asset], resources={}) | 同步执行单个或多个资产 | from dagster import materializeresult = materialize([raw_users], resources={"database": mock_db})assert result.success | - 用于单元测试或快速验证。 - 返回 MaterializeResult 对象。 |
| 传入 Mock 资源 | resources={"db": mock_db} | 隔离外部依赖 | def test_raw_users():mock_data = pd.DataFrame(...)result = materialize([raw_users],resources={"source": StaticPartitionsDefinition(mock_data)}) | - 避免访问真实数据库。 - 资源必须与资产定义中使用的名称一致。 |
| 断言结果 | result.success, result.asset_materializations | 验证执行成功与输出 | assert len(result.asset_materializations) == 1output_df = result.asset_materializations[0].metadata["data"].value | - 可检查材料化元数据。 - 适用于数据质量验证。 |
| 限制 | 仅支持资产,不支持普通 Op Job | - 适用于 @asset 场景。- 普通 Job 使用 execute_in_process。 |
8.2 使用 execute_in_process() 测试 Job 执行
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
execute_in_process() | job.execute_in_process(run_config=..., resources={}) | 在单进程中执行 Job 用于测试 | result = my_job.execute_in_process(run_config={"ops": {"op_a": {"config": {"x": 5}}}},resources={"db": mock_db})assert result.success | - 同步阻塞执行,便于调试。 - 返回 ExecuteInProcessResult。 |
run_config 参数 | 传入与 Dagit 中相同的配置结构 | 测试不同配置分支 | run_config = {"ops": {"process_op": {"config": {"mode": "test"}}}} | - 结构需与 config_schema 匹配。- 可测试异常路径。 |
| 访问输出 | result.output_for_node(node_name, output_name="result") | 获取特定 Op 的输出值 | final_value = result.output_for_node("transform_op")assert final_value == expected | - node_name 为 Op 在 Job 中的名称。- 用于验证业务逻辑。 |
| 资源注入 | resources={...} | 替换真实资源为 Mock | resources={"email_client": FakeEmailClient()} | - 是单元测试的关键。 - 避免发送真实邮件或请求。 |
8.3 Mocking Resources 用于单元测试
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
| MockResource | 自定义函数或类返回模拟值 | 替代数据库、API 等 | def mock_db():yield MockConnection(data=pd.DataFrame(...))@resourcedef mock_email():yield lambda msg: print(f"Mock send: {msg}") | - 实现与真实资源相同接口。 - 可记录调用次数、参数。 |
ResourceDefinition.hardcoded_resource() | 快速创建返回固定值的资源 | 简单场景快速 Mock | from dagster import ResourceDefinitionmock_data = ResourceDefinition.hardcoded_resource(pd.DataFrame(...)) | - 适用于静态数据依赖。 - 无法模拟行为。 |
| 在测试中使用 | 传入 resources 参数 | 注入 Mock 到 Job 或 materialize | result = my_job.execute_in_process(resources={"database": mock_db, "email": mock_email}) | - 确保资源名称匹配。 - 可结合 unittest.mock.Mock 增强功能。 |
| 验证资源调用 | 记录日志或使用 Mock 对象 | 断言外部服务被正确调用 | email_spy = Mock()resources={"email_client": ResourceDefinition.hardcoded_resource(email_spy)}...email_spy.assert_called_with("alert: error") | - 验证副作用(如通知、写入)。 - 提升测试完整性。 |
8.4 日志查看与失败分析
| 方法 | 语法 / 操作 | 用途 | 示例 | 注意事项 |
|---|---|---|---|---|
| Dagit 日志面板 | 进入 Run → 点击 Op 节点 → 查看 Logs | 定位执行失败原因 | 搜索 “ERROR” 或异常堆栈 | - 日志按时间排序,包含上下文信息。 - 支持过滤日志级别。 |
context.log.error() | context.log.error("msg", metadata={...}) | 主动记录错误与上下文 | context.log.error("Query failed", metadata={"sql": query, "error": str(e)}) | - 建议添加结构化元数据。 - 便于快速定位问题。 |
| 异常传播 | 抛出异常中断 Op | 触发 Job 失败并记录 | raise ValueError("Invalid data format") | - Dagster 自动捕获并标记 Run 失败。 - 建议包装原始异常。 |
| 重试策略 | 配置执行器或使用 RetryPolicy(高级) | 应对临时性故障 | @op(retry_policy=RetryPolicy(max_retries=3)) | - 默认不重试。 - 适用于网络请求等场景。 |
| 资产材料化历史 | Dagit → Assets → 点击资产 → Materializations | 分析历史执行趋势 | 查看记录数、耗时变化 | - 可发现性能退化或数据异常。 - 支持导出元数据。 |
第九章:IO Manager 与数据持久化
9.1 IO Manager 的作用与注册
| 概念 | 说明 | 代码示例 | 注意事项 |
|---|---|---|---|
| IO Manager 作用 | 负责资产/Op 输入输出的序列化、存储与加载 | 在资产执行时自动调用 handle_output 和 load_input | - 解耦业务逻辑与存储细节。 - 支持多种格式(Parquet、数据库、API 等)。 |
@io_manager 装饰器 | 定义一个 IO Manager | from dagster import io_manager@io_managerdef my_io_manager(context):return MyIOManager() | - 返回实现 handle_output / load_input 的实例。- 可接受 context 获取配置或资源。 |
| 注册到 Job/AssetGroup | 通过 resource_defs 注册并绑定到资产 | @asset(io_manager_key="parquet_io")def sales_data(): ...resources = {"parquet_io": file_system_io_manager} | - io_manager_key 必须与资源名匹配。- 可在不同资产使用不同 IO Manager。 |
| 默认 IO Manager | 若未指定,使用 pickle_io_manager | 存储为 .pkl 文件 | - 适用于简单对象。 - 不适合生产大规模数据。 |
9.2 自定义 IO Manager 实现数据读写(如 Parquet、数据库)
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
实现 handle_output | def handle_output(context, obj): ... | 将资产输出写入外部系统 | def handle_output(self, context, obj):path = f"/data/{context.asset_key.path[-1]}.parquet"obj.to_parquet(path) | - context.asset_key 获取资产路径。- 可创建目录、压缩、分区等。 |
实现 load_input | def load_input(context): ... | 从外部系统读取数据供下游使用 | def load_input(self, context):path = f"/data/{context.asset_key.path[-1]}.parquet"return pd.read_parquet(path) | - 必须返回与 handle_output 写入格式兼容的数据。- 可添加缓存逻辑。 |
| 自定义 Parquet IO Manager | 结合 pyarrow 或 pandas | 高效存储表格数据 | from pyarrow import parquet as pqimport pyarrow as paclass ParquetIOManager:def handle_output(self, context, obj):table = pa.Table.from_pandas(obj)pq.write_table(table, context.asset_key.path[-1] + ".parquet")def load_input(self, context):return pq.read_table(context.asset_key.path[-1] + ".parquet").to_pandas() | - 支持 schema、压缩、分区。 - 比 Pickle 更高效、跨语言。 |
| 数据库 IO Manager | 写入表或视图 | 与数据库集成 | def handle_output(self, context, df):table_name = context.asset_key.path[-1]df.to_sql(table_name, self.conn, if_exists='replace') | - 需管理连接池。 - 考虑幂等性(replace vs append)。 |
9.3 使用内置 IO Managers(In-Memory、Pickle、Pandas DataFrame)
| IO Manager | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|
in_memory_io_manager | 临时传递数据,不持久化 | resources={"io_manager": in_memory_io_manager}@job(resource_defs=...) | - 适用于快速测试。 - 重启后数据丢失。 |
pickle_io_manager | 默认,序列化任意对象为 .pkl | from dagster import pickle_io_managerresources={"io_manager": pickle_io_manager} | - 简单通用。 - 安全风险(不可信数据),性能较差。 |
fs_io_manager | 基于文件系统的通用 IO Manager | from dagster import fs_io_managerresources={"io_manager": fs_io_manager} | - 存储为 .pkl 到本地磁盘。- 支持 base_dir 自定义路径。 |
db_io_manager | 与 SQLAlchemy 集成,自动管理表 | from dagster_databases import db_io_manager@asset(required_resource_keys={"database"})def users(): ... | - 需配合 database 资源。 - 自动创建表、处理类型映射。 |
pandas_file_system_io_manager | 专为 pandas.DataFrame 优化 | from dagster_pandas import pandas_file_system_io_managerresources={"io_manager": pandas_file_system_io_manager} | - 支持 CSV、Parquet 格式。 - 自动推断格式。 |
9.4 类型元(Type Metadata)与自定义类型存储
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
dagster_type 参数 | @op(out=Out(dagster_type=...)) | 指定输入/输出的 Dagster 类型 | from dagster import PythonObjectTypeclass DataFrameType(PythonObjectType):def init(self):super().init(pd.DataFrame, "DataFrame") | - 用于类型检查和自定义序列化。 - 可定义 serialization_strategy。 |
自定义 SerializationStrategy | SerializationStrategy("name", serialize, deserialize) | 控制特定类型的序列化行为 | from dagster import SerializationStrategyjson_strategy = SerializationStrategy("json",serialize=lambda obj, path: json.dump(obj, open(path, 'w')),deserialize=lambda path: json.load(open(path))) | - 可用于 JSON、YAML、Protobuf 等。 - 提升存储效率与兼容性。 |
| 在 IO Manager 中使用 | 结合 dagster_type 进行差异化处理 | if context.dagster_type.name == "DataFrame":obj.to_parquet(path)else:pickle.dump(obj, open(path, 'wb')) | - 实现多格式混合存储。 - 提升灵活性。 | |
| 类型安全检查 | context.for_type(...) | 在运行时验证类型 | if not context.dagster_type.is_subtype_of(some_type):raise CheckError("Invalid type") | - 增强健壮性。 - 适用于复杂数据管道。 |
第十章:高级主题与最佳实践
10.1 Partitioned Assets(分区资产)
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
@asset + partitions_def | @asset(partitions_def=DailyPartitionsDefinition(...)) | 按时间/值分区管理资产 | from dagster import DailyPartitionsDefinition@asset(partitions_def=DailyPartitionsDefinition(start_date="2025-01-01"))def daily_sales(context):date = context.partition_keyreturn load_for_date(date) | - context.partition_key 获取当前分区(如 "2025-09-30")。- 支持 Daily, Hourly, Static 等。 |
| 手动材料化分区 | Dagit → Assets → 选择分区 → Materialize | 只处理特定分区 | 无代码,UI 操作 | - 避免全量重跑。 - 用于修复单日数据。 |
| 全局 Job 材料化 | define_asset_job(selection=...) | 调度所有分区 | daily_job = define_asset_job("daily_job",selection=AssetSelection.assets(daily_sales)) | - 自动按依赖顺序处理分区。 - 可配置并行度。 |
| 依赖分区对齐 | 上游与下游使用相同 partitions_def | 自动建立分区级依赖 | @assetdef aggregated_sales(daily_sales):return pd.concat(daily_sales.values()) | - daily_sales 是分区资产,自动传递当前分区数据。- 简化分区流水线设计。 |
10.2 Backfill 大规模数据回填
| 方法 | 语法 / 操作 | 用途 | 示例 | 注意事项 |
|---|---|---|---|---|
| Dagit Backfill UI | Assets → Partition → “Backfill All” 或选择范围 | 批量重新材料化历史分区 | 选择日期范围 "2025-01-01" to "2025-06-30" | - 支持并行执行多个 Run。 - 监控资源消耗。 |
dagster backfill CLI | dagster backfill partition --asset-keys ... --start ... --end ... | 命令行触发回填 | dagster backfill partition --asset-keys daily_sales --start 2025-01-01 --end 2025-06-30 --repo my_repo | - 适用于自动化脚本。 - 需 dagster-daemon 运行。 |
| 资源限制 | 配置 max_concurrent_runs | 控制并发数避免系统过载 | 在 dagster.yaml 中设置:run_coordinator:max_concurrent_runs: 5 | - 回填可能消耗大量 I/O 或 CPU。 - 建议分批进行。 |
| 状态监控 | Dagit → Runs → Filter by “Backfill” | 跟踪回填进度与失败 | 检查失败 Run 的日志 | - 失败可重试。 - 建议先小范围测试。 |
10.3 Asset Checks(数据质量校验)
| 方法 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|---|---|---|---|
@asset_check | @asset_check(asset=...) | 定义针对资产的校验规则 | from dagster import asset_check@asset_check(asset=daily_sales)def daily_sales_not_empty(context, daily_sales):assert len(daily_sales) > 0, "No data loaded" | - asset 参数指定目标资产。- 函数参数接收资产值。 |
返回 AssetCheckResult | yield AssetCheckResult(passed=..., metadata=...) | 灵活控制检查结果 | yield AssetCheckResult(check_name="row_count",passed=len(df) >= 100,metadata={"count": len(df)}) | - 可返回多个检查结果。 - 支持软失败( severity=AssetCheckSeverity.WARN)。 |
| 在 Dagit 查看 | Assets → 点击资产 → Checks 标签页 | 可视化校验历史与状态 | 显示通过/失败趋势 | - 检查失败不中断材料化(可配置)。 - 用于告警或通知。 |
| 与 Alert 集成 | 结合 Sensor 或外部工具 | 失败时通知 | 使用 dagster-slack 发送消息 | - 提升数据可靠性。 - 建议关键资产都添加检查。 |
10.4 使用 Dagster Daemon 管理 Sensors/Schedules
| 组件 | 用途 | 配置方式 | 注意事项 |
|---|---|---|---|
dagster-daemon | 后台进程,轮询 Sensors/Schedules 并提交 Runs | 自动启动(开发模式)或通过 dagster-daemon run | - 必须运行才能触发调度。 - 支持高可用部署。 |
run_coordinator | 控制 Run 提交策略 | 在 dagster.yaml 中配置:run_coordinator:module: dagster.core.run_coordinatorclass: QueuedRunCoordinator | - QueuedRunCoordinator 支持排队和并发限制。- 避免资源过载。 |
run_launcher | 控制 Run 执行方式 | 配置 PresetRunLauncher 或 K8sRunLauncher | - 本地:DefaultRunLauncher。- 生产:K8s 或 Celery。 |
| 故障恢复 | Daemon 持久化状态 | 存储在 Dagster Instance(如 PostgreSQL) | - 确保数据库高可用。 - Daemon 重启后继续处理。 |
10.5 生产部署建议(Docker、Kubernetes、Dagster Cloud)
| 部署方式 | 说明 | 建议 | 注意事项 |
|---|---|---|---|
| Docker | 容器化部署核心组件 | 使用官方镜像 dagster/dagster:latest | - 统一环境,避免依赖冲突。 - 结合 docker-compose 管理多服务。 |
| Kubernetes | 生产级弹性伸缩 | 使用 dagster-k8s 和 Helm Chart | - 支持 K8sRunLauncher 分发 Job。- 配置资源请求/限制(CPU/Memory)。 |
| Dagster Cloud | 托管服务,免运维 | 注册账号,推送代码 | - 自动管理 Daemon、Dagit、执行器。 - 支持 Serverless 执行(Docker/K8s)。 |
| 代码组织 | 模块化项目结构 | assets/, jobs/, sensors/, resources/ | - 便于测试与维护。 - 使用 workspace.yaml 注册多个代码位置。 |
| 监控与告警 | 集成外部工具 | Prometheus + Grafana, Slack, PagerDuty | - 监控 Run 失败率、延迟。 - 关键 Asset Check 失败告警。 |
| 权限与安全 | 控制访问 | 使用 Dagster Cloud RBAC 或自建 Auth | - 限制敏感 Job 的执行权限。 - 加密资源配置(如密码)。 |