首页
/ Pathway LLM xpack Rerankers 实战:LLMReranker、CrossEncoderReranker 与 EncoderReranker 的文档重排实现

Pathway LLM xpack Rerankers 实战:LLMReranker、CrossEncoderReranker 与 EncoderReranker 的文档重排实现

2026-09-04 23:00:52作者:魏侃纯Zoe

在 RAG 系统中,向量检索往往只能给出"看起来相似"的文档,而 rerankers 负责对候选文档做二次精排,把真正能回答问题的文档排到前面。本文基于 Pathway Live Data Framework 的用户指南文档 rerankers.md,结合源码 rerankers.py、提示词模块 prompts.py 以及测试 test_rerankers.py,完整讲解 Pathway LLM xpack 提供的三类重排器的用法、参数与底层实现。读完本文,你可以直接在 Pathway 的流式表格管道中为任意 (query, doc) 对计算相关度分数,并将重排结果接入 RAG 问答服务。

为什么需要 Reranker

原文档给出的动机很直接:RAG 系统的第一级稀疏/向量检索通常基于余弦相似度,它把一个段落的全部语义压缩进单个 embedding 向量,速度快但可能丢失细节,检索结果里常常混入与 query 无关的文档。

Reranker 的做法是把 (query, document) 配对交给一个更强的评估模型,让它重新判断给定 document 是否有助于回答 query。Pathway LLM xpack 提供了三类重排器:

重排器 评估方式 分数范围 依赖
LLMReranker 让任意 LLM 按提示词打分 1–5(整数转 float) LLM Chat 实例
CrossEncoderReranker Sentence-Transformers 的 CrossEncoder,对文本对联合建模 0..1(带激活函数时)或原始 logits xpack-llm-local(本地模型)
EncoderReranker SentenceTransformer 编码后计算相似度 余弦相似度 xpack-llm-local(本地模型)

此外源码中还有一个未在文档正文展开的 FlashRankReranker,基于 flashrank 库做重排,可视为第四种可选实现。

安装依赖

根据 LLM xpack 概览文档pyproject.toml 中的 extras 定义:

# 基础 LLM xpack(含 OpenAI/LiteLLM 等 Chat 封装,LLMReranker 可用)
pip install "pathway[xpack-llm]"

# 本地 ML 推理依赖(CrossEncoderReranker / EncoderReranker 需要)
pip install "pathway[xpack-llm-local]"

pyproject.toml 可以看到,xpack-llm-local 这个 extra 实际安装了 sentence_transformerstransformers(>= 4.50.2, < 5.0)。这与源码中的懒加载机制对应:CrossEncoderRerankerEncoderReranker 的构造函数都通过 optional_imports("xpack-llm-local") 延迟导入 sentence_transformers,只有在真正实例化时才要求该依赖存在(见 rerankers.py 第 203 行第 268 行)。因此只跑 LLMReranker 的管道无需安装本地模型栈。

统一的调用模式

三个示例共享同一个数据形态:文档列存 {"text": ...} 的 JSON 对象,query 列存问题文本,然后用 select 把重排器应用到 (docs["text"], prompt) 列对上。原文档中的示例数据如下:

docs = [
    {"text": "John drinks coffee"},
    {"text": "Someone drinks tea"},
    {"text": "Nobody drinks coca-cola"},
]

query = "What does John drink?"

df = pd.DataFrame({"docs": docs, "prompt": query})

这一模式与源码签名一致:三个重排器的 __call__ 都是 reranker(doc: pw.ColumnExpression, query: pw.ColumnExpression) -> pw.ColumnExpression,返回一个 float 类型的分数列(见 rerankers.py 第 104 行)。

LLMReranker:让 LLM 按 1–5 分制打分的点式重排

LLMReranker 把"评估文档与 query 的相关性"这件事直接外包给任意 LLM Chat 实例。文档示例:

from pathway.xpacks.llm import rerankers
from pathway.xpacks.llm import llms
import pandas as pd

