首页
/ Pathway EL Pipeline 模板详解:基于 YAML 声明式配置实现无代码数据抽取与装载

Pathway EL Pipeline 模板详解:基于 YAML 声明式配置实现无代码数据抽取与装载

2026-09-07 17:09:32作者:瞿蔚英Wynne

导读

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 代码

应用代码对应的源码实现

run.py 的签名可以看出,pw.run 还支持 debugmonitoring_leveldefault_loggingterminate_on_errorruntime_typecheckingmax_expression_batch_size 等参数,这些在后续"监控与日志"小节会用到。

环境要求与安装

模板 README 列出的前置条件如下:

  • Python 3.8 或更高版本;
  • Git;
  • Docker(可选,仅当使用容器化方式运行时需要)。

标准安装流程是克隆 Pathway 仓库后进入模板目录:

git clone https://github.com/pathwaycom/pathway.git
cd examples/templates/el-pipeline/

进入目录后即可看到 app.pyapp.yaml 两个文件,按照下一节的配置说明修改 app.yaml,再执行 python app.py 即可启动管线。

配置核心:YAML 声明式语法是如何被执行的

README 强调 app.yaml 使用声明式 YAML 格式定义数据源、数据汇与其他设置。要理解这些配置项,先要明白解析器做了什么。

Pathway 的 YAML 加载逻辑集中在 yaml_loader.py,它注册了两类特殊语法:

  1. $ 前缀变量引用:任何以 $ 开头的标量都会被解析成 Variable(见 yaml_loader.py 的隐式解析器 re.compile(r"\$.*"))。Resolver 会在同层 YAML 映射中查找该变量,形成"先定义、后引用"的复用机制。
  2. ! 前缀的路径构造标签:通过 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 声明数据格式(如 binaryjsoncsvraw 等),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(列名→类型映射),支持 strintfloat 等基础类型,也可扩展更复杂的类型。然后把 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.runpersistence_config 形参:app.py 中正是通过 config.get("persistence_config") 把这个键的值取出并传给 pw.runrun.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 输入连接器

首先定义消息表结构(包含 datemessage 两列),再通过 $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.pyapp.py 会用 logging.basicConfig 初始化带时间戳的 INFO 级日志,然后加载 YAML 并启动计算图;
  • 容器化运行:可选。若选择 Docker 方式,需要确保构建环境满足前述前置条件中的 Docker 依赖。

pw.run 被调用后,会基于全局计算图构造 GraphRunner 并执行所有输出(见 run.py)。默认 monitoring_levelAUTO,Pathway 会根据输出是否为交互式终端在 NONEIN_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.runmonitoring_level 参数(NONE/IN_OUT/ALL,默认 AUTO)控制性能统计的详细程度;debug=True 可启用 table.debug() 输出用于排查。

常见问题排查(Troubleshooting)

模板 README 给出了三条高价值的排查主线,结合配置机制可进一步细化:

  1. 确认所有环境变量均已正确设置$ 大写占位符在解析期直接从环境变量取值(见 yaml_loader.py),缺失会抛出 variable ... is not definedKeyError,运行前可用 echo $KAFKA_HOSTNAME 逐一核对;另外,若定义了变量却从未被引用,解析器会给出 unused YAML variable 告警。
  2. 核对 YAML 中的路径与凭据准确无误!pw.io.fs.readpath!pw.io.csv.writefilename、PostgreSQL 连接参数都必须是真实可达的值;目录不存在、文件无读权限或数据库连不上都会在启动阶段暴露。
  3. 检查日志中的详细报错。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 的最小示例改造源与目标,即可把这套模板复用到你自己的数据集成场景中。

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