Milvus 3.0 开源解读之backfill|亿级 AI 数据,如何做高效特征回填

2026-08-053 分钟阅读

AI 时代,被改变的不只是数据规模,还包括数据本身的性质。

传统业务系统中的数据,大多用于记录已经发生的事实:商品的名称、用户的订单、设备的状态。数据会更新,但字段的含义和生成方式通常相对稳定。写入数据库之后,一条记录往往可以长期作为业务事实被读取和使用。

但进入AI时代, 系统中的大量数据不再是静态事实,而是模型在某个时间点对原始数据做出的计算结果。

Embedding、质量分、分类标签、过滤属性和重排序特征,都会随着模型、数据集和计算逻辑的变化而更新。模型升级后,同一批原始数据往往需要重新计算,并将新的 embedding_v2、质量分或标签回填到已经积累了数亿条记录的 Collection 中。

对一个拥有数亿条数据的向量检索系统来说,这种不断更新的AI数据,会形成一种新的基础设施压力:=

  • 如何避免离线任务长期占用在线写入资源?
  • 如何保证新特征与历史记录准确对应?
  • 如何知道哪些数据已经处理、哪些仍然失败?
  • 如何在 schema 和数据持续变化时,确认计算结果仍然有效?
  • 如何让一次亿级更新具备可重跑、可校验和可发布的工程边界?

很多团队的第一反应是按主键 upsert。这在实时更新和小批量修订中没有问题。但当新特征由 Spark 离线计算、以 Parquet 形式交付,并且需要覆盖数亿乃至数十亿条历史数据时,我们需要的已经不只是更快的upsert ,而是把所有过程版本的数据,作为一个有版本、可校验、可重跑的数据产物,与不同业务需求深度融合。

也是因此,在 最新的3.0中我们在Milvus Spark Connector 中深度优化了 backfill 能力:它没有把 Spark 变成一个更大的在线写入客户端,而是让 Spark 直接在对象存储侧完成历史数据关联和目标列构建,再按照 Milvus 的 segment 边界生成可提交、可审计的结果

这个区别,决定了 backfill 与普通批量写入的根本差异。

接下来,这篇文章将从这个差异出发,介绍AI时代为什么需要backfill 、它如何工作,以及哪些边界需要在生产中提前规划

AI时代的数据变化

在AI时代的生产系统中,AI 基础设施不仅要回答如何写入新数据,还必须回答如何系统性地更新旧数据。

对应的数据更新至少可以分为两类。

第一类是在线实体变化。

用户修改了商品信息,运营调整了一个标签,权限系统更新了一条访问规则。这类变化通常数据量小、延迟要求高,需要尽快对线上查询可见。在线 upsert 很适合处理这类任务。

第二类是离线特征演进。

模型团队训练出了新的 embedding 模型,数据团队重新计算了质量分,检索团队增加了新的过滤或重排序特征。这类变化通常覆盖大量历史数据,结果已经以 Parquet 等离线格式存在,并且需要经过完整性检查后再发布。

两类任务最终都在修改数据,但它们的执行语义并不相同。

在线实体更新强调低延迟和单条记录的即时可见性;离线特征发布强调批量计算、版本一致性、可重跑和结果验收。如果把后者强行转换成数亿次前者,系统不仅会付出额外的网络和写入成本,还会失去离线数据任务本应具备的版本、边界和可观测性。

并且,这一差异会随着模型迭代速度加快而越来越明显。

假设一个商品检索 collection 已经包含以下字段:

Plaintext
item_id | title | embedding_v1 | category | price

模型团队训练出更好的 embedding,数据团队也计算出了质量分。它们的输出通常是一份 Parquet:

Plaintext
item_id | embedding_v2 | quality_score

现在需要把这两个字段补到已经在线运行的 collection 中。

如果只有几万或几十万条数据,Spark 读取 Parquet 后按批次调用 SDK,再按 item_id upsert 回 Milvus,是一种直接而有效的实现。但数据规模从百万增长到亿级、十亿级后,问题会改变。

第一,离线任务会持续占用网络和在线写入资源。

新的 embedding 已经在大数据系统中以 Parquet 形式存在,适合由 Spark 扫描、关联和并行处理。如果把它转换成海量 upsert,数据需要经过客户端切批、序列化、网络传输和在线写入流程。一个原本属于离线数据平台的任务,会长期占用在线服务的网络、写入和计算资源。