docs = [
    {"text": "John drinks coffee"},
    {"text": "Someone drinks tea"},
    {"text": "Nobody drinks coca-cola"},
]

query = "What does John drink?"

df = pd.DataFrame({"docs": docs, "prompt": query})

chat = llms.OpenAIChat(model="gpt-4o-mini", api_key=API_KEY, response_format="{'type': 'json_object'}")
reranker = rerankers.LLMReranker(llm=chat)

input = pw.debug.table_from_pandas(df)
res = input.select(rank=reranker(pw.this.docs["text"], pw.this.prompt))

构造参数

rerankers.py 第 90 行 的构造函数签名看,它接受三个参数:

  • llmllms.BaseChat 子类实例,即执行打分的 Chat 模型。LLMReranker 与 xpack 中所有 Chat 封装(OpenAI、Cohere、LiteLLM 等)解耦,可替换为任意已支持的 LLM;
  • prompt_template(可选):str / 双参可调用对象 / pw.UDF,用于生成给 LLM 的提示词,默认是 prompts.prompt_rerank;字符串模板必须包含 {context}{query} 两个占位符,否则会在校验时抛出 ValueError(见 prompts.py 第 82 行RAGPromptTemplate 校验逻辑);
  • response_parser(可选):pw.UDFCallable[[str], float],负责把 LLM 的原始响应解析为 float,默认是 prompts.parse_score_json

默认提示词与分数解析

默认提示词 prompt_rerank 要求模型扮演 RAG 助手,将文档与问题的相关度评为 1–5 分(5 = 文档足以回答该问题,1 = 完全无关),并且只允许输出 {"score": <int>} 形式的 jsonl。配套的 parse_score_jsonjson.loads 取出 score 字段转成 float,如果响应不是合法 JSON(例如模型返回了纯文本),会抛出 ValueError(f"Expected a json response, got {text}.")

这也解释了文档示例中为什么给 OpenAIChat 传了 response_format="{'type': 'json_object'}"——强制 JSON 输出能显著降低 parse_score_json 解析失败的概率。

内部执行链路与温度控制

LLMReranker.__call__rerankers.py 第 104 行)的完整链路是:

  1. _extract_value UDF 把列值解包为普通 Python 值;
  2. 用提示词模板 UDF 生成 prompt;
  3. temperature=0 调用 LLM(重排场景要求输出稳定,源码里硬编码了 call_kwargs = dict(temperature=0));
  4. 用响应解析 UDF 把响应解析为 float 分数列。

其中第 4 步有一个细节:如果底层 LLM 的执行器是 FullyAsyncExecutor,解析函数会被 _coerce_fully_async 包裹(rerankers.py 第 130 行),使解析逻辑也在异步执行器内完成,保持整条链路的一致性。

测试用例印证

test_rerankers.py 用 mock 的 OpenAIChat 子类验证了以上行为:返回 '{"score": 1}' 时得到分数 1.0,返回 '{"score": 5}' 时得到 5.0;返回非 JSON 文本 'text' 时,pw.debug._compute_tables 触发 ValueError(对应 parse_score_json 的异常分支);使用 pw.udfs.fully_async_executor() 的 LLM 也能正常得到分数(对应上面的 fully-async 分支,见 test_rerankers.py 第 34 行)。

CrossEncoderReranker:基于 CrossEncoder 的文本对打分

CrossEncoderReranker 处理 (query, document) 文本对,计算 0..1 的相关度分数;如果未传入激活函数,则输出原始 logits。文档示例:

from pathway.xpacks.llm import rerankers
import pandas as pd
import torch

docs = [
    {"text": "John drinks coffee"},
    {"text": "Someone drinks tea"},
    {"text": "Nobody drinks coca-cola"},
]

query = "What does John drink?"

df = pd.DataFrame({"docs": docs, "prompt": query})

reranker = rerankers.CrossEncoderReranker(
    model_name="cross-encoder/ms-marco-MiniLM-L-6-v2",
    default_activation_function=torch.nn.Sigmoid(),  # 输出落在 0..1
)

