首页
/ Pathway 连接器体系详解:流式与静态双模式数据接入、更新语义与持久化恢复

Pathway 连接器体系详解:流式与静态双模式数据接入、更新语义与持久化恢复

2026-09-05 15:01:36作者:明树来

本文基于 Pathway 官方用户指南中的连接器总览文档,系统讲解 Pathway Live Data Framework 中连接器(Connector)的定位、全量可用连接器清单、流式/静态两种数据模式的差异与更新语义,并结合仓库中的 Python API 源码说明关键参数(如 modeautocommit_duration_msmax_backlog_size)的实际含义,帮助你为实时 ETL、RAG 与流分析场景选择合适的输入/输出连接器,并正确配置持久化。

一、什么是连接器:Pathway 的“数据入口与出口”

Pathway 是一个 Python 编写的流处理 ETL 框架,其核心运行时用 Rust 实现。要使用 Pathway Live Data Framework,首先要解决的是“如何访问要处理的数据”,而这一职责正是由连接器承担的:

  • 输入连接器(Input connectors):把外部数据源的数据读入框架,并在数据变化时自动把更新推入计算图;
  • 输出连接器(Output connectors):把框架内计算得到的结果变化写出到外部系统。

python/pathway/io/ 包下可以看到完整的连接器模块布局,如 csvkafkas3postgresmongodbelasticsearchdeltalakeicebergnatsmqttmilvuspineconeqdrantweaviateduckdbclickhouseslackpyfilesystem 等,覆盖了任务队列、文件系统、关系型/文档数据库、湖仓格式与向量库等主要数据源类型。

在深入各连接器之前,官方文档要求读者先理解 Pathway 的**流式模式(streaming)与静态模式(static)**两种数据模式,因为连接器是按模式划分且不可混用的,可参考 流式与静态模式文档。

二、可用连接器全景:按“输入/输出 × 流式/静态”分类

下面完整继承原文档中的连接器清单(原始版本以站点链接表格给出),并按“输入/输出连接器”与“流式/静态模式”两个维度重新组织为 Markdown 表格,方便检索。表中“流式”列表示该连接器可作为流式输入/输出使用,“静态”列表示支持一次性批量读写。

输入连接器:

数据源 流式模式 静态模式
Airbyte 支持
Amazon S3 支持 支持
CSV 支持 支持
Debezium 支持
Delta Lake 支持 支持
Elastic Search 支持
File System(本地/对象存储文件系统) 支持 支持
Google Drive 支持 支持
HTTP 支持
Iceberg 支持
JSON Lines 支持 支持
Kafka 支持 支持
Kinesis 支持
MinIO 支持 支持
MongoDB(经 Debezium) 支持
MongoDB(oplog 复制 / MongoDB Atlas) 支持
MQTT 支持
MS SQL Server 支持
NATS 支持
NeonDB 支持
Plain text 支持
PostgreSQL(通过读取 WAL) 支持
Pulsar 支持
Python(自定义 Python 输入连接器) 支持
RabbitMQ 支持
Redpanda 支持 支持
SharePoint 支持
SQLite 支持
Markdown / Pandas(调试用数据构造) 支持

输出连接器:

目标 流式模式 静态模式
CSV 支持 支持
BigQuery 支持
Chroma 支持
ClickHouse 支持
Delta Lake 支持
DuckDB 支持
DynamoDB 支持
Elastic Search 支持
File System 支持
Google PubSub 支持
HTTP 支持
Iceberg 支持
JSON Lines 支持
Kafka 支持
Kinesis 支持
Logstash 支持
Milvus 支持
MongoDB / MongoDB Atlas 支持
MQTT 支持
MS SQL Server 支持
MySQL 支持
NATS 支持
NeonDB 支持
pgvector 支持
Pinecone 支持
PostgreSQL 支持
Pulsar 支持
Qdrant 支持
QuestDB 支持
RabbitMQ 支持
Redpanda 支持
Slack(告警) 支持
SQLite 支持
Weaviate 支持
pw.debug.compute_and_print / compute_and_print_update_stream(调试输出) 支持

可以看到:几乎所有文件系统类连接器(CSV、JSON Lines、Kafka/Redpanda、S3、MinIO、Delta Lake、Google Drive)都同时支持两种模式,而大多数数据库复制与消息队列类连接器(Kinesis、Pulsar、Debezium、oplog 复制等)只提供流式输入;静态输出则非常有限(CSV 与两个 pw.debug 调试函数),因为静态模式的定位本就是“一次性批量计算”。

