首页
/ 在 Pathway Templates 的 YAML 配置中使用自定义组件:mapping tags 导入机制与实战示例

在 Pathway Templates 的 YAML 配置中使用自定义组件:mapping tags 导入机制与实战示例

2026-09-07 13:10:04作者:殷蕙予

本篇技术指南基于 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.pyapp.yaml 同目录下的 utils.py 中:

.
├── app.py
├── app.yaml
└── utils.py

那么 YAML 中就以 !utils.foo 来引用它——此时无需在任何 Python 文件中手动 import utils,加载 YAML 的过程会自动完成模块导入。

底层原理:pw.load_yamlimport_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 中,但直接过滤整个字符串并不方便。你希望在元数据里额外拆出 isincurrency 两个字段,方便后续按基金代码或币种筛选——这类字段无法由连接器自动生成,正是自定义组件的典型用武之地。

第一步:在 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.pyudf 装饰器与 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(即上面连接器产出的表列表)为关键字参数调用该函数,函数返回的新表列表再作为 DocumentStoredocs 输入。这里的模式可以推广为对整条上游数据流做任意 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(无参类引用)同样可被 DocumentStoreparser 参数接受。从 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 解析器唯一认可的缩写是 pathwaypw

一个常见组合是先用外部库构造数据,再把它交给 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_objectyaml_loader.py)的职责:它按路径前缀逐级 importlib.import_module,因此 pandas.DataFrame 中的 pandas 会被当作模块自动加载,DataFrame 再以属性方式取出。也正因如此,模块必须位于当前 Python 环境的可导入路径(sys.path)上;对自研文件来说,与 app.py 放在同一目录、从该目录启动应用通常即可满足条件。

这一特性让 YAML 配置拥有完整的表达能力:你可以在模板配置中组合任意库对象(例如 numpydatetime 等)作为组件的参数,而保持 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]] 返回约定;表级变换须"列表进、列表出" DocumentStoreparser/docs 参数

实际接入模板时,通常的做法是:把自定义 .py 文件放在与 app.yaml/app.py 相同的目录中,在 YAML 里声明 $自定义模块.对象 节点,并通过 $ 变量与模板其余节点(如 $sources$parser$retriever_factory)正确连线。这样,你既可以在 YAML 层完成绝大部分定制,又能在 pathway 库没有覆盖到或模板需要差异化处理的地方,无缝注入自己的 Python 实现。配置语法与更多内嵌组件的组合方式,可继续参考 30.configure-yaml.mdRAG 配置示例;本文涉及的 YAML 加载器、UDF 基类与文档存储实现的对应源码也可在 yaml_loader.pyudfs/init.pydocument_store.py 中进一步研读,相关解析行为在 test_yaml.py 中亦有大量可运行的测试用例佐证。

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