而离线任务追求吞吐,在线系统追求稳定。把两种工作负载放进同一条链路,资源竞争几乎不可避免。

第二,列变化就能解决的问题被升级成了实体更新

本次任务只增加 embedding_v2quality_score 两列。标题、类目、价格和其他历史字段并没有变化。特征计算结果在大数据系统中已经是 Parquet,逐条转发到服务端并不符合它的生产方式;而继续用面向实体的方式处理,意味着系统始终围绕更新一条商品记录组织工作,而不是围绕为历史数据构建两个新字段组织工作。

当数据规模较小时,这种错位可能只表现为效率不高;当数据规模进入亿级后,它会影响任务的执行方式、恢复方式和发布方式。

第三,请求成功无法证明数据发布成功

即使所有请求最终都返回成功,团队仍然需要回答一系列问题:

  • 输入特征中是否存在重复主键?
  • 有多少主键真正匹配到了 collection?
  • 未匹配的历史记录应该保留旧值还是被置空?
  • 任务执行期间 schema 是否发生变化?
  • 哪些 segment 已经完成,哪些需要重跑?
  • 重试是否会产生重复结果?
  • 输出是否仍然适用于当前版本的数据?

也就是说,当任务规模达到亿级时,系统真正需要的是一套数据发布协议,而不只是更高的写入并发。

任务的性质变了,执行模型也必须随之改变。

Backfill:从更新实体转向发布特征列,为什么是最适合AI时代的数据更新方式?

简单地继续扩大 Upsert解决不了AI时代数据更新的需求,对于大规模特征工程,更合理的任务划分是:

  • Spark 负责读取 Parquet、按主键关联、校验输入、并行构建目标字段
  • 对象存储 承载已有数据和新生成的列数据和任务结果
  • Milvus snapshot 提供本次任务所需的 schema、segment 布局和版本信息;
  • segment 成为回填的最小提交与审计单位,提供计算、重试、审计和提交边界

整个过程可以理解为:

Plaintext
离线特征 Parquet
        │
        │ 按主键关联
        ▼
Spark 读取 Milvus 的历史数据快照
        │
        │ 保留每条数据原本所属的 segment 和位置
        ▼
按 segment 并行构建新字段
        │
        ▼
将新列写入对象存储,并提交对应的元数据

这种分工保留了离线数据原本的生产方式。模型和特征团队可以继续以 Parquet 交付结果,数据平台继续用 Spark 调度大规模任务。在线 Milvus 也不必承担全部离线扫描和关联压力。

更重要的是,新特征不再以大量独立请求的形式进入系统,而是以一批有明确版本和物理边界的数据产物的形式进入发布流程。

但是,要让离线计算结果真正成为可发布的产物,首先必须解决一个问题:一个运行数小时的任务,应该以哪个版本的数据为准?

这就是 Snapshot 的作用。

Snapshot:让数据不断变化的离线任务,有一个确定的版本

长时间运行的大规模离线任务最怕输入状态持续变化。

如果 Spark 在任务执行过程中不断从在线服务获取 schema、segment 列表和存储位置,那么一个运行数小时的作业可能在不同阶段看到不同的 collection 状态。最终结果即使计算成功,也未必还能安全提交。

snapshot 在这个过程的作用,是为任务提供一个确定的版本标准。它通常包含:

  • collection schema;
  • 目标字段的类型和字段标识;
  • 主键字段;
  • 本次需要处理的 segment;
  • segment 对应的存储信息;
  • schema 版本。

有了 snapshot,Spark 可以基于一个确定版本读取对象存储中的历史数据,不需要在大规模扫描过程中持续依赖在线 Milvus。它带来了三个实际收益:

计算可复现: 任务失败后,可以基于同一份 snapshot 重跑,而不是重新面对一个已经变化的 collection。

在线服务与离线计算解耦: 大规模扫描、关联和列构建主要发生在 Spark 与对象存储之间,不需要把全部压力传递给在线检索和写入链路。

陈旧结果可以被识别: backfill 结果会携带计算时使用的 schema 版本。提交流程可以比较当前版本和任务版本,拒绝已经过期的产物。

需要明确的是,snapshot 不等于全局事务。它固定的是计算起点,并不会自动协调 schema 变更、compaction、删除和其他并发写入。计算与提交之间的版本检查,以及不同数据任务之间的顺序,仍然需要由上层工作流管理。

