首页
/ Pathway xpack Embedders 实战:为 RAG 与视频检索流水线接入 OpenAI、Gemini、LiteLLM 与本地向量模型

Pathway xpack Embedders 实战:为 RAG 与视频检索流水线接入 OpenAI、Gemini、LiteLLM 与本地向量模型

2026-09-06 20:43:07作者:董斯意

本篇围绕 Pathway Live Data Framework xpack 中的五类文本嵌入器(Embedder)展开:OpenAIEmbedderLiteLLMEmbedderSentenceTransformerEmbedderGeminiEmbedderMarengoEmbedder。读完后你将掌握每类嵌入器的默认模型、构造参数、截断与批处理机制、重试与缓存策略,以及如何在 Python 代码或 Pathway 模板 YAML 中以标签(tag)方式声明嵌入器,并将其用于文档向量化与 Video RAG 索引。

为什么 RAG 流水线需要嵌入器

在把文档存入向量库时,你需要先为文本计算嵌入向量,并把向量与原始文档的引用一起保存;之后对查询文本计算嵌入,即可在向量库中找到与查询最接近的文档。这正是 RAG(Retrieval-Augmented Generation)检索链路的核心步骤。

在 Pathway 中,嵌入器被封装为 pathway.xpacks.llm.embedders 模块下的 UDF(用户自定义函数)。所有嵌入器都继承自 BaseEmbedder(pw.UDF),因此可以直接在表操作上调用,例如:

embeddings = documents.select(embedding=embedder(pw.this.text))

基类的通用能力定义见 embedders.py

  • __call__:把一列文本(pw.ColumnExpression[str])变成嵌入列;
  • get_embedding_dimension():默认实现是请求嵌入器对单个 "." 做嵌入,然后取向量长度。这个能力在构建向量索引工厂时被用来确定索引维度。

模块通过 __all__ 导出的完整类列表比官方文档多一个 BedrockEmbedder(见 embedders.py 导出列表),源码与本文介绍的五类嵌入器使用同一套构造模式,可作为 AWS Bedrock 场景的补充参考。

五类嵌入器一览

xpack 提供以下嵌入器:

嵌入器 底层 SDK 默认模型 说明
OpenAIEmbedder openai text-embedding-3-small 覆盖 OpenAI 全系嵌入模型,支持批处理与上下文截断
LiteLLMEmbedder litellm 无,必须指定 通过 LiteLLM 网关调用任意支持的嵌入模型
SentenceTransformerEmbedder sentence_transformers 必须指定 本地运行 Hugging Face SBERT 模型,无需 API key
GeminiEmbedder google.generativeai models/embedding-001(源码默认值) 调用 Google Gemini 嵌入服务
MarengoEmbedder twelvelabs marengo3.0 TwelveLabs 多模态嵌入,512 维共享文本/图像/音频/视频空间,面向 Video RAG

OpenAIEmbedder:批处理、截断与维度探测

OpenAIEmbedder 的默认模型是 text-embedding-3-small。最小用法是把目录中的文档解析成行后逐行嵌入:

import os
import pathway as pw
from pathway.xpacks.llm.parsers import UnstructuredParser
from pathway.xpacks.llm.embedders import OpenAIEmbedder

files = pw.io.fs.read(
    os.environ.get("DATA_DIR"),
    mode="streaming",
    format="binary",
    autocommit_duration_ms=50,
)

# Parse the documents in the specified directory
parser = UnstructuredParser(chunking_mode="paged")
documents = files.select(elements=parser(pw.this.data))
documents = documents.flatten(pw.this.elements)  # flatten list into multiple rows
documents = documents.select(text=pw.this.elements[0], metadata=pw.this.elements[1])

# Embed each page of the document
embedder = OpenAIEmbedder(api_key=os.environ["OPENAI_API_KEY"])
embeddings = documents.select(embedding=embedder(pw.this.text))

在 Pathway 模板(YAML 应用)中则用类标签声明:

embedder: !pw.xpacks.llm.embedders.OpenAIEmbedder
  model: "text-embedding-3-small"

构造函数参数

