首页
/ Pathway LLM xpack Embedders 实战指南:五种内置嵌入 UDF 的原理、参数与源码级解析

Pathway LLM xpack Embedders 实战指南:五种内置嵌入 UDF 的原理、参数与源码级解析

2026-09-06 21:51:07作者:曹令琨Iris

本文基于 Pathway 仓库中 70.embedders.md 用户指南,系统讲解 pathway.xpacks.llm.embedders 模块提供的五类文本嵌入包装器(OpenAI、LiteLLM、Sentence-Transformer、Gemini、Marengo)的用法、默认参数与底层实现。读完后,你可以直接在向量库管线、文档索引和视频 RAG 场景中复制可运行的嵌入代码,并理解每个包装器在批量处理、截断、重试、缓存与维度探测上的源码级行为。

为什么向量库需要 Embedder

把文档存入向量库时的标准流程是:先对文本计算嵌入向量,再将向量与指向原文档的引用一起存储;查询时计算查询语句的嵌入,然后检索与查询向量最接近的文档。在 Pathway Live Data Framework 中,这一步由 pathway xpack 提供的嵌入包装器完成,它们以 Pathway UDF 的形式嵌入到 select 调用中,与流式处理管线无缝集成。xpack 当前提供的包装器为:

  • OpenAIEmbedder —— 使用任意 OpenAI 嵌入模型;
  • LiteLLMEmbedder —— 使用任意可通过 LiteLLM 调用的模型;
  • SentenceTransformerEmbedder —— 使用 Hugging Face Sentence-Transformers(SBERT)生态中的模型;
  • GeminiEmbedder —— 使用 Google 的嵌入模型;
  • MarengoEmbedder —— 使用 TwelveLabs 的多模态 Marengo 模型,面向 Video RAG 管线。

五类包装器的完整实现集中在 embedders.py,其中 __all__ 还导出了第六个 BedrockEmbedder(见文末补充)。

统一底座:BaseEmbedder 与 UDF 机制

所有包装器都继承自 embedders.py 中的 BaseEmbedder,它直接继承 pw.UDF,因此可以直接在表表达式中调用:

res = t.select(ret=embedder(pw.this.text_column))

两个关键机制值得注意:

  1. 列级批量调用BaseEmbedder.__call__ 接收一个 ColumnExpression[str] 列,Pathway 引擎会自动按批(batch)调用底层 __wrapped__ 方法,把整列文本一次性送入嵌入服务。
  2. 维度自探测get_embedding_dimension() 通过让嵌入器嵌入一个 "." 并返回向量长度来获得维度(见 embedders.py#L78-L85)。这是向量库工厂建立索引时确定向量维度数的途径——你不需要手工指定维度,管线搭建时会向嵌入器发起探测(对某些模型有零请求优化,下文详述)。

各包装器在构造函数中通过 udfs.async_executor(capacity=...) 挂载异步执行器,并统一接受 capacity(并发上限,默认 None 即不限)、retry_strategy(默认 ExponentialBackoffRetryStrategy 指数退避重试)、cache_strategy(默认 None,可传 pw.udfs.DiskCache() 等持久化缓存策略)三个 UDF 层参数。

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))

API Key 必须通过构造参数 api_key 提供,或设置 OPENAI_API_KEY 环境变量;构造函数会先用 openai.AsyncOpenAI(api_key=..., max_retries=0) 校验 key 是否可用(客户端初始化被延迟到首次真正调用时,以避免事件循环问题,见 embedders.py#L212-L230)。

关键构造参数(源码默认值)

OpenAIEmbedder.init 可确认完整参数表:

参数 默认值 说明
model "text-embedding-3-small" 模型 ID,也可以在 UDF 调用时逐行指定
capacity None 最大并发操作数
retry_strategy ExponentialBackoffRetryStrategy() 失败重试策略,None 时退化为 NoRetryStrategy
cache_strategy None 缓存策略,如 pw.udfs.DiskCache()
truncation_keep_strategy "start" 截断保留策略:"start" 保留开头、"end" 保留结尾、None 不截断(可能导致 API 报错)
batch_size 128 单次发送到 API 的最大批量
**openai_kwargs 透传给 OpenAI 的参数(encoding_formatdimensionsusertimeout 等)

截断:用 tiktoken 按模型 token 上限预裁剪

truncation_keep_strategy 不是简单切片,而是基于 tiktoken 按模型真实 token 上限裁剪。truncate_contextconstants.py 中查询模型的最大 token 数——当前查表覆盖 text-embedding-3-smalltext-embedding-3-largetext-embedding-ada-002,均为 8191;若模型不在表中则记录错误日志并跳过截断(可能引发后续 API 异常)。测试用例 test_embedders.py#L45-L69 验证了 start/end 两种策略分别保留文本前段或后段,且短文本不会被截断。

批量与逐行两种调用路径

wrapped_split_batched_kwargs 区分“批内所有行相同”的参数与“逐行不同”的参数:参数一致时整个 batch 一次 embeddings.create 请求发出去;只要 model 等参数逐行不同(例如 t.select(ret=embedder(pw.this.txt, model=pw.this.model))),就改用 asyncio.TaskGroup 对每行单独发请求并发执行。

维度探测的零请求优化

OpenAIEmbedder 覆写了 get_embedding_dimension:当没有显式 dimensions 参数且模型在 OPENAI_EMBEDDERS_DIMENSIONS 查表中(text-embedding-3-small → 1536、text-embedding-3-large → 3072、text-embedding-ada-002 → 1536)时,直接返回查表值,不发任何网络请求;否则回退到基类行为,用真实请求探测。这一点被测试 test_openai_embedder_dimension_of_known_model_sends_no_request 显式验证(断言 create.await_count == 0),并有两条边界:未知模型会真实探测;传入 dimensions=2 这类降维参数后查表失效、改走探测。

LiteLLMEmbedder:任意提供方、必须显式指定模型

LiteLLMEmbedder 是对 litellm.aembedding 的封装(embedders.py#L354-L451)。与 OpenAI 版不同,它没有默认模型——model 必须在构造时或每次 UDF 调用中给出,否则调用会失败。指南示例:

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))