Snapshot 解决了基于哪个版本计算的问题。

接下来还需要解决另一个问题:如此庞大的任务,应该以什么作为计算、失败恢复和提交的边界?

如何设计任务的计算边界?segment是关键

按主键完成关联,只解决了“哪条历史记录应该获得哪个新特征”的问题。接下来,我们要考虑数据要如何写回。

如果把整个 collection 看作一张大表,最直观的做法可能是让 Spark 任意切分数据,然后统一写出。但 Milvus 的物理数据本身按 segment 组织,回填必须尊重这个边界。

因此,Backfill 的设计中,我们选择了以 Segment 作为核心执行边界。在一个Backfill Application中,一个 Spark Task通常负责对应处理一个 Segment:

  • 该任务读取属于这个 segment 的历史行;
  • 将新特征按原始位置排序;
  • 批量写出这个 segment 所需的新列;
  • 记录行数、匹配数、输出位置和版本信息。

这样设计的好处是,在任务结束后,我们不只知道“作业成功了”,还可以知道:

  • 一共处理了多少个 segment;
  • 每个 segment 写出了多少行;
  • 特征文件中的多少主键真正匹配到了 collection;
  • 哪些 segment 失败,需要有针对性地重跑;
  • 新产物对应哪个 schema 版本。

对亿级任务而言,这种可观测性往往比单纯提高吞吐更重要。

Snapshot 固定了版本,Segment 固定了执行边界。接下来,Backfill 还需要明确:在每个 Segment 内,究竟要重建哪些数据?

回填时,如何保证只重建发生变化的字段?

在典型的AI数据任务中,如果只需要增加 embedding_v2quality_score,就没有理由重新写入标题、价格、类目和其他未发生变化的字段。

因此,Milvus 3.0 中的 Backfill 会利用 Storage 层提供的抽象,对指定 Segment 进行按列更新。该能力不依赖具体的底层文件格式,也不需要重写整个 Segment。

Backfill 任务只读取本次计算所需的数据,例如主键、Segment 与行位置信息、必要的旧字段以及参与计算的输入字段;计算完成后,再通过 Storage 接口写入本次新增或发生变化的目标列。其他未发生变化的列仍沿用原有数据,不会被重复构建。

例如,将 embedding_v1 切换为 embedding_v2 时,Backfill 只需要生成并写入 embedding_v2;如果本次任务还会同时计算 quality_score,则一并更新这两个目标字段,而不需要重新写入标题、价格、类目等无关字段。

过程中,对于向量字段,连接器也会在进入写入阶段之前完成常见输入形态的校验与规范化,例如向量维度、稀疏向量索引和数值范围。这样,Spark 作业可以直接消费常见的数组、二进制或结构化输出,而不需要让每个上游特征任务理解底层存储编码。

通过这种方式,Embedding 升级类任务可以被组织为针对目标 Segment 的增量列更新,而不是一次 Collection 级的全量迁移。

它也同样适用于批量补齐新字段;重新计算质量分、标签和过滤属性;迁移到新版本 embedding;修正一批历史特征;以离线结果为准,对目标列做批量更新等任务中,也具备相当不错的表现。

到这里,我们已经明确了基于哪个版本计算、以什么为边界执行以及具体写哪些列。

但还有一个数据语义问题必须提前定义:回填时,缺失数据应该怎么处理?

回填时,缺失数据应该怎么处理?

现实中的特征数据很少能够完美覆盖整个 collection。比如,某些商品可能缺少文本,无法生成新 embedding;某些记录可能被特征任务过滤;也可能只有一部分历史数据需要修正。

因此,backfill 必须明确一个关键问题:

当输入文件中没有某个主键时,目标字段应该发生什么?

针对这种问题,我们当前支持三种模式:

模式含义常见用途
coalesce已有值优先;只有旧值为空时才用新数据补齐补洞、渐进式修复
overwrite只要 Parquet 中存在该主键,就以新值覆盖旧值修正指定范围内的数据
replace以 Parquet 为准;未匹配到的数据也会被置空已确认输入完整的全量重建

默认的 coalesce 相对保守,适合首次落地时降低误覆盖风险。overwrite 表示新特征对已匹配主键具有权威性,但不会影响输入范围之外的数据。replace 则需要格外谨慎:它适用于输入数据确实覆盖完整、且新数据就是唯一权威来源的场景。

