AB
AiBoss站
教程

MiniMax

教程

MiniMax 与 Zilliz Cloud 向量检索实践:实时推荐与海量训练数据去重

介绍如何借助 Zilliz Cloud 承载两类差异极大的 AI 工作负载:面向实时推荐的向量相似检索,以及面向 LLM 训练数据的兆级去重。内容涵盖前置准备、集合与索引设计、MinHash + LSH 去重流程、完整示例与常见限制,适合正在搭建推荐或数据清洗管线的工程师参考。

当 AI 应用从原型走向规模化,向量检索往往会从「一个功能」变成「一条基础设施」。MiniMax 的公开实践说明了这一点:同一套向量底座既要支撑面向数千万月活用户的实时推荐,又要在离线侧处理数十万亿 token 级别的训练数据去重。这两类负载的访问模式、延迟要求和成本结构几乎相反,却可以收敛到同一个检索层上。本文以 MiniMax 的公开案例为线索,整理出一套可复用的做法,适合正在搭建推荐系统、语义搜索或 LLM 数据预处理管线的工程师阅读。文中涉及的产品能力、配额与可用性均可能变化,请以官网当前信息为准。

准备工作

在动手之前,需要先把环境与数据侧的几件事确认清楚。以下条目按依赖顺序排列,缺一项都会让后面的步骤卡住。

账号与集群

  • 一个可用的向量数据库云服务账号,并已创建至少一个集群(Cluster)。集群是计算与存储的载体,后续的 Collection、Index、Replica 都挂在它下面。
  • 确认所选区域(Region)与你的数据源、在线服务处于同一网络可达范围内。跨区域调用会显著抬高延迟,对实时推荐这类场景尤其致命。
  • 如果计划使用按需型(On-Demand)计算来跑离线批处理,需要确认当前账号所属的项目类型与区域是否在支持范围内。这类能力在部分阶段属于预览状态,存在项目类型与区域限制。

凭据与连接信息

  • 集群的访问地址(Endpoint)与 API Key。API Key 通常只在创建时完整展示一次,务必先落盘保存。
  • 如果使用 SDK 接入,准备好对应语言的客户端库。Python 与 Java 是最常见的两种,选择与你的服务栈一致的那一个即可。
  • 确认网络白名单或私有连接配置。生产环境不建议直接暴露公网访问。

数据与模型侧

  • 确定向量维度。维度由你的 Embedding 模型决定,一旦建好 Collection 就不能随意更改。案例中在线推荐使用的是 32 维向量,维度较低意味着单条记录更小、内存占用更低,但语义表达能力也相应受限,需要根据业务权衡。
  • 准备好 Embedding 生成链路。向量检索本身不负责把文本、图片或音频转成向量,这一步必须在上游完成。
  • 对离线去重场景,先明确数据规模量级:是 GB 级、TB 级还是 PB 级。规模直接决定你需要单机跑还是必须分布式水平扩展。

容量与成本预估

  • 在线侧要估算峰值 QPS、可接受的 P99 延迟、以及需要保留的向量总量。案例中的目标是峰值超过 5,000 QPS 时仍保持 30ms 以内的响应。
  • 离线侧要估算单次去重任务的文档数、可接受的运行时长,以及重复执行频率。
  • 计算资源通常按计算单元(Compute Unit)计量,副本(Replica)会成倍放大资源占用。上线前先做小规模压测,再按结果扩容,不要凭感觉配置。

操作步骤

第一步:为在线检索建立 Collection

Collection 是向量数据的逻辑容器,类似关系型数据库里的表。创建时需要指定向量字段的维度与主键字段。

from pymilvus import MilvusClient

client = MilvusClient(
    uri="https://your-cluster-endpoint",
    token="your-api-key",
)

client.create_collection(
    collection_name="talkie_recommend",
    dimension=32,
    metric_type="COSINE",
    auto_id=False,
)

几个关键点:

  • dimension 必须与上游 Embedding 输出严格一致,写错会在插入时报错。
  • metric_type 决定相似度计算方式。余弦距离适合归一化后的语义向量;欧氏距离适合本身带量纲的特征向量。选错指标不会报错,但召回结果会明显变差。
  • 如果业务侧自己生成主键(例如内容 ID),把 auto_id 设为 False,这样后续可以用 Upsert 做增量更新。

第二步:选择索引类型

索引决定了检索速度与召回率之间的平衡。对于大多数在线场景,使用自动索引即可,它会根据数据规模与维度自动选择合适的分层策略,省去手工调参。

index_params = client.prepare_index_params()

index_params.add_index(
    field_name="vector",
    index_type="AUTOINDEX",
    metric_type="COSINE",
)

client.create_index(
    collection_name="talkie_recommend",
    index_params=index_params,
)