实现细节上,它的重试与 OpenAI 版同样放在 __wrapped__ 内部(而非执行器层),使得 get_embedding_dimension 这类绕过执行器的直连调用也被重试覆盖;单次调用以单字符串为输入(litellm.aembedding 返回 ret.data[0]["embedding"])。除 model 外,api_baseapi_versionapi_typecustom_llm_providertimeout(默认 10 分钟)等参数均可通过构造函数或调用时透传。测试 test_litellm_embedder_wrapped_retries_transient_errors 模拟了连接中断后按 FixedDelayRetryStrategy 重试并最终成功的路径。

SentenceTransformerEmbedder:本地模型,离线可用

SentenceTransformerEmbedderembedders.py#L454-L546)接入 Hugging Face Sentence-Transformers 模型库,是唯一完全本地运行的包装器——它通过 optional_imports("xpack-llm-local") 按需导入 sentence_transformers,不依赖任何云端 API。指南示例:

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))

构造参数(SentenceTransformerEmbedder.init):

  • model:模型名或本地路径(必填);
  • device:运行设备,默认 "cpu",有 GPU 时可传 "cuda"
  • batch_size:默认 1024,文档注明更大的 batch 在 GPU 上尤其能降低嵌入耗时;
  • call_kwargs:透传给 SentenceTransformer.encode 的参数(如 normalize),且可在每次 UDF 调用中覆盖;
  • **sentencetransformer_kwargs:透传给 SentenceTransformer 构造函数的参数。

它的 __wrapped__同步函数(本地推理,无网络),同样实现了“批内参数一致则整批 encode、不一致则逐行调用”的分批逻辑(embedders.py#L514-L536),维度探测走基类的真实嵌入 "." 路径。

GeminiEmbedder:Google 嵌入服务

GeminiEmbedder 封装 Google Gemini Embedding 服务(embedders.py#L549-L638),通过 google.generativeaiembed_content 完成嵌入。指南示例:

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))

需要注意的几个实现事实:

  • 源码中 model 的默认值是 "models/embedding-001"embedders.py#L607),指南示例则显式传入 models/text-embedding-004;可选模型清单以 Google Gemini 官方文档为准;
  • API Key 可通过构造参数 api_key、调用参数或 GOOGLE_API_KEY 环境变量提供;
  • 官方文档明确说明:Gemini API 在服务端对超出模型上下文长度的文本做截断,因此该包装器与 OpenAIEmbedder 不同,没有客户端侧的截断逻辑
  • **gemini_kwargs 会原样透传给 genai.embed_content

MarengoEmbedder:Video RAG 的多模态嵌入

MarengoEmbedder 使用 TwelveLabs 的多模态 Marengo 模型,产出 512 维向量,文本、图像、音频、视频嵌入位于同一共享空间,因此是与 TwelveLabsVideoParser 生成的视频描述天然配套的视频 RAG 嵌入器(解析器见 50.parsers.md 中的 TwelveLabsVideoParser 一节)。它要求安装 twelvelabs SDK(pip install "pathway[twelvelabs]")和 TwelveLabs API key——key 显式传入或通过 TWELVELABS_API_KEY 环境变量读取,缺失时 _resolve_twelvelabs_api_key 会直接抛出带明确提示的 ValueError。指南示例:

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))

