首页
/ Pathway LLM xpack 文档索引实战:基于 DocumentStore 与多类型检索器的实时文档检索

Pathway LLM xpack 文档索引实战:基于 DocumentStore 与多类型检索器的实时文档检索

2026-09-04 14:22:26作者:侯霆垣

本篇技术文章围绕 Pathway(Python ETL 流处理框架)LLM xpack 的文档索引(Document Indexing)能力展开,系统讲解向量索引、BM25 全文索引与混合索引的选型与参数配置,并完整演示如何基于 DocumentStore 搭建"解析 → 切分 → 索引 → 查询"的实时检索流水线,最后通过 REST 服务把检索能力暴露为 API。读完本文,你将能够独立完成 Pathway 中文档索引的构建、查询、元数据过滤与对外服务化部署。

什么是文档索引

文档索引(Document Indexing)通过组织与分类文档来支持高效搜索与检索:为文档内容构建一个索引(index)——即内容经过结构化表示后的产物——即可基于查询快速定位相关信息。在大语言模型(LLM)场景中,索引的价值在于组织一个可实时更新的知识库仓库,让 LLM 生成的回答有依据、可溯源。

Pathway 将文档索引划分为两大类别:

  • 基于向量的索引(Vector-based Indexing):使用 embedding 将文档表示为数值向量,通过相似度(最近邻)搜索进行检索;
  • 非向量索引(Non-Vector Indexing):基于传统文本检索方法(如 BM25),无需 embedding,适合精确关键词匹配与全文搜索。

两类索引均可流式构建:源文件新增、变更或删除时,索引会随之增量更新,这正是基于流处理框架(而非一次性批量索引)的核心优势。

Embedding:向量索引的前置条件

Embedding 将文本转换为固定长度的向量,供索引与检索使用。只有使用向量索引(如近似最近邻 ANN 搜索)时才需要 embedding。Pathway LLM xpack 在 embedders 模块 中提供了几种现成的 embedding 模型封装:

  • OpenAIEmbedder:调用 OpenAI 的 embedding API;
  • LiteLLMEmbedder:通过 LiteLLM 接入多提供商模型;
  • GeminiEmbedder:使用 Google Gemini 的 embedding;
  • SentenceTransformerEmbedder:本地运行 Sentence Transformers 模型,无需外部 API。

所有 *Embedder 均继承自 BaseEmbedder(见 embedders.py L77),并实现了 get_embedding_dimension() 方法,这一点在检索器工厂的维度推导中会再次用到(见下文)。

非向量索引:Tantivy BM25

非向量索引基于传统全文检索方法(如 BM25),对精确关键词匹配场景非常友好。Pathway 提供了 TantivyBM25Factory(基于 Rust 的 tantivy 全文检索引擎),其两个关键参数在源码中有明确的默认值:

参数 默认值 含义
ram_budget 50 * 1024 * 1024(50 MB) 索引内存预算上限(字节)。达到上限时索引会将数据块转移到磁盘存储,预算越大索引操作越快,但内存成本越高
in_memory_index True 是否将整个索引放在 RAM 中;为 False 时索引存储在 Pathway 默认的磁盘存储中

BM25 索引的数据列与查询列都必须是字符串类型(源码中的 check_default_bm25_column_types 会强制校验,见 bm25.py L16-L37),这也印证了它完全绕开了 embedding 链路。

检索器(Retriever)家族与参数详解

检索器负责创建并管理索引、高效定位相关文档。Pathway 在 pathway.stdlib.indexing 中提供三类检索器工厂:

  • 向量检索BruteForceKnnFactoryUsearchKnnFactory
  • 非向量检索TantivyBM25Factory
  • 混合检索HybridIndexFactory

BruteForceKnnFactory:精确暴力最近邻

BruteForceKnnFactory 对查询向量与索引中所有向量逐一计算距离,结果精确(非近似),适合中小规模语料或对精度要求极高的场景。其默认参数:

参数 默认值 说明
reserved_space 400 索引初始容量(条目数)
auxiliary_space 1024 * 128 评估查询时内存中可同时保留的距离值上限;若小于当前索引条目数,实际仍与索引规模成比例
metric BruteForceKnnMetricKind.COS 距离度量,默认为余弦相似度
embedder 对文本计算 embedding 的 UDF,文本索引必需
dimensions 由 embedder 推导 向量维度

