Pathway EL Pipeline 模板详解:基于 YAML 声明式配置实现无代码数据抽取与装载
导读
Pathway 是面向流处理、实时分析与 LLM/RAG 管线的数据框架,其 EL(Extract-Load)Pipeline 模板让你在不编写任何 Python 逻辑的前提下,仅通过定制一个 app.yaml 文件,即可把来自不同数据源的数据(Extract)抽取并写入任意目标 Sink(Load)。本文以仓库中的 el-pipeline 模板 为主体,逐步拆解其工程结构、YAML 声明式配置语法、JSON→CSV 与 Kafka→PostgreSQL 两套完整配置,并结合 YAML 解析器实现 与 运行时入口 的源码,说明这些配置被加载与执行的底层原理。读完本文,你将能直接复制该模板,通过改 YAML 快速搭建自己的 EL 数据通道。
项目结构:两文件驱动的数据管线
EL Pipeline 模板的工程极其精简,只包含两个文件(见 examples/templates/el-pipeline):
app.py:基于 Pathway Live Data Framework 编写的应用代码,用 Python 加载 YAML 并启动运行;app.yaml:管线配置,声明数据源(data sources)、数据汇(data sinks)、Schema 与持久化(persistence)等全部设置。
先看应用代码 app.py,核心逻辑只有十几行:
import logging
import pathway as pw
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s %(name)s %(levelname)s %(message)s",
datefmt="%Y-%m-%d %H:%M:%S",
)
def el_pipeline():
with open("app.yaml") as f:
config = pw.load_yaml(f)
persistence_config = config.get("persistence_config")
pw.run(persistence_config=persistence_config)
if __name__ == "__main__":
el_pipeline()
整个管线的"业务逻辑"全部沉淀在配置里:程序先读取同目录的 app.yaml,调用 pw.load_yaml 把声明式配置解析为真实的 Pathway 表与连接器对象,再取出可选的 persistence_config 传给 pw.run 启动计算图。这意味着更换数据源或目标库,通常只需要改 YAML,不需要动 Python 代码。
应用代码对应的源码实现
pw.load_yaml的实现位于 python/pathway/internals/yaml_loader.py,它被 python/pathway/init.py 导出到顶层命名空间(pw.load_yaml)。pw.run的实现位于 python/pathway/internals/run.py,其签名中persistence_config的默认值是None,类型为PersistenceConfig | None,也就是说不声明持久化时完全允许缺省。
从 run.py 的签名可以看出,pw.run 还支持 debug、monitoring_level、default_logging、terminate_on_error、runtime_typechecking、max_expression_batch_size 等参数,这些在后续"监控与日志"小节会用到。
环境要求与安装
模板 README 列出的前置条件如下:
- Python 3.8 或更高版本;
- Git;
- Docker(可选,仅当使用容器化方式运行时需要)。
标准安装流程是克隆 Pathway 仓库后进入模板目录:
git clone https://github.com/pathwaycom/pathway.git
cd examples/templates/el-pipeline/
进入目录后即可看到 app.py 与 app.yaml 两个文件,按照下一节的配置说明修改 app.yaml,再执行 python app.py 即可启动管线。
配置核心:YAML 声明式语法是如何被执行的
README 强调 app.yaml 使用声明式 YAML 格式定义数据源、数据汇与其他设置。要理解这些配置项,先要明白解析器做了什么。
Pathway 的 YAML 加载逻辑集中在 yaml_loader.py,它注册了两类特殊语法:
$前缀变量引用:任何以$开头的标量都会被解析成Variable(见 yaml_loader.py 的隐式解析器re.compile(r"\$.*"))。Resolver会在同层 YAML 映射中查找该变量,形成"先定义、后引用"的复用机制。!前缀的路径构造标签:通过add_multi_constructor("!", ...)注册(见 yaml_loader.py),把类似!pw.io.fs.read的标签转换为 Python 对象。转换逻辑import_object会把pw.前缀映射为pathway包,最终调用对应的连接器构造函数,YAML 键值对作为关键字参数传入(见 yaml_loader.py)。这解释了为什么 YAML 里的写法与 Python 中pw.io.fs.read(path=..., format=..., name=...)的 API 一一对应。
还有一个值得注意的机制:环境变量替换。解析器遇到尚未定义的全大写变量名时,会去 os.environ 中查找同名环境变量作为取值(见 yaml_loader.py)。这正是模板在 Kafka、PostgreSQL 配置中使用 $KAFKA_HOSTNAME、$DB_USER 这类大写占位符的底层原因——运行时它们会被真实环境变量填充,避免把凭据写死在配置文件里。
此外,解析器会在作用域退出时检查是否有未使用的 $ 变量并发出告警(check_unused_variables),帮助你及早发现配置笔误。
声明式配置核心要素
下面逐项拆解模板中使用的四类配置要素。
声明一个数据源(Extract)
数据源通过 Pathway Live Data Framework 的**输入连接器(input connector)**声明。例如文件系统输入连接器:
$source: !pw.io.fs.read
path: ./input_data/
format: binary
name: input_connector
要点解析:
!pw.io.fs.read对应pw.io.fs.read连接器,最终解析为真实的对象;path指定读取目录,format声明数据格式(如binary、json、csv、raw等),name是该连接器的标识名;- 键名
$source以$开头,表示它只是一个内部变量而非输出表。后续 Sink 通过table: $source引用它。
声明一个数据汇(Load)
数据汇由**输出连接器(output connector)**声明。例如把数据写成 CSV 文件:
output: !pw.io.csv.write
table: $source
filename: ./output.csv
name: output_connector
注意区别:这里的键名 output 不带 $ 前缀,因为 pw.run 在 app.py 中只会通过 config.get("persistence_config") 读取持久化配置,其余顶层键更像一个"命名空间"——非 $ 键不会被当作引用变量处理。table: $source 把前面声明的输入表接到该输出上,从而在拓扑上构成完整的 Extract→Load 链路。
连接器的全貌
Pathway Live Data Framework 提供了极为丰富的连接器。以 python/pathway/io/init.py 导出的模块为准,输入/输出侧支持 fs、csv、jsonlines、plaintext、kafka、redpanda、pulsar、mqtt、nats、rabbitmq、pubsub、kinesis、dynamodb、debezium、logstash、postgres、mysql、mssql、sqlite、clickhouse、duckdb、mongodb、elasticsearch、deltalake、iceberg、s3、minio、gdrive、http、python 自定义连接器,以及面向向量库的 chroma、pinecone、qdrant、weaviate、milvus、leann 等。模板中只需把 !pw.io.fs.read/!pw.io.csv.write 换成目标连接器(注意签名差异),并在 YAML 中补上对应参数即可,管线代码无需变动。
定义 Schema
通过 pw.schema_from_types 以"类型即列"的方式声明表结构:
$schema: !pw.schema_from_types
colA: str
colB: int
colC: float
schema_from_types 会依据 Python 类型注解生成一个 Pathway Schema(列名→类型映射),支持 str、int、float 等基础类型,也可扩展更复杂的类型。然后把 Schema 应用到输入连接器上,让连接器知道如何解析与类型校验上游数据:
$source: !pw.io.csv.read
path: ./input_data/
schema: $schema
name: input_connector
对带格式的输入(JSON、CSV 等),声明 Schema 是保证列结构与类型一致、从而让下游直接引用列的关键步骤。
配置持久化(Persistence)
持久化用于保存计算状态,以便从故障中恢复:
persistence_config: !pw.persistence.Config
backend: !pw.persistence.Backend.filesystem
path: ./persistence_storage/
它对应 pw.run 的 persistence_config 形参:app.py 中正是通过 config.get("persistence_config") 把这个键的值取出并传给 pw.run(run.py)。当流式输入需要"断点续跑"——例如 Kafka 消费位点、中间状态表——时,启用该配置能显著提升可靠性。使用文件系统后端时,会把状态持久化到指定目录。需要提醒的是,README 中明确指出使用 persistence 等高级连接器/功能需要 Pathway Live Data Framework 许可证(详见下文许可证小节),仓库中还配套了专门的持久化测试集,例如 python/pathway/tests/test_persistence.py 对状态恢复场景进行了覆盖验证。
示例配置:从 JSON 到 CSV
模板在 README 中给出了一套完整的 JSON→CSV 最小示例,可直接替换进 app.yaml 使用:
# This YAML configuration file is used to set up and configure the EL pipeline.
# It defines the data sources, the data sinks, and the persistence configuration.
# Structure of the data sources using Pathway schema.
$schema: !pw.schema_from_types
colA: str
colB: int
colC: float
# Define the data source using the schema.
$source: !pw.io.fs.read
path: ./input_data/
format: json
schema: $schema
name: input_connector
# Output the data in the data sink of your choice.
output: !pw.io.csv.write
table: $source
filename: ./output.csv
name: output_connector
# Uncomment to use persistence on the file system.
# persistence_config: !pw.persistence.Config
# backend: !pw.persistence.Backend.filesystem
# path: ./persistence_storage/
运行方式:把 JSON 文件放入 ./input_data/,执行 python app.py,管线会持续监听该目录新出现的 JSON 数据,解析后追加写入 ./output.csv。需要持久化时,去掉末尾三行注释即可。
默认配置实战:Kafka → PostgreSQL
仓库中实际交付的 app.yaml 默认配置是 从 Kafka 读取消息并写入 PostgreSQL 的完整链路,比 JSON→CSV 更贴近生产场景,涉及三个子配置块与一组环境变量。
Kafka 输入连接器
首先定义消息表结构(包含 date 与 message 两列),再通过 $rdkafka_settings 变量聚合 librdkafka 级参数,最后声明 Kafka 输入:
$InputStreamSchema: !pw.schema_from_types
date: str
message: str
$rdkafka_settings:
"bootstrap.servers": $KAFKA_HOSTNAME
"security.protocol": "plaintext"
"group.id": $KAFKA_GROUP_ID
"session.timeout.ms": "6000"
"auto.offset.reset": "earliest"
$kafka_source: !pw.io.kafka.read
rdkafka_settings: $rdkafka_settings
topic: $KAFKA_TOPIC
format: "json"
schema: $InputStreamSchema
autocommit_duration_ms: 100
name: input_kafka_connector
参数说明:
bootstrap.servers:Kafka 集群地址;security.protocol:本例为plaintext(明文),生产环境可按需改为ssl/sasl_plaintext等并补充对应凭据参数;group.id:消费组 ID,由KAFKA_GROUP_ID注入;session.timeout.ms:会话超时(6000ms),保持默认即可;auto.offset.reset:无已提交位点时从最早消息开始消费;format: "json":消息体按 JSON 解析;autocommit_duration_ms: 100:Kafka 消费位点自动提交的间隔(毫秒),Pathway 会周期性地把已处理 offset 提交回 broker,配合持久化实现 exactly-once 语义的基础。
PostgreSQL 输出连接器
$postgres_settings 用映射变量集中存放数据库连接参数,随后声明目标表并挂接输出:
$postgres_settings:
"host": $DB_HOSTNAME
"port": $DB_PORT
"dbname": $DB_NAME
"user": $DB_USER
"password": $DB_PASSWORD
$table_name: "messages_table"
output: !pw.io.postgres.write
table: $kafka_source
postgres_settings: $postgres_settings
table_name: $table_name
name: output_postgres_connector
!pw.io.postgres.write 会把上游表 $kafka_source 的每一行写入 PostgreSQL 的 messages_table 表。注意连接参数里的端口 $DB_PORT、密码 $DB_PASSWORD 均为运行时环境变量,模板仓库不会也不能保存你的真实凭据。
持久化配置(默认注释)
默认 app.yaml 将持久化以注释形式给出,需要"断点续传"能力时取消注释:
persistence_config: !pw.persistence.Config
backend: !pw.persistence.Backend.filesystem
path: ./persistence_storage/
必需的环境变量
运行这套默认配置前,需要预先导出以下 8 个环境变量:
| 变量 | 含义 |
|---|---|
KAFKA_HOSTNAME |
Kafka 服务器主机名 |
KAFKA_GROUP_ID |
Kafka 消费组 ID |
KAFKA_TOPIC |
要消费的 Kafka topic |
DB_HOSTNAME |
PostgreSQL 主机名 |
DB_PORT |
PostgreSQL 端口 |
DB_NAME |
PostgreSQL 数据库名 |
DB_USER |
PostgreSQL 用户名 |
DB_PASSWORD |
PostgreSQL 密码 |
在 Linux/macOS 上可一次性导出后运行:
export KAFKA_HOSTNAME=localhost:9092
export KAFKA_GROUP_ID=my_group
export KAFKA_TOPIC=my_topic
export DB_HOSTNAME=localhost
export DB_PORT=5432
export DB_NAME=mydb
export DB_USER=myuser
export DB_PASSWORD=mypassword
python app.py
因为 YAML 解析器的"全大写变量名回退读取 os.environ"机制(yaml_loader.py),$KAFKA_HOSTNAME、$DB_USER 等占位符会在解析阶段被真实环境变量值替换,所以配置文件本身不携带任何敏感信息,可以在仓库与 CI 中安全地版本化。
运行管线
- 直接用 Python 运行:在
examples/templates/el-pipeline/目录下执行python app.py,app.py会用 logging.basicConfig 初始化带时间戳的 INFO 级日志,然后加载 YAML 并启动计算图; - 容器化运行:可选。若选择 Docker 方式,需要确保构建环境满足前述前置条件中的 Docker 依赖。
pw.run 被调用后,会基于全局计算图构造 GraphRunner 并执行所有输出(见 run.py)。默认 monitoring_level 为 AUTO,Pathway 会根据输出是否为交互式终端在 NONE 与 IN_OUT 之间自动选择监控粒度。
Pathway Live Data Framework 许可证与 LICENSE KEY
README 明确提醒:当使用高级连接器或高级功能(例如 SharePoint 连接器、persistence 持久化等)时,项目需要 Pathway Live Data Framework 许可证。基础 EL 用法与仓库内 MIT 许可的示例工程不受此影响,但如果要在商业环境中解锁全部连接器与容灾能力,请按官方许可指南评估。
设置许可证只需把 KEY 导出为环境变量 PATHWAY_LICENSE_KEY:
export PATHWAY_LICENSE_KEY=your_pathway_key
然后照常 python app.py 运行即可。仓库根目录的 LICENSE.txt 说明了底层分发条款(允许非商业使用与大部分商业用途免费使用,代码在特定年限后自动转为 Apache 2.0 开源许可),本模板目录中的代码属于仓库的 MIT 许可示例工程,可在此基础上继续扩展。
监控与日志
- 访问日志:
app.py通过logging.basicConfig把日志级别设为 INFO,并以%(asctime)s %(name)s %(levelname)s %(message)s的统一格式输出到标准错误流,便于 grep 与对接采集器。若想完全接管日志格式,可调用pw.run(default_logging=False)(run.py)以禁止 Pathway 设置自己的日志处理器,再自行配置 handler。 - 监控管线性能:
pw.run的monitoring_level参数(NONE/IN_OUT/ALL,默认AUTO)控制性能统计的详细程度;debug=True可启用table.debug()输出用于排查。
常见问题排查(Troubleshooting)
模板 README 给出了三条高价值的排查主线,结合配置机制可进一步细化:
- 确认所有环境变量均已正确设置。
$大写占位符在解析期直接从环境变量取值(见 yaml_loader.py),缺失会抛出variable ... is not defined的KeyError,运行前可用echo $KAFKA_HOSTNAME逐一核对;另外,若定义了变量却从未被引用,解析器会给出unused YAML variable告警。 - 核对 YAML 中的路径与凭据准确无误。
!pw.io.fs.read的path、!pw.io.csv.write的filename、PostgreSQL 连接参数都必须是真实可达的值;目录不存在、文件无读权限或数据库连不上都会在启动阶段暴露。 - 检查日志中的详细报错。app.py 已将日志切到 INFO 级,
pw.run默认开启日志(除非显式default_logging=False),报错通常带有阶段与堆栈信息,可据此判断是 YAML 解析失败、连接器初始化失败还是运行期数据问题。
小结
这个 EL Pipeline 模板展示了 Pathway 声明式配置的核心能力:! 标签把 YAML 节点映射为真实连接器构造调用,$ 变量提供配置复用与顶层拓扑组织,全大写变量自动回退到环境变量实现敏感信息外置,而 pw.run(persistence_config=...) 把持久化配置无缝接入运行时。借助 yaml_loader.py 的解析机制与 run.py 的运行时参数,你可以仅靠编辑 app.yaml 在文件系统、Kafka、各类数据库与向量库之间快速搭建可运行的抽取-装载管线——先跑通模板自带的 Kafka→PostgreSQL 默认配置,再按 JSON→CSV 的最小示例改造源与目标,即可把这套模板复用到你自己的数据集成场景中。
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