如果对召回率有更严格的要求,也可以选用基于图结构的近似最近邻索引,但需要额外调整构建参数与查询参数。参数调优属于经验性工作,建议用真实数据做离线评估后再定。

第三步:写入与增量更新

向量数据不是一次性导入就完事。推荐场景里内容库持续变化,必须有稳定的增量写入路径。

rows = [
    {"id": 1001, "vector": [0.01, 0.23, ...], "category": "voice_pack"},
    {"id": 1002, "vector": [0.44, 0.07, ...], "category": "message"},
]

client.insert(collection_name="talkie_recommend", data=rows)

# 内容更新时用 upsert 覆盖同主键记录
client.upsert(collection_name="talkie_recommend", data=updated_rows)

写入后数据不会立刻可查,需要等待索引构建完成。可以显式刷新,也可以依赖服务端的一致性级别设置。对一致性要求高的场景,查询时指定更强的一致性等级;对延迟敏感且能容忍短暂延迟的场景,可以放宽这一设置以换取吞吐。

第四步:带过滤条件的相似检索

纯向量检索往往不够用。真实业务通常需要「在某个类目下找最相似的内容」这类组合条件。向量检索与标量过滤可以写在同一次查询里。

results = client.search(
    collection_name="talkie_recommend",
    data=[query_vector],
    limit=20,
    filter='category == "voice_pack"',
    output_fields=["category"],
)

for hits in results:
    for hit in hits:
        print(hit["id"], hit["distance"], hit["entity"]["category"])

过滤字段可以是标量、数组或 JSON 结构。需要注意的是,过滤条件越严格,实际参与相似度计算的数据越少,召回率与延迟都会随之变化。上线前应针对典型过滤组合分别压测。

第五步:用副本支撑高并发

单副本能扛住的 QPS 有限。当峰值流量上来后,横向增加副本是最直接的手段。案例中的在线推荐配置为 8 个计算单元配合 7 个副本,以此在流量尖峰时保持稳定。

副本的作用是分摊查询压力,同时提升可用性:单个副本异常时,其余副本仍可继续服务。代价是资源占用成倍增加,因此副本数应结合压测结果逐步上调,而不是一次性拉满。

第六步:为离线去重准备数据

训练数据去重与在线检索是两种完全不同的负载。它的特点是数据量极大、对单次延迟不敏感、但总运行时长和成本极其敏感。传统做法是用 MapReduce 类框架做全量比对,一个数据集跑几周甚至几个月是常态,这会直接拖慢模型迭代节奏。

更可行的思路是把去重也变成一次相似度检索:先用 MinHash 把每篇文档压缩成紧凑的签名,再用局部敏感哈希(LSH)把可能相似的文档分到同一个桶里,只对桶内候选做精细比对。这样避免了全量两两比较。

from datasketch import MinHash, MinHashLSH

lsh = MinHashLSH(threshold=0.8, num_perm=128)

for doc_id, text in documents:
    m = MinHash(num_perm=128)
    for token in tokenize(text):
        m.update(token.encode("utf8"))
    lsh.insert(doc_id, m)

# 查询某篇文档的近重复候选
query = MinHash(num_perm=128)
for token in tokenize(target_text):
    query.update(token.encode("utf8"))

near_duplicates = lsh.query(query)

关键参数说明:

  • num_perm 是签名长度。值越大,相似度估计越准确,但内存占用与计算量同步上升。128 是常见的起点。
  • threshold 是判定为近重复的相似度阈值。调低会召回更多候选,同时引入更多误判;调高则可能漏掉真正的重复内容。
  • 分词粒度直接影响效果。字符级 n-gram 对语言差异不敏感,适合多语种混合语料;词级切分对中文需要额外分词处理。

第七步:把去重引擎放进检索系统

如果单独维护一套去重服务,就要额外承担部署、扩缩容和数据同步的成本。更简洁的做法是把 MinHash + LSH 的计算放进向量数据库的索引体系里,与向量插入、索引构建、近似查询共用同一套 API。这样数据量增长时,去重能力可以跟着水平扩展,而不需要另起一套管线。

案例中的做法正是如此:没有新建独立的去重服务,而是让去重引擎运行在索引系统内部,复用同一批接口。结果是处理速度相比原有 MapReduce 方案提升约一倍,成本降到原来的三分之一到五分之一,十亿文档规模的去重从「偶尔跑一次的大工程」变成了可以日常执行的操作。

第八步:扩展到多模态与更多数据操作

同一套向量底座可以复用到更多场景。图片向量、音频向量、文本向量在结构上没有本质区别,只是维度与 Embedding 模型不同。把去重能力扩展到其他数据预处理环节,也能加快数据集迭代速度。

离线侧除了去重,还有几类常见的相似度操作值得一并考虑:

  • 数据集之间的近似重复检查,避免训练集与评测集泄漏。
  • 按领域或质量条件筛选样本,构造候选数据集。
  • 检索与评测失败样本相近的样本,用于定向补充数据。
  • 在数据湖上直接做探索性检索,而不必先把数据整体搬进数据库。

