进阶 dagster.io 2026-10-10 01:48:56 · 6 阅读

第19章 从资产到样本:重塑现代 MLOps 编排的 Metaxy 框架

了解如何使用 Metaxy 在 Dagster 上构建具备样本级粒度的多模态数据管道。 我叫 Daniel Gafni,是 Anam 的 MLOps 工程师。 Anam 在打造一个实时交互式数字人平台,支撑产品的核心组件之一就是我们自研的视频生成模型。 我们用定制数据集训练它,这需要对视频和音频数据做各种预处理:用 ML 模型提取 embedding、调用外部 API 做标注和数据合成等等。当然,我们用 Dagster 把这些步骤编排成数据资产。 这篇博客将分享我们如何用 Dagster 加上我们新开源的 Metaxy 框架,解决这些管道的样本级版本管理问题。 Metaxy 把 Dagster 这类通常只在表(或资产)级别操作的编排器,与 Ray 等底层计算引擎连接起来,让我们在每一步只处理真正需要处理的样本,一个都不多。

一个小改动

几个月前,我们决定对数据准备管道做一个很小的改动。当时我们用的是自研的数据版本管理系统,会为每个样本记录版本。系统基于手动指定的 code_version 及其上游步骤来计算每一步的指纹(是不是听着很熟悉?)。它和 Dagster 数据版本管理的唯一区别是粒度:我们为数据集的每一行计算版本,而 Dagster 只支持资产级版本管理。 我们要做的改动非常简单:用新的分辨率裁剪视频。这意味着要改一下裁剪阶段的 code_version,下游步骤会自动重新计算。但随之而来的是一个令人头疼的结果。裁剪步骤之后,管道分成了两条分支: 其中一半的下游步骤根本不用裁剪后的视频帧,只处理音频部分。但我们的版本管理系统不知道这个细节,照样把它们全部重算。这意味着要在整个训练数据集上重新跑音频 ML 模型:又贵又完全没必要。 那一刻我意识到,我们这种朴素的数据版本管理思路有问题。Metaxy——也就是这篇博客要介绍的项目——就这样诞生了。

多模态数据一瞥

随着 AI 吞噬软件世界,越来越多的团队和组织开始接触多模态数据管道。与传统数据工程工作流不同,这些管道处理的不仅仅是表格,还有文本、图像、音频、视频、向量 embedding、医疗数据等等。 多模态数据管道可能非常独特,需求和复杂度因场景而异。无论是通过 HTTP 调用 AI API、在本地跑 ML 推理,还是只是调用 ffmpeg,有一点是共通的:计算和 I/O 成本会迅速飙升。 传统(表格型)数据管道重新执行通常花不了多少钱。当然,大数据是存在的,Spark 作业可以查询 PB 级的表格数据,但实际上真正遇到这类问题的团队很少。这也是 Small Data 运动成功的一大原因:Snowflake 扫描的中位数读取量不到 100MB,80% 的组织数据量不足 10TB!所以重跑表格管道通常没什么问题,而且比起实现增量处理,重跑要简单得多。 多模态管道完全是另一回事。它们需要的计算量和数据搬运量要高出几个数量级。一不小心对整个数据集重跑了 Whisper 语音转录步骤?恭喜,1 万美元打水漂了! 因此在多模态管道中,增量方案不是可选项,而是刚需。而事实证明,这玩意儿实现起来相当复杂。

Metaxy 简介

Metaxy 补上了缺失的一环,把只认数据集(或分区)的 Dagster 世界,和必须处理单个样本的计算世界连接了起来。 Metaxy 有两个独特之处: 它能够追踪部分数据更新。 它不依赖特定基础设施,可以接入任何用 Python 写的数据管道。 它采用了与 Dagster 相同的数据版本计算方式,并做了扩展: 支持批量处理,一次处理数百万行 可以跑在远程数据库上,也可以在本地运行 不绑定特定的 dataframe 引擎或数据库 感知数据字段:Metaxy 里的每个数据版本不是一个字符串,而是一个字典。

数据字段

Metaxy 的主要目标之一是实现细粒度的部分数据版本管理,让部分更新能被正确识别。在 Metaxy 中,除了普通的元数据列,用户还可以定义并管理版本的数据字段: 数据字段可以任意定义,由用户自行设计。关键在于:数据字段描述的是数据本身(比如 mp4 文件),而不是表格元数据。 然后跑几行 Python 代码: with store: increment = store.resolve_update("video/face_crop") 这背后做了大量工作: 联接上游步骤的状态表 为每一行计算期望的数据版本——这是最复杂的一步。复杂在于每个版本是个字典,字典里每个字段可能只依赖各自不同的上游字段子集 加载该步骤的状态表,与期望版本比对 向用户返回新增、过期和孤立的样本 拿到 increment 对象后(顺便说一句,它可以是惰性的!),用户可以决定每类样本怎么处理:通常 new 和 stale 会被重新处理,orphaned 可能被删除。 Metaxy 通过感知部分数据依赖,解决了前面那个"小改动"问题。 看 3 个 Metaxy feature(Metaxy 把每一步产出的数据叫 feature):视频文件(video/full)、Whisper 转录(transcript/whisper)、以及按人脸裁剪的视频文件(video/face_crop): 音频和帧的信息路径用不同颜色标出。可以看到 feature 之间存在清晰的字段级(部分数据)依赖。每个字段的版本由它依赖的字段版本计算得出,同一 feature 的字段版本再合并成 feature 版本。 显然,transcript/whisper 的 text 字段只依赖 video/full 的 audio 字段。如果我们调整了 video/full 的分辨率,transcript/whisper 根本不用重算。 Metaxy 能识别这类"无关"更新,跳过那些字段未受上游变更影响的下游 feature 的重算。这是通过为每个 feature 的每个样本的每个字段单独记录数据版本实现的: idmetaxy_provenance_by_fieldvideo_001{"audio": "a7f3c2d8", "frames": "b9e1f4a2"}video_002{"audio": "d4b8e9c1", "frames": "f2a6d7b3"}video_003{"audio": "c9f2a8e4", "frames": "e7d3b1c5"}video_004{"audio": "b1e4f9a7", "frames": "a8c2e6d9"}

