Article

任务调度 Dagster

更新于:2026-07-13

第一章:Dagster 简介与核心概念

1.1 什么是 Dagster?

概念说明注意事项
Dagster一个现代数据编排(data orchestration)框架,用于构建、测试、监控和部署数据管道。它强调”资产驱动”(Asset-centric)的开发模式,支持从简单 ETL 到复杂数据平台的构建。- 不是传统 DAG 工具的简单替代,而是面向数据资产的编程模型。
- 适用于数据工程、机器学习、分析流水线等场景。
核心设计理念以数据”资产”为核心组织逻辑,强调可观察性、可测试性与可维护性。- 强调开发者体验(DX),提供强大可视化工具 Dagit。
- 支持增量构建与依赖自动推导。
DagitDagster 自带的 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 可通过参数接收输入,通过 yieldreturn 输出数据。- 输入类型需与上游输出匹配。
- 支持结构化数据传递。
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 进行可视化调试

命令语法用途代码示例注意事项
启动 Dagitdagit -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 装饰器@op
def my_op(context):
  ...
定义一个最小执行单元(操作)@op
def 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可定义空函数用于占位或结构测试用于测试组合逻辑@op
def placeholder_op():
  pass
- 不推荐在生产中使用。
- 可用于流程原型设计。

3.2 组合多个 Op 构建 Job

方法语法用途代码示例注意事项
数据流连接在 Job 中调用 Op 函数建立 Op 之间的依赖关系@job
def my_job():
  result = op_a()
  op_b(result)
- 调用顺序决定执行顺序。
- 上游 Op 的返回值自动传入下游 Op 的第一个参数。
并行执行多个 Op 无依赖时可并行提高执行效率@job
def parallel_job():
  op_a()
  op_b()
  op_c()
- Dagster 自动并行调度无依赖 Op。
- 需执行器支持(如多进程)。
显式依赖使用 .after() 强制顺序控制无数据依赖的 Op 执行顺序@job
def ordered_job():
  op_b(after=op_a())
- 适用于需要顺序执行但无数据传递的场景。
- 不如数据流直观,建议优先使用数据依赖。

3.3 使用 @job 装饰器定义执行流程

方法语法用途代码示例注意事项
@job 装饰器@job
def my_job():
  ...
将 Op 组合为可执行的 Job@job
def 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

操作语法 / 步骤用途示例注意事项
启动 Dagitdagit -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 的输出@op
def process_data(data):
  return data.upper()
- 参数名需与上游 Op 返回值变量名匹配(在 Job 中赋值)。
- 类型提示可增强可读性。
输出返回使用 returnyield Output()向下游传递数据@op
def generate_data():
  return ["a", "b", "c"]

# 或使用 Output
from dagster import Output
@op
def 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")
- 必须使用 yieldOutput(output_name=...)
- 在 Job 中通过 .out("out1") 引用特定输出。

4.2 数据在 Op 之间的传递

方法语法用途代码示例注意事项
直接返回传递上游 return,下游接收参数最简单的数据流@op
def op_a():
  return "hello"

@op
def op_b(greeting):
  return f"{greeting} world"

@job
def my_job():
  result = op_a()
  op_b(result)
- 参数名 greeting 匹配变量名 result
- 类型自动传递,无需声明。
命名输出传递使用 .out("name") 引用特定输出多输出场景下的精确传递@job
def multi_out_job():
  out1, out2 = split_data()
  consumer_op(out1=out1, out2=out2.out("out2"))
- 必须使用 .out("output_name") 访问命名输出。
- 混合使用位置和命名参数时注意语法。
结构化数据传递传递 dict、list、pandas DataFrame 等处理复杂数据@op
def load_csv():
  import pandas as pd
  return pd.read_csv("data.csv")

@op
def 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 time
  time.sleep(context.op_config["delay"])
- default_value 避免必填。
- 在 Dagit 中仍可覆盖。

4.4 Op 上下文(context)的使用

方法语法用途代码示例注意事项
context.logcontext.log.info()/error()/debug()记录日志信息@op
def my_op(context):
  context.log.info("Starting process...")
  context.log.debug("Debug info")