input = pw.debug.table_from_pandas(df)
res = input.select(
    rank=reranker(pw.this.docs["text"], pw.this.prompt), text=pw.this.docs["text"]
)
pw.debug.compute_and_print(res)

构造参数与实现细节

rerankers.py 第 194 行 看,构造签名为 CrossEncoderReranker(model_name, *, cache_strategy=None, **init_kwargs)

  • model_name:CrossEncoder 模型名。文档示例使用 cross-encoder/ms-marco-MiniLM-L-6-v2;源码 docstring 中推荐的轻量模型是 cross-encoder/ms-marco-TinyBERT-L-2-v2
  • cache_strategy:Pathway UDF 的缓存策略,可传入 CacheStrategy 实例(如 DiskCache)对重复的 (query, doc) 打分做缓存,降低本地模型推理开销;
  • **init_kwargs:其余参数直接透传给 sentence_transformers.CrossEncoder 构造函数——示例里的 default_activation_function=torch.nn.Sigmoid() 就是通过这里传入的,它把原始 logits 压到 0..1 区间,使不同模型/不同调用之间分数可比。

实际打分逻辑在 __wrapped__ 中(rerankers.py 第 208 行):把 (query, doc) 组成单元素对 [[query, doc]] 调用 self.model.predict(...),返回第一个分数。由于它继承自 pw.UDF,在表格管道中每行独立调用、可水平并行,**kwargs 在调用时还能覆盖构造时的默认参数。

模型行为细节(如可用预训练 CrossEncoder 列表)可参考 sentence_transformers 官方的 CrossEncoder 文档。

EncoderReranker:基于 SentenceTransformer 的相似度重排

EncoderReranker 用 SentenceTransformer 编码器分别编码 query 与 doc,再计算相似度。文档示例:

from pathway.xpacks.llm import rerankers
import pandas as pd

docs = [
    {"text": "John drinks coffee"},
    {"text": "Someone drinks tea"},
    {"text": "Nobody drinks coca-cola"},
]

query = "What does John drink?"

df = pd.DataFrame({"docs": docs, "prompt": query})

reranker = rerankers.EncoderReranker(
    model_name="all-mpnet-base-v2",
)

input = pw.debug.table_from_pandas(df)
res = input.select(
    rank=reranker(pw.this.docs["text"], pw.this.prompt), text=pw.this.docs["text"]
)

实现上(rerankers.py 第 259 行),它构造 SentenceTransformer(model_name, **init_kwargs),打分时执行:

embeddings = self.model.encode([query, doc], normalize_embeddings=True, **kwargs)
return embeddings[0] @ embeddings[1].T

即把 query 和 doc 各自编码为归一化向量后做点积——归一化向量的点积等价于余弦相似度,因此分数落在 -1..1 之间、典型相关文本对为正数。构造参数与 CrossEncoderReranker 相同:model_name + cache_strategy + 透传给 SentenceTransformer**init_kwargs。源码 docstring 给出的推荐模型是 BAAI/bge-large-zh-v1.5(面向中文场景);文档示例则用英文通用的 all-mpnet-base-v2。与 CrossEncoder 不同,它没有 query-doc 联合交互,推理更便宜,但判别精度通常弱于 CrossEncoder——从源码结构看,这正是它被定位为"measure similarity"类重排器的原因。

rerank_topk_filter:按分数取 Top-K

单行逐条打分解决"每篇文档多相关"的问题,而"从一堆候选里只留最好的 k 篇"则由模块级 UDF rerank_topk_filter 完成。它的签名是 rerank_topk_filter(docs: list[Doc], scores: list[float], k: int = 5) -> tuple[list[Doc], list[float]]:按分数降序排序后截取前 k 篇,返回 (docs, scores) 元组并打印保留数量与分数日志。

