Pathway 连接器体系详解:流式与静态双模式数据接入、更新语义与持久化恢复
本文基于 Pathway 官方用户指南中的连接器总览文档,系统讲解 Pathway Live Data Framework 中连接器(Connector)的定位、全量可用连接器清单、流式/静态两种数据模式的差异与更新语义,并结合仓库中的 Python API 源码说明关键参数(如 mode、autocommit_duration_ms、max_backlog_size)的实际含义,帮助你为实时 ETL、RAG 与流分析场景选择合适的输入/输出连接器,并正确配置持久化。
一、什么是连接器:Pathway 的“数据入口与出口”
Pathway 是一个 Python 编写的流处理 ETL 框架,其核心运行时用 Rust 实现。要使用 Pathway Live Data Framework,首先要解决的是“如何访问要处理的数据”,而这一职责正是由连接器承担的:
- 输入连接器(Input connectors):把外部数据源的数据读入框架,并在数据变化时自动把更新推入计算图;
- 输出连接器(Output connectors):把框架内计算得到的结果变化写出到外部系统。
在 python/pathway/io/ 包下可以看到完整的连接器模块布局,如 csv、kafka、s3、postgres、mongodb、elasticsearch、deltalake、iceberg、nats、mqtt、milvus、pinecone、qdrant、weaviate、duckdb、clickhouse、slack、pyfilesystem 等,覆盖了任务队列、文件系统、关系型/文档数据库、湖仓格式与向量库等主要数据源类型。
在深入各连接器之前,官方文档要求读者先理解 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 计算引擎驱动,连接器只负责把数据事件喂入图中。
两个需要特别注意的流式语义:
- 更新以 commit 为触发单位,保证原子性。 实践中,更新由 commit 触发,每次 commit 保证一批更新的原子性。输入连接器普遍提供
autocommit_duration_ms参数控制 commit 间隔:在 python/pathway/io/csv/init.py 的read()签名中可以看到其默认值为 1500 毫秒,即每 1500 ms 内收到的更新会被打包成一次提交推入计算图。 - 计算是“永不结束”的。 因为流式数据在概念上是无限的,程序会一直运行直到进程被终止——这是框架的正常行为,不是死循环 bug。
输出的是“变化”而非全量
流式模式下,输出连接器是访问计算结果的唯一途径。但注意:输出的不是完整表,而是每次更新产生的变化(delta)。每条变化表示为一行,包含:
- 表本身各列的字段,即被修改的值;
time列:更新的逻辑时间,每次新 commit 递增;diff列:表示该更新是“新增”还是“删除”,只取两个值——1表示新增,-1表示删除。
当一个字段值从旧值更新为新值时,会表示为两行:一行删除旧值(diff = -1),一行添加新值(diff = 1)。这要求下游系统(如写入数据库或消息队列)按 upsert/删除语义消费这些变化,理解这一点对正确接入输出连接器至关重要。
想体验完整的实时流式应用,仓库提供了两个入门模板:基于 CSV 输入的首个实时应用与基于 Kafka 的线性回归模板。
反压参数 max_backlog_size
从 python/pathway/io/csv/init.py 的 read() 签名可以看到,输入连接器还普遍支持 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_ms、max_backlog_size 等参数把吞吐、延迟与内存占用调到适合自身负载的平衡点。
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 StartedRust0623
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