在 Pathway Templates 的 YAML 配置中使用自定义组件:mapping tags 导入机制与实战示例
本篇技术指南基于 Pathway 开发者文档 35.custom-components.md,系统讲解如何在使用 YAML 配置文件定制 Pathway 模板(如各类 RAG / LLM 应用模板)时,把自己编写的 Python 组件接入配置:包括 ! mapping tags 的导入规则、模块解析原理,以及"补充元数据、替换 Parser、实例化第三方库对象"三类真实可运行示例。读完本文,你将掌握在不改动模板 Python 代码的前提下,仅通过 YAML 就引入任意自研函数、类与外部库对象的完整方法,并能理解其底层由 pw.load_yaml 驱动的加载机制。
前提说明:本文聚焦"自定义组件的导入",不覆盖 Pathway YAML 配置的基础语法(如
$变量、环境变量、Schema 定义)。如需了解基础语法,请先阅读 30.configure-yaml.md。
组件导入机制:! 之后写什么
Pathway 的 YAML 配置使用自定义解析器,允许通过 mapping tags(YAML 中以 ! 开头的标记)在配置文件中引用真实的 Python 对象。凡是想在配置中出现的对象,都必须在 ! 后给出 完整的模块名(full module name),例如:
llm: !pw.xpacks.llm.llms.OpenAIChat
对于 Pathway 自身,pw 是唯一的缩写形式,解析器会把 pw. 前缀自动替换为 pathway。例如 !pw.xpacks.llm.llms.OpenAIChat 最终等价于导入 pathway.xpacks.llm.llms 模块并取出 OpenAIChat 类。
无模块名时默认指向 builtins
如果在 ! 后只给出对象名而不带任何模块路径,解析器会假定它来自 Python 内置命名空间(builtins)。例如用 !str 把对象转换为字符串:
one_as_string: !str
object: 1
上面这段 YAML 等价于在 Python 中调用 str(object=1),其结果 one_as_string 的值是字符串 "1"。这印证了 mapping tag 的本质:标签代表一个可调用对象(函数或类),其后的 mapping 会被展开为该对象的 **kwargs。
自研模块的命名即文件名
因此,任何你想要在配置里引用的对象,都必须能够通过"模块名"被导入。对你自己编写的代码来说,模块名通常就是它所在的 Python 文件名。例如,你要使用的函数 foo 定义在与 app.py、app.yaml 同目录下的 utils.py 中:
.
├── app.py
├── app.yaml
└── utils.py
那么 YAML 中就以 !utils.foo 来引用它——此时无需在任何 Python 文件中手动 import utils,加载 YAML 的过程会自动完成模块导入。
底层原理:pw.load_yaml 与 import_object
这套行为由 Pathway 的 YAML 加载器实现,核心源码在 python/pathway/internals/yaml_loader.py:
- 解析器类
PathwayYamlLoader通过add_multi_constructor("!", ...)注册了对所有!前缀 tag 的处理(见 yaml_loader.py); - 其中关键函数
import_object(见 yaml_loader.py)负责把!后的字符串解析成真实对象,其逻辑是:把以pw./pw:开头的路径改写为pathway前缀;随后从最长的模块前缀开始调用importlib.import_module逐步尝试;若某个前缀无法作为模块导入,则把剩余部分当作对象属性链用getattr逐级取出; - 最终,公开 API
load_yaml(stream)(见 yaml_loader.py)读入整个配置文件并对!对象完成构造——这也是模板的 Python 启动脚本(如模板中的main.py)解析 YAML 时实际调用的入口。
从源码结构还可以看到,import_object 使用 str.partition(":") 处理模块与属性,因此除了点号路径,解析器在内部同样支持显式的 module:attribute 写法;不过文档化且最常用的仍是点号全限定名。另外,标签对应的对象必须是可调用的:若后面跟一个空节点则触发调用(构造实例),若对象本身不可调用则会抛出 constructor ... is not callable 的 YAML 解析错误。
空 mapping {} 与"对象即值"的区别
熟悉 mapping tag 语义还需要区分两种写法(详见 30.configure-yaml.md):
cache_strategy: !pw.udfs.DefaultCache {}
!pw.udfs.DefaultCache {} 表示传入空 mapping,即调用 DefaultCache() 得到一个新实例;而若直接引用某个不可调用的值(例如枚举成员),如 metric: !pw.indexing.BruteForceKnnMetricKind.COS,解析器会把对象原样作为值使用,不进行调用。自定义组件完全复用同一套规则。
实战一:给输入文档补充连接器无法生成的元数据
假设你正在做基金数据的 RAG,数据集中每个文件名遵循 {ISIN}_{基金名}_{货币}.pdf 的格式,例如 US123456789_Fund_Alpha_USD.pdf。虽然文件名本身会出现在连接器产出的 metadata 中,但直接过滤整个字符串并不方便。你希望在元数据里额外拆出 isin 与 currency 两个字段,方便后续按基金代码或币种筛选——这类字段无法由连接器自动生成,正是自定义组件的典型用武之地。
第一步:在 utils.py 中编写处理函数
先在同一目录下新建 utils.py,写入两个函数:
import pathway as pw
@pw.udf
def add_isin_and_currency(metadata: pw.Json) -> dict:
metadata_dict = metadata.as_dict()
path = metadata_dict["path"]
filename = path.split('/')[-1]
isin = filename.split('_')[0]
currency = filename.split('_')[-1][:-4]
metadata_dict["isin"] = isin
metadata_dict["currency"] = currency
return metadata_dict
def augment_metadata(sources: list[pw.Table]) -> list[pw.Table]:
return [t.with_columns(_metadata=add_isin_and_currency(pw.this._metadata)) for t in sources]
两个函数各司其职:
add_isin_and_currency是一个被@pw.udf装饰的标量 UDF,输入是文件系统连接器产生的_metadata列(类型为pw.Json)。它先as_dict()取出元数据字典,再从path中截出文件名,按下划线切分,得到isin(第一段)与currency(末段去掉.pdf扩展名),写回字典返回。UDF 机制的实现位于 python/pathway/internals/udfs/init.py:udf装饰器与UDF基类负责把普通 Python 函数包装为可作用于 Pathway 列表达式的算子,并把类型注解(如pw.Json、-> dict)转化为 Pathway 类型信息。augment_metadata则是面向 Table 集合的批处理函数,专门用于 YAML 编排:接收一个sources列表(每个元素是一张由连接器产生的表),对其中每一张表执行with_columns,用 UDF 重写_metadata列。pw.this._metadata表示当前列,这样模板在传给下游时,文档表仍保留全部原有列,只是元数据被增强了。
第二步:在 app.yaml 中接线
把 augment_metadata 接到模板数据流中:先让 $sources 从本地文件系统读取数据(开启 with_metadata: true 才会生成 _metadata 列),随后经自定义函数加工得到 $sources_with_metadata,再交给 $document_store:
$sources:
# File System connector, reading data locally.
- !pw.io.fs.read
path: data
format: binary
with_metadata: true
$sources_with_metadata: !utils.augment_metadata
sources: $sources
$document_store: !pw.xpacks.llm.document_store.DocumentStore
docs: $sources_with_metadata
parser: $parser
splitter: $splitter
retriever_factory: $retriever_factory
关键点在于 $sources_with_metadata 这一行:!utils.augment_metadata 会以 sources=$sources(即上面连接器产出的表列表)为关键字参数调用该函数,函数返回的新表列表再作为 DocumentStore 的 docs 输入。这里的模式可以推广为对整条上游数据流做任意 Python 级变换——只要函数签名是"列表进、列表出",即可作为管道中的中转节点,而不必修改模板本体。
对应到源码,DocumentStore 构造函数接受 docs: pw.Table | Iterable[pw.Table] 作为首个参数(见 document_store.py),因此传入"经过增强的表列表"完全合法。模板的运行入口代码与说明可参考本仓库中的对应项目目录 examples/projects/question-answering-rag(其 main.py 演示了从 YAML 启动这类 RAG 应用的方式)。
实战二:为 DocumentStore 定制 Parser(解析 Kafka JSON 消息)
第二个典型场景是:pw.xpacks.llm 自带的各种 Parser 不满足你的输入格式。例如,你有一条 Kafka 消息流,每条消息是一个 JSON dict,其中 "content" 键存放需要被索引的文本;而 DocumentStore 默认的 Parser 无法从这种结构里取出正文。此时你可以写一个专属 Parser 并把模板中默认的 parser 替换掉。
第一步:实现类式 UDF 作为 Parser
在 json_parser.py 中,用继承 pw.UDF 的方式实现:
import pathway as pw
import json
class JsonParser(pw.UDF):
def __wrapped__(self, message: str) -> list[tuple[str, dict]]:
message_dict = json.loads(message)
return [(message_dict["content"], {})]
这里有两个约定需要理解:
- 子类化
pw.UDF时只需实现__wrapped__方法(见 python/pathway/internals/udfs/init.py 中的说明与示例),UDF.__init__会自动包装它,使其成为可应用到 Pathway 表上的算子; - Parser 的返回类型被约定为
list[tuple[str, dict]],即"文本块 + 该块元数据"的列表。这与 LLM xpack 中内置 Parser 的契约完全一致——例如Utf8Parser.__wrapped__返回list[tuple[str, dict]](见 parsers.py),UnstructuredParser也遵循同样的返回结构(见 parsers.py 及其_chunk方法)。本例中JsonParser对每条消息只产出一个(text, metadata)对:正文取自message_dict["content"],元数据为空字典。你完全可以在此基础上向元数据里填充 Kafka 分区号、时间戳等信息。
第二步:在 YAML 中完成替换
$sources:
- !pw.io.kafka.read
# Kafka configuration, check the pathway.io.kafka API docs for needed arguments
format: plaintext
$parser: !json_parser.JsonParser {}
$document_store: !pw.xpacks.llm.document_store.DocumentStore
docs: $sources
parser: $parser
splitter: $splitter
retriever_factory: $retriever_factory
注意 $parser: !json_parser.JsonParser {} 末尾的 {}:因为 JsonParser 构造不需要参数,加上空 mapping 才能让加载器在解析阶段调用构造函数得到一个实例。这与内置 Parser 的用法一致——在 30.configure-yaml.md 的完整示例中,$parser: !pw.xpacks.llm.parsers.UnstructuredParser(无参类引用)同样可被 DocumentStore 的 parser 参数接受。从 document_store.py 的签名可以看到,parser 参数被标注为 Callable[[bytes], list[tuple[str, dict]]] | pw.UDF | None,同时兼容普通函数与 UDF 类,这就是自定义 Parser 能够"即插即用"的接口保证。
替换完成后,DocumentStore 内部在对每条 Kafka 消息执行解析(对应源码中 self.parsed_docs = ...parser(...) 的处理管线)时,就会调用你的 JsonParser,后续的 splitter 与检索索引构建流程保持不变。
实战三:在 YAML 中直接使用外部库对象(以 pandas 为例)
自定义组件不限于自研代码,任何已安装的第三方库对象同样可以出现在 ! 标签中。此时只需要注意命名规则:必须使用库的完整模块名。比如 pandas 要以 pandas 而不是惯常的别名 pd 来引用——Pathway 的 YAML 解析器唯一认可的缩写是 pathway 的 pw。
一个常见组合是先用外部库构造数据,再把它交给 Pathway 的调试工具建表:
$dataframe: !pandas.DataFrame
"data": ["foo", "bar", "baz"]
$sources:
- !pw.debug.table_from_pandas
df: $dataframe
这里的 !pandas.DataFrame 会被解析为调用 pandas.DataFrame(**{"data": [...]}),即 pandas.DataFrame(data=["foo", "bar", "baz"]);随后 $sources 通过 pw.debug.table_from_pandas 把该 DataFrame 包装为一张 Pathway 表,供下游算子使用。
为什么不需要手动 import?
与直觉相反,使用外部库对象时,你无需在 app.py 或其他任何地方先行 import pandas——pw.load_yaml 在解析 YAML 时已经替你完成了导入。这正是前述 import_object(yaml_loader.py)的职责:它按路径前缀逐级 importlib.import_module,因此 pandas.DataFrame 中的 pandas 会被当作模块自动加载,DataFrame 再以属性方式取出。也正因如此,模块必须位于当前 Python 环境的可导入路径(sys.path)上;对自研文件来说,与 app.py 放在同一目录、从该目录启动应用通常即可满足条件。
这一特性让 YAML 配置拥有完整的表达能力:你可以在模板配置中组合任意库对象(例如 numpy、datetime 等)作为组件的参数,而保持 Python 端零改动。
小结:自定义组件接入的规则清单
综合以上三类示例,可以把"在 Pathway Templates 的 YAML 中使用自定义组件"归纳为以下要点:
| 规则 | 说明 | 示例 |
|---|---|---|
| 全限定名引用 | ! 后写 模块.对象,不写模块名则默认为 builtins |
!utils.augment_metadata、!str |
pw 为唯一缩写 |
pw 自动映射到 pathway 包;第三方库必须写全名(如 pandas 而非 pd) |
!pw.xpacks.llm.llms.OpenAIChat |
| 模块名 = 文件名 | 自研代码以其所在 .py 文件名为模块名,并保证可被导入 |
utils.py → !utils.foo |
| mapping 触发调用 | 标签后跟 mapping 时,以 mapping 作为关键字参数调用对象 | $sources_with_metadata: !utils.augment_metadata\n sources: $sources |
空构造用 {} |
无参可调用对象想触发构造需显式传 {} |
!json_parser.JsonParser {} |
| 免手动 import | pw.load_yaml 通过 importlib 自动完成模块加载(见 yaml_loader.py) |
!pandas.DataFrame |
| 组件接口需匹配约定 | 如替换 Parser 须遵循 list[tuple[str, dict]] 返回约定;表级变换须"列表进、列表出" |
DocumentStore 的 parser/docs 参数 |
实际接入模板时,通常的做法是:把自定义 .py 文件放在与 app.yaml/app.py 相同的目录中,在 YAML 里声明 $自定义模块.对象 节点,并通过 $ 变量与模板其余节点(如 $sources、$parser、$retriever_factory)正确连线。这样,你既可以在 YAML 层完成绝大部分定制,又能在 pathway 库没有覆盖到或模板需要差异化处理的地方,无缝注入自己的 Python 实现。配置语法与更多内嵌组件的组合方式,可继续参考 30.configure-yaml.md 与 RAG 配置示例;本文涉及的 YAML 加载器、UDF 基类与文档存储实现的对应源码也可在 yaml_loader.py、udfs/init.py 与 document_store.py 中进一步研读,相关解析行为在 test_yaml.py 中亦有大量可运行的测试用例佐证。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0627
Hy4-previewHy4 preview 是由腾讯混元团队研发的新一代混合专家(MoE)旗舰模型。模型总参数量 770B,每个 token 激活 49B,主干共包含78层,第一层采用标准 FFN,其余 77 层均为 MoE 结构,每层包含 256 个路由专家与 1 个共享专家,每个 token 激活 top-8 路由专家及共享专家。主干之外原生内置 1 层 MTP(总参数量 10B,激活 0.7B)以支持投机解码。Python00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
GLM-5.3-FlashGLM-5.3-Flash (320B-A18B),是GLM-5系列的首个原生多模态模型。320B总参数,能力超过GLM-5.2Jinja00
Spark-X2.5-4BSpark-X2.5-4B 旨在让强大的 AI 更实用、更高效、更易获得。在广泛日常任务中表现强劲,涵盖对话、写作、翻译、推理、编码、工具调用以及智能体工作流,并在同等规模的开源模型中取得领先成绩。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00