关于 dimensions 有一个值得注意的实现细节:工厂的 __post_init__ 会自动推导维度(nearest_neighbors.py L422-L428)——若提供了 embedder 而未显式给出 dimensions,框架会调用 embedder.get_embedding_dimension()(对 BaseEmbedder)或用一个单点文本试算向量长度(对任意 pw.UDF)自动确定维度;两者都未提供则抛出 ValueError。因此下文的文档示例只需传入 embedder 即可。

构建 BruteForceKnnFactory 的最小示例(这是 DocumentStore 的关键组件):

from pathway.stdlib.indexing.nearest_neighbors import BruteForceKnnFactory
from pathway.xpacks.llm.embedders import OpenAIEmbedder
import os

embedder = OpenAIEmbedder(api_key=os.environ["OPENAI_API_KEY"])
retriever_factory = BruteForceKnnFactory(
    embedder=embedder,
)

UsearchKnnFactory:基于 HNSW 的近似最近邻

UsearchKnnFactory 基于 USearch 库实现 HNSW(Hierarchical Navigable Small World)算法的 k 近邻索引,适合更大规模语料下换取更快的查询速度(代价是近似而非精确)。其参数除 dimensions/embedder 外:

参数 默认值 说明
reserved_space 400 索引初始容量
metric USearchMetricKind.COS 距离度量,默认余弦
connectivity 0 HNSW 图中单个节点的最大出边数;设 0 交由 usearch 自动配置
expansion_add 0 插入元素时投入的计算量(越大定位越准、代价越高);0 表示自动
expansion_search 0 查询时投入的计算量(越大结果越准、代价越高);0 表示自动

需要说明的一点限制:从源码看,UsearchKnnBruteForceKnn 目前仅实现了 query_as_of_now 变体(即查询时基于"当前时刻"的索引快照返回结果),其流式增量版本 query() 会抛出 NotImplementedErrornearest_neighbors.py L118-L121)。DocumentStore.retrieve_query 内部调用的正是 query_as_of_now(见下文),因此这一限制不影响文档检索主流程。

HybridIndexFactory:RRF 融合混合索引

HybridIndexFactory 将任意多个子索引组合成一个混合索引,使用**倒数排名融合(Reciprocal Rank Fusion, RRF)**合并结果:对每个子索引返回的每一行 d 赋予分数 1 / (k + rank(d)),再跨所有索引求和,最后按总分返回最优行。关键约束与参数:

  • retriever_factories:子索引工厂列表,必须至少提供 2 个,否则构造时即抛出 ValueError
  • k:RRF 平滑常数,默认 60

典型用法是将向量索引与 BM25 全文索引组合,兼顾语义相似与关键词精确匹配:

from pathway.stdlib.indexing import HybridIndexFactory, BruteForceKnnFactory, TantivyBM25Factory
from pathway.xpacks.llm.embedders import OpenAIEmbedder
import os

hybrid_factory = HybridIndexFactory(
    retriever_factories=[
        BruteForceKnnFactory(embedder=OpenAIEmbedder(api_key=os.environ["OPENAI_API_KEY"])),
        TantivyBM25Factory(),
    ],
    k=60,
)

RRF 打分与去重、截断逻辑完整实现在 HybridIndex._combine_results 中(hybrid_index.py L35-L122),包括按 (query_id, matched_id) 分组累加分数、按总分排序后按 k 截断。

构建 DocumentStore:索引流水线的核心

要操作索引、检索相关文档,需要创建 DocumentStore 对象。它会处理文档的解析(parsing)后处理(post-processing)切分(splitting),然后基于处理后的文本构建索引(retriever),并充当查询索引的统一接口。

DocumentStore 的构造函数参数(源码签名,document_store.py L75-L102):

参数 说明
docs 来自各类 connector 的表,须包含 bytes 类型的 data 列(通常将 connector 的 format 设为 "binary"/"raw");可含 _metadata 列用于过滤,许多 connector 提供 with_metadata=True 来返回该列
retriever_factory 上文介绍的索引工厂之一
parser 将文件内容解析为文档列表的可调用对象;默认 Utf8Parser
splitter 切分长文档的可调用对象;默认 NullSplitter(不切分)
doc_post_processors 可选的可调用对象列表,签名 (text: str, metadata: dict) -> tuple[str, dict],用于修改解析结果与元数据

最小示例:

from pathway.xpacks.llm.document_store import DocumentStore
from pathway.xpacks.llm.splitters import TokenCountSplitter
import pathway as pw