结合 OpenAIEmbedder 源码,其构造参数如下:

参数 默认值 说明
model "text-embedding-3-small" 模型 ID;设为 None 后可在每次 UDF 调用中按行指定(model=pw.this.model
api_key OPENAI_API_KEY 环境变量 也可通过 **openai_kwargs 传入
capacity None(不限并发) 异步执行器允许的最大并发操作数
retry_strategy ExponentialBackoffRetryStrategy() 失败重试策略;传 None 则退化为 NoRetryStrategy
cache_strategy None pw.udfs.DiskCache(),重启时避免重复嵌入
truncation_keep_strategy "start" 超上下文文本的截断保留策略,"start"/"end"/None(不截断)
batch_size 128 单次发给 API 的最大批大小,更大批次可缩短嵌入耗时
**openai_kwargs - 透传给 embeddings.create,如 encoding_formatfloatbase64)、usertimeoutextra_headers/extra_query/extra_body

参数既可以写在构造函数里,也可以在 UDF 调用时覆盖;若要在 UDF 调用中指定 model,需先构造时把 model 设为 None

批处理与逐行参数分流

__wrapped__ 的批处理逻辑(embedders.py)把 UDF 调用中的 kwargs 拆成两类:

  • 整批取值相同的参数进入 constant_kwargs,与 input 一起作为一个 batch 发给 client.embeddings.create
  • 各行取值不同的参数(如按行指定的 model)进入 per_row_kwargs,此时无法批处理,改用 asyncio.TaskGroup 为每行并发发起单条请求。

_split_batched_kwargs 负责这一分流判断。这意味着“同一批里混用多个 OpenAI 嵌入模型”是受支持的,只是会退化为逐行并发请求。

上下文截断机制

truncation_keep_strategy 不是 None 时,truncate_context 静态方法会用 tiktoken 按模型的分词器编码文本,超出模型上限的 token 数即被裁掉,保留开头("start")或结尾("end"):

  • 各模型上限定义在 constants.pytext-embedding-3-smalltext-embedding-3-largetext-embedding-ada-002 均为 8191 tokens;
  • 若所用模型不在上表中,源码会记录 error 日志并跳过截断(原样返回文本),文档提示这可能导致 API 报错;
  • 集成测试 test_embedders.py 分别用 65600 个 "A"+65600 个 "B" 的超长文本验证了 start/end 两种策略的保留方向,并验证短文本不会被改动。

维度探测:不发请求也能知道向量维度

OpenAIEmbedder 覆写了 get_embedding_dimension()embedders.py):

  • 模型在 OPENAI_EMBEDDERS_DIMENSIONS 查表内且未显式传 dimensions 时直接查表返回:text-embedding-3-small 为 1536 维、text-embedding-3-large 为 3072 维、text-embedding-ada-002 为 1536 维,不发起网络请求;
  • 否则(未知模型,或用 dimensions 参数缩短了向量)会真实请求一次,用 embedder 对 "." 的嵌入长度 得到维度。

对应测试 用 mock 客户端断言:已知模型 await_count == 0(零请求),未知模型或传了 dimensions=2 时恰好请求一次。对流水线构建者而言,这意味着用 text-embedding-3-small 建索引时不依赖 OpenAI API 的可达性。

重试发生在请求层

OpenAIEmbedderLiteLLMEmbedderBedrockEmbedder 都把客户端自身的重试禁用(max_retries=0 等),把重试收敛到 __wrapped__ 内用 self.retry_strategy.invoke(...) 包裹真实请求。源码注释说明了原因:让 get_embedding_dimension 这类绕过执行器、直接调用 __wrapped__ 的路径也被重试策略覆盖。测试用例 模拟了首次请求抛 RuntimeError("connection reset by peer") 后第二次成功的场景,断言 create.await_count == 2

LiteLLMEmbedder:模型必须显式指定

LiteLLMEmbedder 用于经 LiteLLM 调用的任意嵌入模型,构造时必须指定 model,没有默认值。基本用法:

from pathway.xpacks.llm import embedders