其 docstring 内置了一个可复制的 doctest 场景:先用 pw.reducers.tuple 把每行的 docs 与 scores 聚合成列表,再应用过滤器,最后取元组解包:

docs_table = table.reduce(
    doc_list=pw.reducers.tuple(pw.this.docs),
    score_list=pw.reducers.tuple(pw.this.reranker_scores),
)
docs_table = docs_table.select(
    docs_scores_tuple=rerankers.rerank_topk_filter(
        pw.this.doc_list, pw.this.score_list, 2
    )
)

单元测试 test_rerankers.py 第 64 行 验证了排序正确性:10 个文档中分数 9.5 出现了两次(下标 5 和 9),取 top-3 的结果为下标 [5, 9, 6],对应分数 [9.5, 9.5, 5.555],说明稳定排序下并列分数按原顺序保留。

在 RAG 问答服务中接入 Reranker

重排器最重要的实战场景是 Pathway 的模块化 RAG。BaseRAGQuestionAnswerer 的构造函数接收两个可选参数:

  • reranker:任意 pw.UDF 重排器实例(上面三种均可),默认 None(即标准 RAG,不做重排);
  • rerank_topk:重排后保留的文档数。默认 None 表示禁用重排。

注意构造函数中的校验逻辑(question_answering.py 第 559 行):rerankerrerank_topk 必须同时提供,只给其一时会打印 "Incomplete reranker configuration" 警告并整体禁用重排,而不是抛错——排查"重排不生效"问题时应先检查这一配置。

重排发生在 answer_query 中:向量库按 search_topk(默认 6)取出候选文档后,若配置了 reranker 则进入 _apply_rerankingquestion_answering.py 第 588 行)。其流水线是典型的 Pathway 表格操作:

  1. flatten 把每行的 docs 列表炸开成逐行记录;
  2. 对每行调用 self.reranker(pw.this.docs["text"], pw.this.prompt) 得到 reranker_score,随后 await_futures() 等待异步打分完成;
  3. -reranker_scoresort_keygroupby(query_id, sort_by=sort_key) 重新聚合,把打分后的文档(附带 reranker_score 字段)按分数从高到低排好并打包回 tuple;
  4. _limit_documents(docs, k=rerank_topk) 只保留 top-k 篇,再交给 context processor 组装提示词、调用 LLM 生成答案。

这样,一个 RAG 服务只需在构造 BaseRAGQuestionAnswerer 时传入例如 reranker=CrossEncoderReranker(model_name="cross-encoder/ms-marco-MiniLM-L-6-v2", cache_strategy=DiskCache())rerank_topk=3,即可获得"向量粗排 6 篇 → CrossEncoder 精排 → 保留 3 篇入上下文"的两阶段检索。

选型小结

维度 LLMReranker CrossEncoderReranker EncoderReranker
分数语义 1–5 整数(提示词约定) 0..1(配 Sigmoid)或 logits 余弦相似度(-1..1)
每文档成本 一次 LLM 调用(temperature=0) 一次本地 CrossEncoder 推理 两次 encode + 点积
依赖 任意 BaseChat xpack-llm-local xpack-llm-local
可定制点 提示词模板 / 响应解析器 **init_kwargs 透传 CrossEncoder **init_kwargs 透传 SentenceTransformer
适用场景 无本地 GPU/模型环境、追求语义级判断 需要 query-doc 联合建模的精确重排 已有 embedding 基础设施、追求低延迟

三者都实现了"对 (query, doc) 列对逐行打分、返回 float 分数列"的统一接口,因此可以互换地使用在 select 表达式或 BaseRAGQuestionAnswererreranker 参数中;打分之后再配合 rerank_topk_filter 或 RAG 内置的 rerank_topk 截断,即构成完整的两阶段检索。实现与验证代码集中在 python/pathway/xpacks/llm/rerankers.pypython/pathway/xpacks/llm/prompts.pypython/pathway/xpacks/llm/question_answering.pypython/pathway/xpacks/llm/tests/test_rerankers.py,可直接在仓库中继续深入。

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