data_sources = pw.io.fs.read(
    "./sample_docs",
    format="binary",
    with_metadata=True,
)
text_splitter = TokenCountSplitter()
store = DocumentStore(
    docs=data_sources,
    retriever_factory=retriever_factory,
    splitter=text_splitter,
)

可以看到,构建 DocumentStore 需要准备好 splitter 并定义数据源。关于切分器的更多说明可参考 splitters 模块

内部流水线:build_pipeline 做了什么

构造 DocumentStore 时会自动执行 build_pipeline()document_store.py L320-L407),从源码结构看,其处理链是:

  1. 清洗与合并:把多个输入表统一选出 data(bytes)与 _metadata 两列并拼接;若表缺少 _metadata 列会发出警告并置为空 dict(此时过滤功能对该表失效);随后为每个文件注入 _file_id(文件 id 的字符串形式),用于后续追踪该文件被切成了多少块;
  2. PARSING:用 parser(默认 Utf8Parser)把 bytes 解析成文本与元数据;
  3. POST PROCESSING:依次执行所有 doc_post_processors
  4. CHUNKING:用 splitter(如 TokenCountSplitter)把长文本切块,元数据会传播到每个切块;
  5. INDEXING:调用 retriever_factory.build_index(chunked_docs.text, metadata_column=...) 真正建立索引,embedder 就在此环节对每个切块计算向量。

此外该流水线还会同步维护两张"观测表":progress_table(每个文件的切块数与是否已解析完成)和 stats(文件总数、最后修改时间、最后索引时间),它们正是下文 statistics_queryinputs_query 的数据来源。

准备查询与执行检索

Preparing Queries

把查询保存在 CSV 文件中,列定义如下:

必填 说明
query 你的问题
k 要检索的文档数量
metadata_filter 按元数据过滤文件(JMESPath 表达式)
filepath_globpattern 按路径 glob 模式缩小文件范围

示例:

printf "query,k,metadata_filter,filepath_globpattern\n\"Who is Regina Phalange?\",3,,\n" > queries.csv

这四种列在源码中被定义为一个预置 Schema——DocumentStore.RetrieveQuerySchema,其中 query: strk: int,另外两列为可空字符串。利用它连接 CSV:

query = pw.io.fs.read(
    "queries.csv",
    format="csv",
    # 查询表预定义 schema
    schema=DocumentStore.RetrieveQuerySchema
)

Retrieval

随后直接对 store 对象执行 retrieve_query,即可看到哪些文档切块可能包含回答查询所需的信息:

result = store.retrieve_query(query)

retrieve_query 实现 看,其返回的 result 列为 JSON,包含按距离升序排列的条目,每条为 {"text": ..., "metadata": ..., "dist": ...}dist 是距离,越小越相似)。test_document_store.py 中的 _test_vs 用例完整演示了这条链路:用 BruteForceKnnFactory + 假 embedding 模型建库,再按 RetrieveQuerySchema 发起查询并断言命中文本存在、dist 接近 0,可作为最小可运行的验证蓝本。

按文件过滤:metadata_filter 与 filepath_globpattern

DocumentStore 允许基于文件元数据或其路径缩小搜索范围,对应查询中的两个可选字段:

  • metadata_filter:以 JMESPath 风格表达式过滤 modified_atownercontains 等元数据字段(实际为 JMESPath 布尔表达式);
  • filepath_globpattern:按 glob 路径模式缩小文件范围。

示例——同时限定 owneralbert 且路径匹配 **/phoebe*

printf 'query,k,metadata_filter,filepath_globpattern\n"Who is Regina Phalange?",3,owner==`albert`,**/phoebe*\n' > queries.csv
query k metadata_filter filepath_globpattern
"Who is Regina Phalange?" 3 owner==`albert` **/phoebe*
query = pw.io.fs.read(
    "queries.csv",
    format="csv",
    schema=DocumentStore.RetrieveQuerySchema
)

result = store.retrieve_query(query)

两个过滤条件在内部如何协作?源码中的 _get_jmespath_filter UDF(document_store.py L34-L46)会把 metadata_filterfilepath_globpattern 合并为一条 JMESPath 表达式:路径条件被改写为 globmatch('<pattern>', path),两部分以 && 连接后交给底层索引作为 metadata_filter 使用。也就是说,glob 路径过滤最终也是以"元数据过滤"的形式下发到索引层的。