各连接器的教程文档集中在 docs/2.developers/4.user-guide/20.connect/99.connectors/ 目录下,例如 CSV 连接器教程数据库连接器教程Kafka 连接器教程从 Kafka 切换到 Redpanda自定义 Python 输入连接器 等。

如果表中暂时没有你需要的连接器,官方也欢迎反馈需求以推动新连接器落地。

三、流式模式的连接器:更新语义与自动传播

在流式模式下,输入连接器持续等待新的更新(new update)。每当收到更新,它会被推入计算图(dataflow),一路传播到输出连接器,由输出连接器把“结果的变化”写出。

这里的关键行为是:输入连接器创建的表,以及所有基于它构建的计算,都会在收到新更新时自动刷新。例如当目录中出现一个新的 CSV 文件时,不需要任何人工干预,下游所有计算和输出都会自动纳入这条数据——这正是 Pathway Live Data Framework 的核心卖点。从源码结构看,这种“自动重算”由 Rust 运行时上的 Differential Dataflow 计算引擎驱动,连接器只负责把数据事件喂入图中。

两个需要特别注意的流式语义:

  1. 更新以 commit 为触发单位,保证原子性。 实践中,更新由 commit 触发,每次 commit 保证一批更新的原子性。输入连接器普遍提供 autocommit_duration_ms 参数控制 commit 间隔:在 python/pathway/io/csv/init.pyread() 签名中可以看到其默认值为 1500 毫秒,即每 1500 ms 内收到的更新会被打包成一次提交推入计算图。
  2. 计算是“永不结束”的。 因为流式数据在概念上是无限的,程序会一直运行直到进程被终止——这是框架的正常行为,不是死循环 bug。

输出的是“变化”而非全量

流式模式下,输出连接器是访问计算结果的唯一途径。但注意:输出的不是完整表,而是每次更新产生的变化(delta)。每条变化表示为一行,包含:

  • 表本身各列的字段,即被修改的值;
  • time 列:更新的逻辑时间,每次新 commit 递增;
  • diff 列:表示该更新是“新增”还是“删除”,只取两个值——1 表示新增,-1 表示删除。

当一个字段值从旧值更新为新值时,会表示为两行:一行删除旧值(diff = -1),一行添加新值(diff = 1)。这要求下游系统(如写入数据库或消息队列)按 upsert/删除语义消费这些变化,理解这一点对正确接入输出连接器至关重要。

想体验完整的实时流式应用,仓库提供了两个入门模板:基于 CSV 输入的首个实时应用基于 Kafka 的线性回归模板

反压参数 max_backlog_size

python/pathway/io/csv/init.pyread() 签名可以看到,输入连接器还普遍支持 max_backlog_size 参数:它限制“同一时刻正在参与计算的条数上限”,达到上限时读取暂停,直到部分条目处理完成才恢复。官方文档提示,默认 None(不限制)在“初始大量批量数据 + 后续少量增量”的场景下可能让内存无限增长,因此对大型数据源建议显式设置。该机制在架构文档 How Pathway Live Data Framework Connectors Work 中有更深入的引擎级说明(按 mini-batch 粒度追踪在途数据量)。

四、静态模式的连接器:批量计算与调试

静态模式下,计算以批处理方式进行:一次性读入全部数据、处理、写出,不存在“更新”的概念。官方文档明确强调:该模式主要用于调试和测试

静态输出方面,除将输出表转储为 CSV 文件的 CSV 连接器外,Pathway 提供 pw.debug.compute_and_print 函数:它构建计算图、摄取全部数据,并打印图中指定的表。原文档给出的手工表静态模式示例如下:

import pathway as pw

t = pw.debug.table_from_markdown(
    """
    | name  | age
 1  | Alice | 15
 2  | Bob   | 32
 3  | Carole| 28
 4  | David | 35 """
)

pw.debug.compute_and_print(t)

运行输出(每行前的乱码前缀是 Pathway 为每行生成的唯一行标识符):

            | name   | age
^YYY4HAB... | Alice  | 15
^Z3QWT29... | Bob    | 32
^3CZ78B4... | Carole | 28
^3HN31E1... | David  | 35

table_from_markdown / table_from_pandas 等构造函数在连接器总表中被归为“静态输入”,它们让开发者无需真实数据源即可对计算图进行快速验证。

五、常见陷阱:两种模式的连接器不可混用

原文档专门用一节警示兼容性问题:流式与静态两种模式互不兼容,不能把两种模式的连接器混在同一个管线中,因为它们操作的数据本质不同(数据流 vs 静态数据)。