embedder = embedders.LiteLLMEmbedder(
    model="text-embedding-3-small", api_key=API_KEY
)
# Create a table with one column for the text to embed
t = pw.debug.table_from_markdown(
    """
text_column
Here is some text
"""
)
res = t.select(ret=embedder(pw.this.text_column))

YAML 模板写法:

embedder: !pw.xpacks.llm.embedders.LiteLLMEmbedder
  model: "text-embedding-3-small"

实现要点(LiteLLMEmbedder 源码):

  • 构造参数与 OpenAIEmbedder 相同地支持 capacityretry_strategycache_strategy,其余参数(api_baseapi_versionapi_keyapi_typecustom_llm_providerlitellm_call_id 等)通过 **llmlite_kwargs 透传给 litellm.aembedding
  • __wrapped__ 是逐文本嵌入(单条输入),空字符串会被替换为 "." 再发送,避免空输入报错;
  • OpenAIEmbedder 不同,它没有内置的上下文截断逻辑,超长文本的行为取决于 LiteLLM 后端。

SentenceTransformerEmbedder:本地离线嵌入

SentenceTransformerEmbedder 使用 Hugging Face 的 Sentence Transformers(SBERT)模型,完全在本地运行,无需任何 API key。模型名在构造时指定:

import pathway as pw
from pathway.xpacks.llm import embedders

embedder = embedders.SentenceTransformerEmbedder(model="intfloat/e5-large-v2")

# Create a table with text to embed
t = pw.debug.table_from_markdown('''
txt
Some text to embed
''')

# Extract the embedded text
t.select(ret=embedder(pw.this.txt))
embedder: !pw.xpacks.llm.embedders.SentenceTransformerEmbedder
  model: "intfloat/e5-large-v2"

源码 可以看到它的构造参数:

参数 默认值 说明
model 必填 模型名或本地路径
device "cpu" 运行设备,如 "cuda"
batch_size 1024 传给 encode 的批大小,GPU 上加大批长通常更快
call_kwargs {} 透传给每次 encode 调用,可在 UDF 调用时逐次覆盖
**sentencetransformer_kwargs - 透传给 SentenceTransformer 初始化

与云端嵌入器一样,它也实现了 constant_kwargs/per_row_kwargs 分流:整批参数一致时一次性 encode 整批,参数逐行不同时退化为逐条 encodeget_embedding_dimension 沿用基类实现——对 "." 实际编码一次后取向量长度(本地模型无需网络)。

GeminiEmbedder:Google 嵌入服务

GeminiEmbedder 封装 Google Gemini Embedding Services,可用于 Google 提供的文本嵌入模型(如文档示例中的 models/text-embedding-004):

import pathway as pw
from pathway.xpacks.llm import embedders

embedder = embedders.GeminiEmbedder(model="models/text-embedding-004")

# Create a table with a column for the text to embed
t = pw.debug.table_from_markdown('''
txt
Some text to embed
''')

t.select(ret=embedder(pw.this.txt))
embedder: !pw.xpacks.llm.embedders.GeminiEmbedder
  model: "models/text-embedding-004"

源码层面值得注意的两点(GeminiEmbedder):

  1. 构造默认模型model 的源码默认值是 "models/embedding-001"(文档示例用 text-embedding-004,两者均为 Google 可用模型,按构造时的实参生效)。
  2. 截断策略差异:docstring 明确说明 Gemini API 在文本超过模型上下文长度时会自行截断内容,因此该嵌入器不提供 OpenAIEmbedder 那样的 truncation_keep_strategy 参数。
  3. API key 可通过构造参数 api_key 传入,或设置 GOOGLE_API_KEY 环境变量;其余 **gemini_kwargs 透传给 genai.embed_content

MarengoEmbedder:面向 Video RAG 的多模态嵌入

MarengoEmbedder 使用 TwelveLabs 的 Marengo 多模态模型嵌入文本,产出 512 维向量,位于文本/图像/音频/视频共享的嵌入空间。这使得它能与 TwelveLabsVideoParser 生成的视频描述向量直接比对,是 Video RAG 流水线的天然检索端嵌入器。