可用哪些元数据字段取决于所用 connector。可查看该 connector read 函数文档中 with_metadata 参数对应的元数据字段。例如 CSV connector 在 with_metadata=True 时可提供 created_atmodified_atownersizepathseen_at 等字段用于过滤。

Finding documents:inputs_query

如果只想按 glob 模式和元数据查找文件、而不涉及任何向量/全文检索,可以使用 inputs_query 方法。查询表只需要两列 metadata_filterfilepath_globpattern,遵循 DocumentStore.InputsQuerySchema。该 Schema 还有一个 return_status 布尔列(默认 False):置为 True 时结果会为每个文件附加 _indexing_status 字段,取值 INDEXED(已完成切块解析)或 INGESTED(仅已读入、尚未完成索引),可用于监控索引进度——这一状态来自 build_pipeline 中维护的 progress_table

import pathway as pw
from pathway.xpacks.llm.document_store import DocumentStore

inputs_queries = pw.debug.table_from_rows(
    schema=DocumentStore.InputsQuerySchema,
    rows=[(None, "**/*.py", False)],
)
input_results = store.inputs_query(inputs_queries)

通过 REST Server 暴露 DocumentStore

REST 服务器可以把 DocumentStore 作为服务对外暴露,供 API 请求访问。当需要把文档检索集成到更大的系统、尤其是从外部进程访问时非常有用。

from pathway.xpacks.llm.servers import DocumentStoreServer

PATHWAY_PORT = 8765
server = DocumentStoreServer(
    host="127.0.0.1",
    port=PATHWAY_PORT,
    document_store=store,
)
server.run(threaded=True, with_cache=False)

DocumentStoreServer 在源码中注册了三个端点(均支持 GET/POST):

  • /v1/retrieve:对应 retrieve_query,执行相似度检索;
  • /v1/statistics:对应 statistics_query,返回文件数、最后修改/索引时间等统计信息;
  • /v1/inputs:对应 inputs_query,返回输入文件清单(可含 _indexing_status)。

run() 的参数(servers.py L43-L89):threaded=True 表示在新线程中运行引擎(不阻塞当前进程,方便在 notebook 中边写边查);with_cache=True 时会对设置了 cache_strategy 的 UDF 启用持久化缓存(默认后端为本地 ./Cache 目录的 filesystem backend),对 embedding 这类高开销 UDF 可显著降低重复调用成本。

服务器运行后,即可向 API 发送请求:

curl -X POST http://localhost:8765/v1/retrieve \
     -H "Content-Type: application/json" \
     -d '{
           "query": "Who is Regina Phalange?",
           "k": 2
         }'

如果不想手写 HTTP 请求,仓库还提供了 DocumentStoreClient,封装了 /v1/retrieve/v1/statistics/v1/inputs 三个端点的 Python 调用:

from pathway.xpacks.llm.document_store import DocumentStoreClient

client = DocumentStoreClient(host="127.0.0.1", port=8765)
results = client.query(
    query="Who is Regina Phalange?",
    k=2,
    metadata_filter="owner==`albert`",
    filepath_globpattern="**/phoebe*",
)

客户端返回的结果已按 dist 升序排好序,可直接接入下游 RAG 问答流程。

小结

Pathway 的文档索引模块把 RAG 场景中"索引 + 检索"环节做成了流式组件:

  • 索引选型:向量侧可选 BruteForceKnnFactory(精确)或 UsearchKnnFactory(HNSW 近似);关键词侧用 TantivyBM25Factory(可调 ram_budgetin_memory_index);混合需求用 HybridIndexFactory(RRF 融合,k 默认 60,至少 2 个子索引);
  • 流水线DocumentStore 将"connector 读取 → 解析 → 后处理 → 切分 → 建索引"组织为一条可增量更新的 DAG,元数据(含 _file_id)贯穿全程,支撑进度追踪与过滤;
  • 查询retrieve_query 执行相似检索,inputs_query 纯按元数据/glob 找文件,statistics_query 查看库内统计;
  • 服务化DocumentStoreServer 暴露 /v1/retrieve/v1/statistics/v1/inputs 三个 REST 端点,配合 DocumentStoreClient 可无缝接入外部系统。

相关实现与验证入口:DocumentStore 源码KNN 工厂BM25 工厂混合索引REST 服务器文档存储测试

登录后查看全文
热门项目推荐
相关项目推荐