可组合性的挑战

前面提到过,增量管道形态各异:它们往往运行在独特的环境中,需要特定的基础设施、不同的云厂商(包括 Neocloud),或者 Ray、Modal 这类扩展引擎。 Metaxy 的目标之一是像 Dagster 一样通用中立,支持这些多样的场景。它必须做到可插拔,才能被不同的用户和组织使用。 事实证明,这是可行的!Metaxy 95% 的元数据管理工作实现方式与数据库无关,甚至可以借助 Polars 或 DuckDB 在本地运行! 这要归功于 Ibis 和 Narwhals 两个项目投入的巨大心血。Ibis 为 20 多种数据库提供统一的 Python 接口(不是常见的 DataFrame API),Narwhals 则为各种 DataFrame 引擎(Pandas、Polars、DuckDB、Ibis 等)做同样的事,把一切收敛到 Polars API 的一个子集上。 Narwhals(或 Polars)表达式是构建程序化查询的利器。Metaxy 的版本引擎大部分都用 Narwhals 表达式实现,只有少数特殊部分不得不下沉到特定后端。 这一点怎么强调都不为过。全新一代可组合的数据工具完全可以构建在 Narwhals 之上——当然,没有 Apache Arrow,这一切都不可能。

Metaxy 与 Dagster

Metaxy 自然是为配合 Dagster 使用而设计的。它不仅从 Dagster 的数据版本设计中汲取了灵感,还天然继承了 Dagster 的其他特性:声明式和面向资产。看看这个 API,任何 Dagster 用户都会觉得无比眼熟: import metaxy as mx spec = mx.FeatureSpec( key="video/crop", id_columns=["id"], deps=["video/raw"], description="Videos cropped to 720x480.", metadata={"team": "ML"}, ) 和 Dagster 一样,Metaxy 有声明式 DSL 来定义 DAG,节点代表数据,在 Metaxy 里叫 Feature。Metaxy Feature 直接映射为 Dagster Asset。 上面的 spec 可以挂到一个 feature 类上: class VideoCrop(mx.BaseFeature, spec=spec): path: str duration: float frame_count: int 然后这些 feature 定义可以轻松集成到 Dagster 资产中: import dagster as dg import metaxy as mx from metaxy.ext.dagster import metaxify @metaxify() @dg.asset(metadata={"metaxy/feature": "video/crop"}) def video_crops(store: dg.ResourceParam[mx.MetadataStore]): with mx.BufferedMetadataWriter(store) as writer: changes = store.resolve_update("video/crop") for sample in changes.new: # handle each sample here, e.g. run `ffmpeg` to crop videos # or do something else # once done, submit metadata to Metaxy # BufferedMetadataWriter flushes them automatically writer.put({"video/crop": results}) 就这么简单! 这里的 @metaxify 装饰器干了不少重活,把 Metaxy 掌握的所有信息——Metaxy 了解整个 feature/资产图——注入到原本非常简陋的 Dagster 资产中,包括正确的血缘关系、Dagster 的 code_version、表结构元数据等等。 Metaxy 还以其他几种方式与 Dagster 集成。值得一提的是 MetaxyIOManager,它支持对接任何 Metaxy 支持的元数据存储,如 DuckDB、DeltaLake、ClickHouse——甚至可以在开发和生产之间按需切换,而代码一行不用改! 由于这种集成无侵入且简单,Metaxy 相当于 Dagster 在多模态数据集上编排单个样本的扩展。 详情请看 @metaxify 文档。

结语

一直以来,在 Dagster 管道中管理单个文件都是件麻烦事。Metaxy 让这件事变简单了,Dagster 用户可以专注于数据转换本身。
Dagster 与 Metaxy 完美互补:前者负责资产级别的编排,并启动调用 Metaxy 代码的流程;后者则处理行级别、以子样本为粒度细化的编排。 Metaxy 也能脱离 Dagster 独立使用!相信未来我们会发现它的更多潜力。 阅读文档,并通过 `uv pip install metaxy[dagster]` 安装吧! 我们很高兴能帮更多用户通过 Metaxy 解决元数据管理难题,欢迎在 GitHub 上联系我们。 致谢 感谢 Georg Heiler 通过讨论和代码为项目做出贡献,以及 Metaxy 所基于的开源项目:Narwhals、Ibis 等。 有疑问或建议?欢迎在 Slack 或 GitHub 发起讨论。 对我们感兴趣,想了解更多岗位?请查看开放职位。 想获取更多类似内容?关注我们的 LinkedIn。

评论 (0)