- 日志在 Dagit 中可查看。
- debug 级别需启用调试模式。
context.op_configcontext.op_config["key"]访问配置参数见 4.3 示例- 仅在定义了 config_schema 时可用。
- 键不存在会报错,建议用 .get()
context.resourcescontext.resources.db_conn访问注入的资源@op
def query_db(context):
  conn = context.resources.database
  return conn.execute("SELECT * FROM table")
- 资源需在 Job 中定义并配置。
- 是依赖注入的核心机制。
context.instancecontext.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 装饰器@asset
def my_asset():
  ...
定义一个可观察、可重构建的数据资产@asset
def raw_users():
  import pandas as pd
  return 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"])。
输入依赖参数名匹配上游资产名自动建立资产依赖关系@asset
def 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)
AssetOutAssetOut(description=..., io_manager_key=...)配置单个资产输出属性@multi_asset(
  outs={
    "report": AssetOut(description="Monthly summary")
  }
)
def generate_report():
  ...
- 可设置描述、IO Manager、资产键等。
- 替代 Out 用于资产场景。
AssetGroupAssetGroup(assets=[...])将多个资产组织为可调度单元from dagster import AssetGroup
my_group = AssetGroup(
  assets=[raw_users, cleaned_users]
)
- 用于模块化组织资产。
- 可在 workspace.yaml 中注册。
资产组名称AssetGroup(name="...")自定义组名AssetGroup(name="etl_group", assets=[...])- 在 Dagit 中用于分类显示。
- 避免名称冲突。

5.3 资产依赖自动推导与手动指定

方法语法用途代码示例注意事项
自动依赖推导参数名匹配资产函数名系统自动建立依赖图@asset
def 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.py
def cleaned_users(): ...

# in assets/analytics.py
@asset
def user_stats(cleaned_users): ...
- 需确保模块可导入。
- 推荐使用 AssetSelection 进行 Job 构建。

5.4 资产元数据与观察者模式

方法语法用途代码示例注意事项
MaterializeResultyield MaterializeResult(metadata={...})返回资产材料化结果及元数据@asset
def summary_report():
  count = 100
  yield MaterializeResult(
    metadata={
      "record_count": count,
      "preview": MetadataValue.md("...")
    }
  )
- 替代 return
- 支持 text, md, url, table 等元数据类型。
MetadataValueMetadataValue.text(), md(), table()构造结构化元数据from dagster import MetadataValue
MetadataValue.table(
  records=[{"a": 1}, {"a": 2}]
)
- 在 Dagit 中渲染为富文本。
- 提升数据可观测性。
观察者模式(Observation)@observable_source_asset定义外部系统状态资产from dagster import observable_source_asset
@observable_source_asset
def api_status(context):
  return DataVersion("timestamp")
- 用于监听外部系统(如 API、文件)。
- 不执行计算,仅报告版本。

5.5 使用 define_asset_job 构建资产同步任务

方法语法用途代码示例注意事项
define_asset_jobdefine_asset_job(name="...", selection=...)创建基于资产依赖的 Jobfrom dagster import define_asset_job
all_assets_job = define_asset_job(
  name="all_assets_job",
  selection="*"
)
- 自动生成执行顺序。
- selection 控制包含哪些资产。
AssetSelectionAssetSelection.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 装饰器@resource
def my_resource(context):
  yield ...
定义可重用的外部服务连接@resource
def database(context):
  conn = create_connection()
  yield conn
  conn.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>获取注入的资源实例@op
def use_db_op(context):
  db = context.resources.database
  db.execute("SELECT * FROM users")
- 资源名必须与 resource_defs 中的键一致。
- 资源未定义会报错。

6.3 资源生命周期管理(setup/cleanup)

概念说明注意事项
setup 阶段yield 之前的代码- 用于初始化连接、认证、文件打开等。
- 若失败,Job 不会启动。
cleanup 阶段yield 之后的代码- 无论成功或失败都会执行。
- 用于关闭连接、删除临时文件、释放锁。
异常处理setup 失败不触发 cleanuptry:
  conn = connect()
  yield conn
finally:
  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。
ResourceDefinitionResourceDefinition.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 os
if os.path.exists("/data/latest.csv"):
  yield RunRequest(run_key="file_v1")
