如何在 Milvus 外部表文本字段上生成 BM25 稀疏向量派生字段以支持关键词检索
【免费下载链接】milvusMilvus is a high-performance, cloud-native vector database built for scalable vector ANN search项目地址: https://gitcode.com/GitHub_Trending/mi/milvus
你有一批 Parquet 文本文件放在对象存储(如 MinIO/S3)上,不想把原始数据复制进 Milvus,但希望这些文本能像普通集合一样做关键词检索。Milvus 的外部表(External Collection)支持在 schema 里声明一个 BM25 函数:原始文本列保持外部引用,Milvus 只在 refresh 时计算并物化出一个SPARSE_FLOAT_VECTOR派生字段(sparse vector output),刷新后即可为该字段建稀疏倒排索引,并用原始文本直接发起 BM25 检索。
本文按仓库端到端测试 TestExternalTableBM25Function 与 外部表函数输出设计文档 的操作路径展开,使用 Go 客户端完成:建表 → 触发 refresh → 验证行数 → 建索引 → load → 用文本搜索。
准备条件
- 一个可连接 Milvus 的客户端环境(下文命令基于仓库内 Go 客户端 client/milvusclient/external_table.go);
- 对象存储中已有若干 Parquet 文件,且列名与下文 schema 的
external_field映射一致。测试数据是 10 个文件、每个文件 3000 行,列包含id(int64)和text(string); external_source填对象存储 URI,external_spec至少声明格式,测试与文档示例使用{"format": "parquet"}。CI 中外部数据源 scheme 的白名单目前只有 MinIO(见 test_milvus_client_external_table.py 中ALLOWED_EXTERNAL_SOURCE_SCHEMES)。
下面的示例中,ext_bm25_docs是集合名,minio://127.0.0.1:9000/my-bucket/ext-data/是示例 external_source,替换为你自己的集合名与 Parquet 所在的前缀即可。
第一步:创建外部集合并声明 BM25 函数输出字段
外部集合的 schema 规则(见 外部表设计文档 与函数输出设计文档):
- 普通用户字段必须设置
external_field,把 Milvus 字段名绑定到外部源中的物理列名; - 函数输出字段不能设置
external_field,它由 Milvus 在 refresh 时生成; - 外部集合由系统注入虚拟主键(
__virtual_pk__),不支持用户自定义主键、autoID、dynamic field、partition key。
schema := entity.NewSchema(). WithName("ext_bm25_docs"). WithExternalSource("minio://127.0.0.1:9000/my-bucket/ext-data/"). WithExternalSpec(`{"format": "parquet"}`). WithField(entity.NewField().WithName("id").WithDataType(entity.FieldTypeInt64). WithExternalField("id")). WithField(entity.NewField().WithName("text").WithDataType(entity.FieldTypeVarChar). WithMaxLength(1024).WithExternalField("text"). WithEnableAnalyzer(true).WithAnalyzerParams(map[string]any{"tokenizer": "standard"})). WithField(entity.NewField().WithName("sparse_vec").WithDataType(entity.FieldTypeSparseVector)). WithFunction(entity.NewFunction().WithName("bm25_fn"). WithInputFields("text").WithOutputFields("sparse_vec").WithType(entity.FunctionTypeBM25)) err := mc.CreateCollection(ctx, client.NewCreateCollectionOption("ext_bm25_docs", schema))CreateCollection时 Proxy 会先执行函数校验,把sparse_vec标记为函数输出字段,再做外部 schema 校验,因此输出字段不会被误判为缺少external_field的普通字段。创建后用DescribeCollection确认 schema 中Functions有 1 项、ExternalSource与所填一致(测试中的校验方式)。
第二步:手动触发 refresh 并等待完成
外部集合的数据只在 refresh 时同步,且没有自动检测源数据变化的机制。调用RefreshExternalCollection获得 job ID,然后轮询GetRefreshExternalCollectionProgress,直到状态变为RefreshStateCompleted;如果变为RefreshStateFailed,进度信息中的Reason字段会给出失败原因:
refreshResult, err := mc.RefreshExternalCollection(ctx, client.NewRefreshExternalCollectionOption("ext_bm25_docs")) // refreshResult.JobID 用于跟踪进度 for { p, err := mc.GetRefreshExternalCollectionProgress(ctx, client.NewGetRefreshExternalCollectionProgressOption(refreshResult.JobID)) if p.State == entity.RefreshStateCompleted { break // 刷新完成 } if p.State == entity.RefreshStateFailed { // 失败原因在 p.Reason 中 break } time.Sleep(3 * time.Second) }refresh 过程中,DataNode 会扫描外部源的分片,对需要新建/重建的 segment 流式读取外部输入列、执行 BM25 函数、把输出列写入 StorageV3 列组,并持久化 BM25 stats(manifest 中bm25.<fieldID>键)。若集合之前建过、旧 segment 的 manifest 里缺少函数输出列,这些 segment 会被判定失效并重建。
第三步:验证刷新后的行数
刷新完成后用GetCollectionStats核对行数,应与上传的 Parquet 总行数一致(测试中 10 文件 × 3000 行 = 30000):
stats, err := mc.GetCollectionStats(ctx, client.NewGetCollectionStatsOption("ext_bm25_docs")) rc, _ := strconv.ParseInt(stats["row_count"], 10, 64) // rc 应等于外部 Parquet 的总行数第四步:为 sparse 字段建稀疏倒排索引并加载集合
BM25 输出字段在 refresh 之后就是普通 schema 字段,可以为它建稀疏倒排索引(测试使用index.NewSparseInvertedIndex(entity.BM25, 0.1),其中0.1是测试传入的索引构建参数),然后加载集合:
sparseIdx := index.NewSparseInvertedIndex(entity.BM25, 0.1) idxTask, err := mc.CreateIndex(ctx, client.NewCreateIndexOption("ext_bm25_docs", "sparse_vec", sparseIdx)) err = idxTask.Await(ctx) loadTask, err := mc.LoadCollection(ctx, client.NewLoadCollectionOption("ext_bm25_docs")) err = loadTask.Await(ctx)第五步:用原始文本做 BM25 关键词检索
搜索时直接把原始文本作为查询向量传给sparse_vec字段,无需自己计算向量:
res, err := mc.Search(ctx, client.NewSearchOption("ext_bm25_docs", 10, []entity.Vector{entity.Text("machine learning")}). WithANNSField("sparse_vec"). WithOutputFields("id")) // 结果命中行、分数在 res[0](hits.Len()、hits.Scores)测试语料是 10 个不同句子循环重复 30000 行,因此查询"machine learning"时 top 结果的id % 10全部等于 2(该行对应的句子类别),这是该测试语料特有的校验方式(文档示例),不要把它当成固定分数阈值;换成自己的数据后,验证点应改为:top 结果应集中在文本与查询词匹配的行上。
PyMilvus 的等价写法(可选分支)
设计文档给出的 PyMilvus 侧 schema 声明方式:
schema.add_field("doc", DataType.VARCHAR, external_field="doc_text") schema.add_field("sparse", DataType.SPARSE_FLOAT_VECTOR) schema.add_function( name="bm25_fn", input_fields=["doc"], output_fields=["sparse"], function_type=FunctionType.BM25, )Python 测试中 refresh 的调用形态是refresh_external_collection(client, collection_name=coll)返回 job id,随后用get_refresh_external_collection_progress轮询,终态字符串为RefreshCompleted/RefreshFailed(失败时看reason字段)。
限制与常见问题
- BM25 输出字段不能直接取出:
CanRetrieveRawFieldData对 BM25 输出返回 false,检索结果里不要请求sparse_vec字段,只输出原始列(如id);MinHash/TextEmbedding 等非 BM25 输出字段则可以取出。 - 外部集合是只读的:insert/delete/upsert/import/flush 均被拒绝(错误信息形如 "insert operation is not supported for external collection");数据变更必须在源头做,然后再次手动 refresh。
- refresh 覆盖源变更场景:向 Parquet 前缀新增文件后,重新调用 refresh 才会把新分片组织进 segment;未变化的 segment 会被保留。
- schema 限制:外部集合不支持 dynamic field、用户主键、autoID、partition key、clustering key、struct 字段;varchar 字段声明
external_field时若映射为空会报field 'xxx' ... must have external_field mapping。 - text_match 是独立开关:在 varchar 字段上加
enable_match=True可开启text_match(...)过滤,其文本索引 stats 由独立的后台任务生成,未就绪时 load 会以TextIndexNotFound失败并可重试。本文的 BM25 检索路径不依赖该开关。
参考文档
- External Table Function Output 设计文档:函数输出字段、refresh 时机、BM25 stats 与 manifest 规则、行为矩阵
- External Table 设计文档:外部集合 API、refresh job 生命周期、配置参数
- 端到端测试:本文 Go 步骤的完整可运行来源
- PyMilvus 外部表测试:Python 客户端的 schema 校验与 refresh 用法
【免费下载链接】milvusMilvus is a high-performance, cloud-native vector database built for scalable vector ANN search项目地址: https://gitcode.com/GitHub_Trending/mi/milvus
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考