基于Dagster的机器学习管道资产化实践:从特征工程到模型部署
1. 项目概述:当数据管道遇上机器学习
最近几年,无论是数据工程师还是算法工程师,都面临着一个共同的痛点:机器学习项目从实验到生产的过程,总是充满了各种“意外”。你可能在Jupyter Notebook里跑得飞快的模型,一旦要集成到每天定时运行的数据流水线里,就变得磕磕绊绊——数据版本对不上、特征计算逻辑不一致、模型再训练流程混乱,更别提监控和回滚了。这些问题消耗的精力,有时甚至超过了模型开发本身。
这正是Dagster这类现代数据编排框架的价值所在。它不仅仅是一个任务调度器,更是一个面向数据资产、强调可观测性和开发体验的平台。将机器学习工作流构建在Dagster之上,意味着你可以用声明式的方式定义数据如何流动、模型如何训练与部署,并且每一步都清晰可见、易于测试和复用。这篇内容,就是带你快速上手,看看如何用Dagster为你的机器学习项目搭建一个既健壮又灵活的基础设施。无论你是想管理一个简单的批处理预测任务,还是构建一个复杂的特征工程与模型训练管道,Dagster都能提供一套统一的思维模型和工具集。
2. 核心设计思路:从“脚本集合”到“资产化工作流”
传统机器学习项目往往由一堆松散的脚本组成: train.py , predict.py , evaluate.py ,再配上一个 requirements.txt 和一个简陋的 README 。这种模式的维护成本会随着项目复杂度的提升而指数级增长。Dagster引入的核心范式转变,是**“资产”(Asset)** 驱动的开发。在Dagster的视角里,你的机器学习流水线产出的不是一堆临时文件,而是一系列明确定义的、有依赖关系的资产。
2.1 资产化思维解析
什么是机器学习项目中的资产?一个训练好的模型文件( model.pkl )、一份处理好的特征数据集( features.parquet )、一份模型评估报告( evaluation_report.html ),甚至是一个记录本次实验参数的JSON文件,都可以被定义为一个资产。每个资产都知道:
- 它由谁生产 :即哪个计算任务(Op)负责生成它。
- 它依赖谁 :生成它需要哪些上游资产作为输入。
- 它的计算逻辑是什么 :代码实现。
- 它长什么样 :可以通过元数据(Metadata)描述其模式、统计信息、版本等。
这种思维带来的直接好处是 可观测性 。在Dagster的UI中,你可以一目了然地看到整个流水线的资产图谱:哪些资产已成功生成、它们的依赖关系如何、生成耗时多长、甚至预览资产的内容(如数据样本)。当预测结果出现异常时,你可以快速追溯到是哪个特征数据出了问题,或者是哪一版的模型导致的。
2.2 与常见MLOps工具的比较
你可能会问,这和Airflow、Luigi或Prefect有什么区别?与Airflow这种以“任务”(Task)为中心的调度器相比,Dagster的“资产”中心化设计更贴合数据与机器学习项目的本质。Airflow关心的是“任务A是否成功运行”,而Dagster更关心“特征表B这个资产是否已就绪且可用”。这种抽象使得数据血缘更加清晰,也更容易实现增量计算和局部重跑。
与MLflow这类实验追踪工具的关系则是互补的。MLflow擅长记录每次实验的参数、指标和模型,而Dagster擅长编排产生这些实验结果的整个工作流。一个典型的模式是:用Dagster编排数据获取、清洗、特征工程和模型训练任务,在训练Op中调用MLflow的API来记录实验,最后将训练好的模型(作为Dagster资产)注册到模型仓库。Dagster确保了流程的可靠执行,MLflow确保了实验的可追溯性。
3. 环境搭建与核心概念快速上手
理论说了不少,我们直接动手,从一个最简单的“鸢尾花分类”模型训练管道开始,感受Dagster的运作方式。
3.1 初始化项目与安装
首先,创建一个新的项目目录并安装Dagster。建议使用虚拟环境来管理依赖。
# 创建并进入项目目录
mkdir ml-with-dagster && cd ml-with-dagster
# 创建虚拟环境(以conda为例)
conda create -n dagster-ml python=3.9 -y
conda activate dagster-ml
# 安装dagster及其网页界面dagster-webserver
pip install dagster dagster-webserver
# 为了机器学习示例,安装scikit-learn和pandas
pip install scikit-learn pandas
接下来,创建项目的基本结构。一个典型的Dagster项目包含一个定义工作流的Python文件(例如 ml_pipeline.py )和一个配置文件 workspace.yaml 来告诉Dagster服务器从哪里加载代码。
# workspace.yaml
load_from:
- python_file:
relative_path: ml_pipeline.py
3.2 理解三大核心概念:Op, Graph, Asset
在编写第一个管道前,需要理解三个核心构建块:
- Op :这是计算的基本单元,相当于一个纯函数。它接受输入,执行一些计算,并返回输出。在ML场景中,
load_data、train_model、evaluate_model都可以是一个Op。 - Graph :将多个Op按照依赖关系连接起来,形成一个有向无环图(DAG)。它定义了“如何做”的逻辑流程。
- Asset :这是数据的具体产物。一个Graph可以被封装成一个“资产化”的作业(Job),其输出被显式定义为资产。这是我们推荐的方式。
让我们从一个基于Graph和Op的传统方式开始,再升级到Asset方式,以便理解演进过程。
第一步:用Op和Graph构建基础管道
# ml_pipeline.py
from dagster import op, graph, Out, In, job
import pandas as pd
from sklearn.datasets import load_iris
from sklearn.model_selection import train_test_split
from sklearn.ensemble import RandomForestClassifier
from sklearn.metrics import accuracy_score
import pickle
# 定义Op:加载数据
@op(out={"iris_data": Out()})
def load_iris_data():
iris = load_iris()
data = pd.DataFrame(iris.data, columns=iris.feature_names)
data['target'] = iris.target
return data
# 定义Op:分割数据集
@op(
ins={"data": In()},
out={"train_set": Out(), "test_set": Out()}
)
def split_data(data):
X = data.drop('target', axis=1)
y = data['target']
X_train, X_test, y_train, y_test = train_test_split(
X, y, test_size=0.2, random_state=42
)
return (X_train, y_train), (X_test, y_test)
# 定义Op:训练模型
@op(
ins={"train_set": In()},
out={"model": Out()}
)
def train_model(train_set):
X_train, y_train = train_set
model = RandomForestClassifier(n_estimators=100, random_state=42)
model.fit(X_train, y_train)
return model
# 定义Op:评估模型
@op(
ins={"model": In(), "test_set": In()},
out={"accuracy": Out()}
)
def evaluate_model(model, test_set):
X_test, y_test = test_set
y_pred = model.predict(X_test)
acc = accuracy_score(y_test, y_pred)
print(f"模型准确率: {acc:.4f}")
return acc
# 将Op连接成Graph
@graph
def iris_ml_graph():
data = load_iris_data()
train_set, test_set = split_data(data)
model = train_model(train_set)
evaluate_model(model, test_set)
# 将Graph封装成可执行的Job
iris_ml_job = iris_ml_graph.to_job(name="iris_ml_job")
现在,运行这个管道。在终端启动Dagster的网页界面:
dagster-webserver -f ml_pipeline.py
打开浏览器访问 http://localhost:3000 ,你就能在“Jobs”标签页看到 iris_ml_job 。点击它,然后点击“Launchpad”来启动一次运行。在UI中,你可以实时看到每个Op的执行状态、日志输出,以及最终的准确率结果。
注意 :这个传统的
@graph和@op模式清晰易懂,但它有一个关键局限:中间产物(如train_set,model)没有被持久化为可追踪的资产。如果只想重新评估模型而不重跑训练,或者想查看昨天生成的特征数据,这种模式就无能为力了。接下来,我们将其升级到更强大的Asset模式。
4. 构建资产化的机器学习管道
Asset模式是Dagster发挥其优势的推荐方式。我们将上面的流程重新定义为一系列资产。
4.1 定义资产:从原始数据到模型评估
在Asset模式下,我们关注的是“什么”被生产出来,而不是“如何”生产(尽管“如何”也包含在定义中)。我们使用 @asset 装饰器。
# ml_pipeline_assets.py
from dagster import asset, AssetIn, Output, MetadataValue
import pandas as pd
from sklearn.datasets import load_iris
from sklearn.model_selection import train_test_split
from sklearn.ensemble import RandomForestClassifier
from sklearn.metrics import accuracy_score, classification_report
import json
from pathlib import Path
# 资产1:原始鸢尾花数据集
@asset(
name="raw_iris_dataset", # 资产名称
description="从sklearn加载的原始鸢尾花数据集"
)
def raw_iris_dataset():
iris = load_iris()
data = pd.DataFrame(iris.data, columns=iris.feature_names)
data['target'] = iris.target
# 为资产附加元数据,在UI中展示
metadata = {
"样本数": len(data),
"特征列表": MetadataValue.json(list(data.columns)),
"预览": MetadataValue.md(data.head().to_markdown())
}
return Output(value=data, metadata=metadata)
# 资产2:处理后的训练集和测试集
@asset(
name="train_test_splits",
ins={"raw_data": AssetIn("raw_iris_dataset")}, # 明确依赖上游资产
output_requireds=False, # 允许返回多个输出
)
def train_test_splits(raw_data):
X = raw_data.drop('target', axis=1)
y = raw_data['target']
X_train, X_test, y_train, y_test = train_test_split(
X, y, test_size=0.2, random_state=42, stratify=y
)
# 返回一个字典,Dagster会自动将其解构为多个资产
return {
"X_train": X_train,
"y_train": y_train,
"X_test": X_test,
"y_test": y_test
}
# 资产3:训练好的模型
@asset(
name="trained_model",
ins={"training_data": AssetIn(key="train_test_splits", output_asset_key="X_train"), # 依赖特定部分
"training_labels": AssetIn(key="train_test_splits", output_asset_key="y_train")}
)
def trained_model(training_data, training_labels):
model = RandomForestClassifier(n_estimators=100, random_state=42)
model.fit(training_data, training_labels)
# 将模型保存到文件系统,并将其路径作为资产的一部分
model_path = Path("model.pkl")
with open(model_path, 'wb') as f:
pickle.dump(model, f)
metadata = {
"模型类型": "RandomForestClassifier",
"参数": MetadataValue.json({"n_estimators": 100, "random_state": 42}),
"模型文件路径": str(model_path.absolute())
}
# 返回模型对象本身,同时附加元数据
return Output(value=model, metadata=metadata)
# 资产4:模型评估报告
@asset(
name="model_evaluation_report",
ins={"model": AssetIn("trained_model"),
"test_data": AssetIn(key="train_test_splits", output_asset_key="X_test"),
"test_labels": AssetIn(key="train_test_splits", output_asset_key="y_test")}
)
def model_evaluation_report(model, test_data, test_labels):
y_pred = model.predict(test_data)
acc = accuracy_score(test_labels, y_pred)
clf_report = classification_report(test_labels, y_pred, output_dict=True)
# 生成一份详细的报告字典
report = {
"accuracy": acc,
"classification_report": clf_report,
"test_sample_count": len(test_labels)
}
# 将报告保存为JSON文件
report_path = Path("evaluation_report.json")
with open(report_path, 'w') as f:
json.dump(report, f, indent=2)
metadata = {
"准确率": f"{acc:.4f}",
"详细报告路径": str(report_path.absolute()),
"报告预览": MetadataValue.json(clf_report["weighted avg"]) # 在UI中预览加权平均指标
}
return Output(value=report, metadata=metadata)
更新 workspace.yaml 指向新的资产定义文件,然后重启 dagster-webserver 。现在,在UI中你会看到一个全新的“Assets”标签页。这里展示的正是你定义的资产图谱: raw_iris_dataset -> train_test_splits -> trained_model -> model_evaluation_report 。你可以点击任何一个资产,查看它的元数据、内容预览(如果支持)和生成它的最新运行记录。
4.2 资产分区的威力:管理多版本实验
机器学习中经常需要基于不同日期、不同参数进行多次实验。Dagster的 资产分区 功能完美契合这一场景。它允许你根据一个维度(如日期、实验ID)将资产划分为多个独立的部分,每个部分可以单独进行材料化(计算)。
假设我们想按不同的随机种子( random_state )进行多次训练实验:
from dagster import DailyPartitionsDefinition, AssetSelection, define_asset_job
# 定义一个基于“实验日期”的分区(这里简化用日期,实际可用实验ID)
experiment_partitions_def = DailyPartitionsDefinition(start_date="2024-01-01")
@asset(
name="trained_model_partitioned",
partitions_def=experiment_partitions_def, # 为该资产指定分区定义
ins={"training_data": AssetIn("train_test_splits", partition_mapping=None)}, # 分区映射
config_schema={"random_seed": int} # 接受配置参数
)
def trained_model_partitioned(context, training_data):
# 通过context获取当前分区键(如日期)
partition_key = context.asset_partition_key_for_output()
# 从配置中获取随机种子
random_seed = context.op_config["random_seed"]
X_train = training_data["X_train"]
y_train = training_data["y_train"]
model = RandomForestClassifier(n_estimators=100, random_state=random_seed)
model.fit(X_train, y_train)
# 将模型保存到包含分区键的路径中
model_path = Path(f"models/model_{partition_key}_seed{random_seed}.pkl")
model_path.parent.mkdir(parents=True, exist_ok=True)
with open(model_path, 'wb') as f:
pickle.dump(model, f)
context.add_output_metadata({
"partition": partition_key,
"random_seed": random_seed,
"model_path": str(model_path)
})
return model
# 定义一个专门用于材料化(运行)特定分区资产的Job
model_training_job = define_asset_job(
name="model_training_job",
selection=AssetSelection.keys("trained_model_partitioned"), # 选择特定资产
partitions_def=experiment_partitions_def
)
现在,你可以在UI中按分区来查看和运行 trained_model_partitioned 资产。例如,你可以只运行“2024-05-01”这个分区的训练,并传入 {"random_seed": 123} 的配置。这为管理A/B测试、超参数搜索等场景提供了清晰的结构。
5. 高级模式与生产化考量
一个可用于生产的机器学习管道,还需要考虑依赖管理、配置化、测试和监控。
5.1 资源抽象:管理外部依赖
机器学习管道通常依赖数据库、对象存储、模型注册中心等外部服务。Dagster的 资源 系统允许你将这些依赖抽象出来,并在不同环境(开发、测试、生产)中灵活切换。例如,定义一个模型存储资源:
from dagster import resource, ConfigurableResource
import boto3
from botocore.exceptions import ClientError
class ModelStorageResource(ConfigurableResource):
"""抽象模型存储,可以是本地文件系统或S3"""
bucket_name: str = ""
use_s3: bool = False
def save_model(self, model_obj, key: str):
if self.use_s3:
s3 = boto3.client('s3')
import io
buffer = io.BytesIO()
pickle.dump(model_obj, buffer)
buffer.seek(0)
s3.upload_fileobj(buffer, self.bucket_name, key)
return f"s3://{self.bucket_name}/{key}"
else:
path = Path(key)
path.parent.mkdir(parents=True, exist_ok=True)
with open(path, 'wb') as f:
pickle.dump(model_obj, f)
return str(path.absolute())
def load_model(self, key: str):
# ... 实现加载逻辑
pass
# 在资产定义中使用资源
@asset(required_resource_keys={"model_storage"})
def trained_model_with_resource(context, training_data):
model = RandomForestClassifier()
model.fit(training_data["X_train"], training_data["y_train"])
# 使用资源来保存模型,而不是硬编码路径
model_key = f"models/{context.asset_partition_key_for_output()}.pkl"
model_path = context.resources.model_storage.save_model(model, model_key)
context.add_output_metadata({"model_path": model_path})
return model
在启动管道时,通过配置文件提供资源的具体实现:
# dagster.yaml (生产环境)
resources:
model_storage:
config:
use_s3: true
bucket_name: "my-ml-models-prod"
5.2 配置化与参数管理
硬编码模型参数(如 n_estimators=100 )是不灵活的。Dagster允许通过 config_schema 将参数外置。
from dagster import Field
from pydantic import BaseModel
class ModelConfig(BaseModel):
n_estimators: int = 100
max_depth: int = None
random_state: int = 42
@asset(
ins={"training_data": AssetIn("train_test_splits")},
config_schema=ModelConfig.to_config_schema() # 使用Pydantic模型生成配置模式
)
def trained_model_configurable(context, training_data):
config = ModelConfig(**context.op_config) # 获取配置
model = RandomForestClassifier(**config.dict())
model.fit(training_data["X_train"], training_data["y_train"])
context.add_output_metadata({"config_used": config.dict()})
return model
在UI中启动该资产时,会提供一个表单让你动态填写或修改这些参数。你也可以通过代码或API以编程方式传入配置。
5.3 测试与数据验证
可靠的管道需要测试。Dagster鼓励为每个Op/Asset编写单元测试,因为它本质上是纯函数或带有明确依赖的函数。
# test_ml_assets.py
import pandas as pd
from sklearn.datasets import load_iris
from .ml_pipeline_assets import raw_iris_dataset, train_test_splits
def test_raw_iris_dataset():
# 测试资产函数
result_output = raw_iris_dataset()
data = result_output.value
# 断言数据形状和列
assert isinstance(data, pd.DataFrame)
assert data.shape == (150, 5)
assert 'target' in data.columns
def test_train_test_splits():
# 模拟上游资产
raw_data = raw_iris_dataset().value
# 调用资产函数
splits = train_test_splits(raw_data)
# 断言返回了正确的键和数据类型
assert set(splits.keys()) == {"X_train", "y_train", "X_test", "y_test"}
assert len(splits["X_train"]) == 120 # 80% of 150
此外,可以在资产上附加 数据验证 。例如,确保特征数据没有空值,或模型准确率高于某个阈值。
from dagster import AssetCheckResult, AssetCheckSpec, asset_check
@asset_check(asset="model_evaluation_report")
def check_model_accuracy(model_evaluation_report):
acc = model_evaluation_report["accuracy"]
passed = acc > 0.9
return AssetCheckResult(
passed=passed,
metadata={"accuracy": acc, "threshold": 0.9}
)
如果检查失败,在UI中该资产会显示警告标志,帮助你快速发现问题。
6. 常见问题与实战调试技巧
在实际使用Dagster构建ML管道时,你可能会遇到一些典型问题。以下是一些排查思路和技巧。
6.1 依赖解析失败与循环依赖
问题 :启动Dagster UI时,资产图谱显示解析错误,或提示循环依赖。 排查 :
- 仔细检查
AssetIn中指定的资产键名是否完全匹配@asset装饰器中的name参数。Dagster的键名是大小写敏感的。 - 确保没有形成A依赖B,B又依赖A的循环。Dagster的资产图谱必须是DAG(有向无环图)。使用UI中的“资产图谱”视图可以直观发现循环。
- 如果使用了动态输出(
output_requireds=False并返回字典),下游资产在引用时需要指定output_asset_key参数,格式为字典的键。
6.2 资产材料化缓慢或失败
问题 :运行作业时,某个资产长时间卡住或失败。 排查 :
- 查看日志 :这是第一步。在Dagster UI的运行详情页,点击失败的Op/资产,查看其标准输出和标准错误日志。大部分Python错误(如导入错误、数据格式错误)都会在这里显示。
- 检查资源 :如果资产使用了资源(如数据库、S3),确认资源配置是否正确,网络是否连通,权限是否足够。
- 内存/计算瓶颈 :对于计算密集型的训练任务,可能会因内存不足而失败。考虑使用Dagster的
@op配置将任务分发到Kubernetes或Docker容器中执行,以获取更多资源。 - 增量计算 :如果每次都要全量处理数据,速度必然慢。思考你的资产是否支持增量更新。例如,
raw_iris_dataset可能是全量的,但daily_features可以设计为只处理新日期的数据。Dagster的分区功能和I/O管理器可以协助实现增量处理逻辑。
6.3 在本地开发与生产部署的差异
问题 :管道在本地运行良好,但部署到生产环境(如Kubernetes)后出现问题。 解决思路 :
- 环境一致性 :使用Dagster的
DockerRunLauncher或K8sRunLauncher,确保代码在相同的容器镜像中运行。将依赖包版本精确锁定在requirements.txt或Pipfile中。 - 配置分离 :绝对不要将生产数据库的密码硬编码在代码里。使用Dagster的资源系统和环境变量(通过
run_config或dagster.yaml)来管理不同环境的配置。开发环境可以用本地SQLite,生产环境用云数据库。 - 文件路径 :代码中避免使用硬编码的绝对路径(如
/home/user/data.csv)。使用资源抽象(如前文的ModelStorageResource)或Dagster提供的fs_io_manager、s3_pickle_io_manager等I/O管理器来处理数据的读写,它们会自动处理不同环境下的路径问题。
6.4 与现有ML代码库的集成
问题 :已有大量训练和推理脚本,如何快速“Dagster化”? 渐进式迁移策略 :
- 从最关键的环节开始 :不要试图一次性重写所有代码。先选择痛点最明显、价值最高的环节,比如每天定时运行的预测任务或模型再训练流程,用Dagster将其包装成一个资产或作业。
- 包装现有函数 :最简单的办法是将现有的Python函数(如
train_model())直接套上@op或@asset装饰器。先让它在Dagster的框架里跑起来,获得调度和监控能力,再逐步重构其内部实现以更好地利用Dagster的特性(如资源、配置)。 - 使用
execute_in_process进行测试 :在迁移过程中,你可以使用dagster.execute_in_process在本地测试单个作业的运行,而无需启动完整的Web服务器,这非常适合快速迭代和调试。
将机器学习工作流迁移到Dagster,初期可能会感觉增加了些许复杂性,但它带来的秩序、可观测性和可维护性,对于中长期的项目发展至关重要。它迫使你以资产和数据流的视角来思考问题,这种思维模式本身就能帮助发现流程中的脆弱环节。从一个小而重要的管道开始尝试,逐步体验它如何让你的机器学习项目变得更加可靠和易于协作。
更多推荐




所有评论(0)