- 在函数体中实现检查逻辑。
- 可结合 context.cursor 记录上次状态。
run_keyRunRequest(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 编程模型

方法语法用途代码示例注意事项
RunRequestRunRequest(run_key, run_config, tags)请求执行一次 Jobyield RunRequest(
  run_key="file_123",
  run_config={"ops": {"load": {"config": {"path": "/data/123.csv"}}}},
  tags={"source": "s3"}
)
- 必须 yield 返回。
- run_config 覆盖 Job 默认配置。
SkipReasonyield SkipReason("message")表示本次检查无触发必要if not new_file_detected():
  yield SkipReason("No new file arrived")
- 提升可读性,记录跳过原因。
- 在 Dagit Sensor 日志中可见。
context.cursorcontext.cursorcontext.update_cursor()持久化 Sensor 状态last_seen = context.cursor
new_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 materialize
result = 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) == 1
output_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={...}替换真实资源为 Mockresources={"email_client": FakeEmailClient()}- 是单元测试的关键。
- 避免发送真实邮件或请求。

8.3 Mocking Resources 用于单元测试

方法语法用途代码示例注意事项
MockResource自定义函数或类返回模拟值替代数据库、API 等def mock_db():
  yield MockConnection(data=pd.DataFrame(...))

@resource
def mock_email():
  yield lambda msg: print(f"Mock send: {msg}")
- 实现与真实资源相同接口。
- 可记录调用次数、参数。
ResourceDefinition.hardcoded_resource()快速创建返回固定值的资源简单场景快速 Mockfrom dagster import ResourceDefinition
mock_data = ResourceDefinition.hardcoded_resource(pd.DataFrame(...))
- 适用于静态数据依赖。
- 无法模拟行为。
在测试中使用传入 resources 参数注入 Mock 到 Job 或 materializeresult = 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_outputload_input- 解耦业务逻辑与存储细节。
- 支持多种格式(Parquet、数据库、API 等)。
@io_manager 装饰器定义一个 IO Managerfrom dagster import io_manager
@io_manager
def 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_outputdef 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_inputdef 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 pq
import pyarrow as pa

class 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默认,序列化任意对象为 .pklfrom dagster import pickle_io_manager
resources={"io_manager": pickle_io_manager}
- 简单通用。
- 安全风险(不可信数据),性能较差。
fs_io_manager基于文件系统的通用 IO Managerfrom dagster import fs_io_manager
resources={"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_manager
resources={"io_manager": pandas_file_system_io_manager}
- 支持 CSV、Parquet 格式。
- 自动推断格式。

9.4 类型元(Type Metadata)与自定义类型存储

方法语法用途代码示例注意事项
dagster_type 参数@op(out=Out(dagster_type=...))指定输入/输出的 Dagster 类型from dagster import PythonObjectType
class DataFrameType(PythonObjectType):
  def init(self):
    super().init(pd.DataFrame, "DataFrame")
- 用于类型检查和自定义序列化。
- 可定义 serialization_strategy
自定义 SerializationStrategySerializationStrategy("name", serialize, deserialize)控制特定类型的序列化行为from dagster import SerializationStrategy
json_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_key
  return 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自动建立分区级依赖@asset
def aggregated_sales(daily_sales):
  return pd.concat(daily_sales.values())
- daily_sales 是分区资产,自动传递当前分区数据。
- 简化分区流水线设计。

10.2 Backfill 大规模数据回填

方法语法 / 操作用途示例注意事项
Dagit Backfill UIAssets → Partition → “Backfill All” 或选择范围批量重新材料化历史分区选择日期范围 "2025-01-01" to "2025-06-30"- 支持并行执行多个 Run。
- 监控资源消耗。
dagster backfill CLIdagster 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 参数指定目标资产。
- 函数参数接收资产值。
返回 AssetCheckResultyield 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_coordinator
  class: QueuedRunCoordinator
- QueuedRunCoordinator 支持排队和并发限制。
- 避免资源过载。
run_launcher控制 Run 执行方式配置 PresetRunLauncherK8sRunLauncher- 本地: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 的执行权限。
- 加密资源配置(如密码)。