前置依赖:twelvelabs SDK(pip install "pathway[twelvelabs]")以及 TwelveLabs API key——未显式传入 api_key 时从 TWELVELABS_API_KEY 环境变量读取。

import pathway as pw
from pathway.xpacks.llm import embedders

embedder = embedders.MarengoEmbedder(
    cache_strategy=pw.udfs.DiskCache(),  # don't re-embed documents on restarts
)

# Create a table with text to embed
t = pw.debug.table_from_markdown('''
txt
Some text to embed
''')

t.select(ret=embedder(pw.this.txt))
embedder: !pw.xpacks.llm.embedders.MarengoEmbedder
  model: "marengo3.0"
  cache_strategy: !pw.udfs.DiskCache

MarengoEmbedder 源码 可补充以下实现细节:

  • 单请求嵌入:Marengo 每次请求只嵌入一条文本,因此构造时固定 max_batch_size=1__wrapped__ 内用 asyncio.gather 并发发起单条请求,保持热路径不阻塞;
  • 并发上限capacity 默认 16,注释说明这是刻意留在 API 限速以下的值,账号额度更高时可自行调大;
  • 维度免探测get_embedding_dimension 默认直接返回构造参数 embedding_dimension(默认 512,常量 MARENGO_EMBEDDING_DIMENSION),建索引时无需访问 API;若构造时传 embedding_dimension=None,会改用一次性请求探测真实维度(源码 解释:探测走独立的短生命周期事件循环与一次性客户端,避免污染热路径缓存的异步客户端);
  • 缓存建议:docstring 建议生产环境配 pw.udfs.DiskCache(),避免流水线重启后重新嵌入全部文档——文档示例中 cache_strategy=pw.udfs.DiskCache() 正是此用途;
  • 重试测试 验证了首请求失败后按 retry_strategy 重试的行为。

嵌入器在模板 YAML 中的组织方式

除本文档的独立片段外,仓库中的真实应用配置展示了嵌入器在完整流水线里的位置,例如 integration_tests/rag_evals/app.yaml 用变量标签声明后再被解析器/向量库引用:

$embedder: !pw.xpacks.llm.embedders.OpenAIEmbedder
  model: "text-embedding-3-small"

docs/2.developers/7.templates/30.configure-yaml.mddocs/2.developers/7.templates/39.yaml-snippets/20.rag-configuration-examples.md 给出了 OpenAI、LiteLLM、SentenceTransformer、Gemini 四种嵌入器在 RAG 配置中的完整 YAML 组合示例,可作为模板开发的直接参照。

嵌入器与向量库的衔接在 vector_store.py 中体现:DefaultKnnFactory 接收 embedder 回调,InMemoryVectorStore 持有 embedder 引用用于对文档和查询做嵌入。

选型与使用要点小结

  • 优先 OpenAI:需要批处理效率与可控截断时,OpenAIEmbedder 是唯一实现了批大小(默认 128)与 tiktoken 截断的云端嵌入器;
  • 多云/私有网关LiteLLMEmbedder 覆盖任意 LiteLLM 支持的模型,注意它没有默认模型与内置截断;
  • 完全离线SentenceTransformerEmbedder 本地推理,无网络依赖,device="cuda" 可启用 GPU;
  • Google 生态GeminiEmbedder 简单直接,超长文本由 API 侧截断;
  • 视频检索MarengoEmbedder 的 512 维多模态共享空间使其成为 TwelveLabsVideoParser 流水线的配套检索嵌入器,生产环境务必考虑 DiskCache
  • 通用能力:所有嵌入器均可通过 get_embedding_dimension() 供向量索引工厂确定维度;云端嵌入器的重试被收敛在请求层,ExponentialBackoffRetryStrategy 是默认策略;cache_strategycapacity 的语义在各类构造参数中一致。

以上行为均以本仓库源码与测试为准:实现位于 python/pathway/xpacks/llm/embedders.py,模型常量位于 python/pathway/xpacks/llm/constants.py,行为验证位于 python/pathway/xpacks/llm/tests/test_embedders.py

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