但无论选择哪种模式,我们输入 Parquet 中的主键必须唯一。重复主键会让 join 结果膨胀,使同一条历史记录产生多份输出。当前 connector 在检测到重复 join key 时会立即 fail fast 并报错,避免任务继续运行并将数据问题带到发布阶段。

不过,计算完成仍然不等于数据发布完成。生产环境还需要一套从输入冻结、任务执行到结果验收的完整流程。

生产环境中,完成backfill 的几个注意事项

Backfill 是一个物理列构建过程,不是对在线 Collection 的万能事务。

为了让新列与旧列保持行位置一致,backfill 在读取时需要保留原始的物理顺序。即使某些记录在逻辑上已经被删除,回填阶段也不能简单把它们过滤掉;否则新列的行号会整体偏移,反而破坏数据对齐。

因此,生产上建议遵循几个原则:

  1. 先完成 schema 变更。 待回填字段应已存在于任务所用 snapshot 的 schema 中。
  2. 再固定输入。 使用 snapshot 作为回填任务的起点。
  3. 把版本纳入提交检查。 回填结果带有 schema 版本,提交方应拒绝过期结果。
  4. 显式安排并发窗口。 对同一个 collection,schema 变更、compaction、大量删除和 backfill 不应被当作互不相关的任务;需要由工作流定义它们的先后关系。
  5. 先从小范围验证。 先选一个 partition 或少量 segment,检查匹配率、空值比例和新字段质量,再扩大范围。

这些原则可以让一次重计算从“长时间运行的大脚本”变成可验证的数据发布流程,后续我们还会将schema 变更与 snapshot 合并精简,敬请期待。

一个典型的作业形态

从使用角度看,一个 Backfill 作业主要包含两类输入:一份特征 Parquet 和一份 Milvus snapshot。

作业既可以通过 spark-submit 独立运行,也可以嵌入现有的 Spark 工作流。

spark-submit \
  --class com.zilliz.spark.connector.operations.backfill.BackfillApp \
  spark-connector-assembly-<version>.jar \
  --parquet s3a://feature-bucket/features/embedding-v2.parquet \
  --snapshot s3a://<snapshot_metadata_path> \
  --s3-endpoint s3.us-west-2.amazonaws.com \
  --s3-bucket milvus-bucket \
  --use-iam \
  --mode overwrite \
  --output-result s3a://milvus-bucket/backfill/embedding-v2-result.json

特征数据与 Milvus 数据可以位于不同的对象存储 bucket,并使用不同的访问配置。对数据平台来说,这一点很重要:特征团队、训练平台和向量库未必共享同一个存储账户或 region。

运行后,建议把 result JSON 作为工作流的一部分保存和检查,而不是只依赖 Spark 作业的退出码。至少应关注:处理的 segment 数量、源数据与特征数据的匹配率、每个 segment 的写入行数、输出版本,以及任务使用的 schema 版本。

结语

长期来看,向量库中的字段演进会越来越常见:embedding 会迭代,标签会重算,质量特征会持续改善,新的过滤和排序信号也会不断加入。

如果每一次变化都依赖海量在线 upsert,离线计算和在线服务会彼此牵制。基于 Spark 的 backfill 则提供了一条更贴近大数据工作方式的路径:让离线任务在对象存储上完成关联与计算,只重建需要变化的列,并按 Milvus 的 segment 边界留下可审计、可校验的结果。

Tips 何时选 Backfill,何时继续用 Upsert?

最简单的判断标准是看任务的本质:

场景更合适的选择
用户刚修改了一条实体,需要很快可见在线 upsert
少量记录的实时修订在线 upsert
离线模型为数亿条数据生成了新 embeddingSpark backfill
批量补齐标签、质量分、过滤字段Spark backfill
需要可重跑、按 segment 审计和版本校验Spark backfill
需要改变主键、重排数据或整体迁移模型专门的数据迁移/重导入方案

可以把两者看成互补关系:

  • Upsert 面向在线世界,强调一条实体现在就要变;
  • Backfill 面向离线世界,强调一个新特征要系统性地覆盖历史数据。

当数据规模进入大数据量级,最有价值的并不是把 upsert 调得更快,而是承认这已经是一个数据工程问题,并为它选择合适的计算引擎、存储路径和提交边界。

AI Assistant