MarengoEmbedder 实现 可确认以下默认值与设计取舍:

  • model 默认 "marengo3.0"(常量 DEFAULT_MARENGO_MODEL),embedding_dimension 默认 512;
  • capacity 默认 16——源码注释说明该值刻意低于 API 限速,账号允许时可自行调高;
  • max_batch_size=1:Marengo 每次请求只嵌入一段文本,因此批处理层把批拆成单条,再在 wrapped 内用 asyncio.gather 并发发出多个请求,保持嵌入热路径非阻塞;
  • 维度探测(get_embedding_dimension)默认直接返回 512,不需要在管线搭建时连接 TwelveLabs API;只有构造时显式传 embedding_dimension=None 才会用一次性的独立请求探测(探测用短命事件循环中的临时客户端,与热路径缓存的 aclient 严格隔离)。

cache_strategy=pw.udfs.DiskCache() 是生产环境的推荐配置:嵌入结果按内容哈希落盘,管线重启后不会重复为相同文档付费调用 API。测试 test_marengo_embedder_wrapped_retries_transient_errors 验证了 Marengo 请求在瞬时故障下的重试行为。

源码中的第六个:BedrockEmbedder

当前仓库的 embedders.py 还导出了用户指南列表之外的 BedrockEmbedder实现见 L641-L799):封装 AWS Bedrock 上的 Amazon Titan 与 Cohere 嵌入模型,默认 model_id="amazon.titan-embed-text-v2:0",支持 region_name、显式 AWS 凭证或 IAM 默认凭证链,并按模型类型自动切换 Titan/Cohere 两种请求与响应格式;相关初始化与重试行为在 test_embedders.py 中有参数化测试覆盖。

与向量库、模板的集成方式

传入向量库组件

上述任何包装器都可以直接作为 embedder 参数传给 Live Data Framework 的向量库工厂。例如 20.llm-app.md30.docs-indexing.md 中的文档索引示例:

from pathway.xpacks.llm.embedders import OpenAIEmbedder
from pathway.xpacks.llm import vector_store

embedder = OpenAIEmbedder(api_key=os.environ["OPENAI_API_KEY"])
vs = vector_store.ElasticsearchVectorStore(
    address=ADDRESS,
    api_key=ES_API_KEY,
    index_name="index",
    document_field="text",
    embedder=embedder,
)

vector_store.py 的结构看,向量库工厂在创建索引时会调用 get_embedding_dimension() 确定维度——这正解释了前文“维度探测”对管线搭建体验的意义。

模板中的 YAML 声明式配置

Pathway 模板(template)体系允许用 YAML 标签直接实例化嵌入器,五类包装器在指南中给出的写法分别是:

embedder: !pw.xpacks.llm.embedders.OpenAIEmbedder
  model: "text-embedding-3-small"
embedder: !pw.xpacks.llm.embedders.LiteLLMEmbedder
  model: "text-embedding-3-small"
embedder: !pw.xpacks.llm.embedders.SentenceTransformerEmbedder
  model: "intfloat/e5-large-v2"
embedder: !pw.xpacks.llm.embedders.GeminiEmbedder
  model: "models/text-embedding-004"
embedder: !pw.xpacks.llm.embedders.MarengoEmbedder
  model: "marengo3.0"
  cache_strategy: !pw.udfs.DiskCache

模板中嵌入器还可以作为变量复用,例如 MCP Server 模板(见 40.live-data-framework-mcp-server.md)中的 $embedder 变量被文档索引与检索两处共同引用。

选型与实操要点小结

  • 云端 OpenAI 系:首选 OpenAIEmbedder,享受批量请求、tiktoken 截断、维度查表零请求探测;需要逐行不同模型时自动退化为并发单请求。
  • 多提供方统一入口LiteLLMEmbedder 适合把多种 provider 收敛到同一套管线代码,代价是没有默认模型、无批量发送。
  • 离线/本地推理SentenceTransformerEmbedder 零 API 成本,注意 devicebatch_size 对吞吐的影响。
  • Google 生态GeminiEmbedder 简洁直接,截断由服务端处理,模型选择以 Gemini 官方文档为准。
  • Video RAGMarengoEmbedder 与 TwelveLabs 视频解析器配套,512 维多模态共享空间,生产环境务必配 DiskCache
  • 所有包装器都接受 capacity/retry_strategy/cache_strategy 三个 UDF 层参数,重试默认指数退避,缓存默认关闭——长文档库场景建议显式开启缓存以避免重启后重复嵌入。

相关文档可进一步阅读:向量库指南解析器指南LLM 应用示例;实现与测试参见 embedders.pyconstants.pytest_embedders.py

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