一个典型错误场景:你想在管线中用 pw.debug.compute_and_print(table) 检查某张表 table 的中间值是否正确,于是把该行插进两次 select 之间,然后以流式输入连接器运行程序。结果——程序会陷入死循环。原因是 compute_and_print 会等待数据全部摄取完毕才打印表;这对有限的静态数据成立,但对持续不断产生更新的流式数据永远不会成立。

因此在用静态数据/静态调试函数排查管线时,务必确认整条链路(输入、调试节点、输出)都处于静态模式。这一点也可以从 CSV 连接器源码得到印证:read()mode 参数文档字符串明确说明 "streaming" 模式会“等待指定目录的更新,跟踪文件的增删改”,而 "static" 模式“只考虑现有数据并在一次 commit 中摄取全部”(见 python/pathway/io/csv/init.py)。

六、连接器中的持久化(Persistence)

无论流式还是静态模式,连接器都可以持久化已读取的数据及部分中间计算结果,以便程序在重启后从上次终止的位置继续,而无需从头重放。典型用途:

  • 程序追加新数据后需要重跑(re-runs with added data);
  • 希望程序能“幸存”于代码崩溃(crash recovery)。

启用方式是在 pw.run 方法中指定持久化配置。若连接器开启了持久化,Pathway 会保存其辅助数据(auxiliary data),使程序可以断点续跑。持久化的使用细节(存储后端、恢复与带新数据重启)参见 持久化文档

从架构层面看,持久化要求每条记录携带“位置元数据”(offset,如 Kafka 的分区+偏移量、文件流式读取的字节游标、Delta Lake 的版本号),重启时引擎通过 seek 方法让读取器跳转到检查点位置;每个 worker 会写各自的 Write-Ahead Log,恢复时合并所有 worker 的 frontier。这些机制对所有连接器统一生效,属于框架自动提供、连接器无需自行实现的能力,详见 连接器工作原理

七、数据格式与人工数据流

格式支持

不同连接器支持不同的数据格式(CSV、JSON 等),但有一条统一约定:所有连接器都支持 binary 格式。当需要完全自定义解析逻辑时,可以直接读取二进制字节再自行处理。

用 demo 模块生成人工数据流

在真实场景中获取可用的数据流进行测试有时很困难。Pathway 提供了 demo 模块来模拟流入的数据流:可以从零开始自定义数据流,也可以基于一个 CSV 文件生成流,方便对实时处理逻辑进行实验与测试。

八、实战教程索引

原文档最后列出的连接器教程,在仓库中对应以下文档,按学习路径排序:

教程 仓库路径
CSV 连接器 csv_connectors
数据库连接器(PostgreSQL WAL 等) database-connectors
Kafka 连接器 kafka_connectors
从 Kafka 切换到 Redpanda switching-to-redpanda
Python 输入连接器 custom-python-connectors
Python 输出连接器 python-output-connectors
Google Drive 连接器 gdrive-connector

此外,仓库 ETL 模板目录 提供了一个完整的实时数据处理管线示例(欺诈/异常活动检测 + 翻滚窗口),展示了输入连接器、窗口计算与输出连接器如何组合成一条端到端的流式管线。

九、小结

  • 连接器是 Pathway 的数据边界:输入连接器把外部变化引入计算图,输出连接器把结果变化写出;二者都区分流式与静态两种模式。
  • 流式模式下计算永不结束、更新由 commit 原子化(autocommit_duration_ms 控制提交窗口,默认 1500 ms)、输出的是带 time/diff 的变化行(1 增、-1 删,值更新表现为两行)。
  • 静态模式用于一次性批量计算与调试,pw.debug.compute_and_print 是调试核心,但严禁与流式连接器混用,否则会因为“等待全量数据”而死循环。
  • 持久化通过 pw.run 的持久化配置开启,配合 offset/检查点机制实现崩溃恢复与断点续跑。
  • 选型建议:文件/对象存储类源(CSV、JSON Lines、S3、Delta Lake 等)优先使用同时支持双模式的连接器;实时消息与 CDC 类源(Kafka、Kinesis、Debezium、WAL 复制等)走流式输入;向量库类(Milvus、Pinecone、Qdrant、Weaviate、pgvector 等)作为流式输出端用于 RAG 场景。

掌握以上分类与语义后,你就可以在 python/pathway/io/ 中定位具体连接器模块,按对应教程完成接入,并用 autocommit_duration_msmax_backlog_size 等参数把吞吐、延迟与内存占用调到适合自身负载的平衡点。

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