这些操作的共同点是:从大量数据中找出「相似的」或「符合条件的」记录,再交给下游处理。它们与在线推荐共享同一套检索原语。

一个完整示例

下面把在线检索与离线去重串成一条最小可运行的链路。假设场景是:为内容平台构建推荐候选召回,同时清理训练语料中的重复文档。

阶段一:初始化在线集合

from pymilvus import MilvusClient

client = MilvusClient(
    uri="https://your-cluster-endpoint",
    token="your-api-key",
)

client.create_collection(
    collection_name="content_recall",
    dimension=32,
    metric_type="COSINE",
    auto_id=False,
)

index_params = client.prepare_index_params()
index_params.add_index(
    field_name="vector",
    index_type="AUTOINDEX",
    metric_type="COSINE",
)
client.create_index(
    collection_name="content_recall",
    index_params=index_params,
)

阶段二:批量导入内容向量

import random

rows = []
for content_id in range(1, 10001):
    rows.append({
        "id": content_id,
        "vector": [random.random() for _ in range(32)],
        "category": random.choice(["voice_pack", "message", "starter"]),
    })

client.insert(collection_name="content_recall", data=rows)
client.flush(collection_name="content_recall")

阶段三:执行一次带过滤的召回

query_vector = [random.random() for _ in range(32)]

results = client.search(
    collection_name="content_recall",
    data=[query_vector],
    limit=10,
    filter='category == "voice_pack"',
    output_fields=["category"],
)

for hits in results:
    for hit in hits:
        print(f"id={hit['id']} score={hit['distance']:.4f}")

阶段四:对训练语料做近重复检测

from datasketch import MinHash, MinHashLSH

lsh = MinHashLSH(threshold=0.8, num_perm=128)
seen = {}

for doc_id, text in corpus:
    m = MinHash(num_perm=128)
    for token in text.split():
        m.update(token.encode("utf8"))

    duplicates = lsh.query(m)
    if duplicates:
        seen[doc_id] = duplicates
    else:
        lsh.insert(doc_id, m)

print(f"发现 {len(seen)} 篇近重复文档")

阶段五:清理并输出干净语料

drop_ids = set()
for doc_id, dup_ids in seen.items():
    drop_ids.add(doc_id)
    drop_ids.update(dup_ids)

clean_corpus = [
    (doc_id, text)
    for doc_id, text in corpus
    if doc_id not in drop_ids
]

print(f"原始 {len(corpus)} 篇,清理后 {len(clean_corpus)} 篇")

这条链路跑通后,把在线部分的副本数按压测结果上调,把离线部分的阈值与签名长度按语料特征调优,就能覆盖大部分实际需求。

注意事项

关于资源与配额

  • 计算资源按计算单元计量,副本会成倍放大占用。扩容前先压测,确认瓶颈在查询侧还是索引构建侧。
  • 按需型计算能力在部分阶段属于预览状态,存在项目类型与区域限制,使用前需确认当前账号是否符合条件。
  • 部分能力对单次查询可返回的记录数有上限,超出上限的候选生成需求需要另行设计分批策略。
  • 价格、配额与区域可用性会随时间调整,务必以官网当前信息为准。

关于性能设计

  • 不要只盯着近似最近邻算法本身的耗时。生产环境的延迟由写入、更新、过滤条件、召回率、峰值 QPS、副本数、索引与查询资源共同决定。
  • 过滤条件与向量检索组合使用时,过滤的选择性会显著影响实际扫描量。高选择性过滤可能让延迟不降反升,需要实测。
  • 索引构建是资源密集型操作。大批量导入时,索引构建与在线查询会争抢资源,建议错峰执行或做资源隔离。

关于去重质量

  • MinHash 的签名长度与相似度阈值需要针对语料调优。阈值过低会误删内容,过高则漏删重复。
  • 分词方式对多语种语料影响很大。字符级 n-gram 通常比词级切分更稳健。
  • 去重不只是为了省存储。冗余数据会导致模型过拟合、泛化能力下降,因此去重质量直接影响训练效果。
  • 数据规模上升到十亿文档量级后,签名存储与候选匹配会成为内存瓶颈,需要在算法层面做针对性优化。

关于数据管线演进

  • 训练数据处理已经不是一次性的预处理,而是分布在预训练、中期训练、监督微调、偏好优化等多个阶段反复执行的操作。每个阶段都需要生成、更新和复用不同的数据集。
  • 随着数据资产增多,「用哪些数据」「某个数据集包含什么」「如何抽取特定条件的样本」这类问题会越来越复杂,相似度检索在这些环节都能派上用场。
  • 把去重、检索、筛选统一到同一套基础设施上,可以减少维护多套管线带来的同步成本与运维复杂度。