Pathway 自定义 Python 连接器:实现任意数据源接入流式处理
本篇指南介绍如何在 Pathway 实时数据框架中创建自定义 Python 连接器(ConnectorSubject),将静态文件、第三方库等任意数据源转换为流式输入,并深入解析 pw.io.python.read 的完整参数与底层线程/缓冲区机制。读完本篇,你可以掌握:从数据源读数据并通过 next 推入引擎缓冲区的完整模式、外部库(如 Tweepy)与 Pathway 的对接方式,以及 autocommit_duration_ms、max_backlog_size 等关键参数的作用与默认值。
核心概念:ConnectorSubject 是数据源与引擎之间的桥梁
创建自定义连接器的核心步骤是:继承 pw.io.python.ConnectorSubject 抽象类,并实现其中的 run 方法——该方法负责从数据源读取数据并将数据喂入缓冲区。
ConnectorSubject 充当数据源与 Pathway 引擎之间的桥梁,它内置了 next 方法,允许你将数据推入缓冲区。从 ConnectorSubject 的类定义 可以看到几个关键事实:
run是唯一的抽象方法(@abstractmethod),必须实现;- 该函数由 Pathway 引擎在独立线程中启动;
- 当
run函数终止后,连接器即被视为结束,Pathway 不会再等待它的新消息; - 如果连接器不删除记录,可设置类属性
deletions_enabled为False,有助于提升性能; - 若
read时设置了max_backlog_size,当事件队列达到上限时,next等方法的调用会被阻塞,直到队列大小降回限制以下。
简单场景:把静态文件变成流
假设你有一个文件 cats.jsonl,每行包含一条 JSON 数据:
{ "key": 1, "genus": "otocolobus", "epithet": "manul" }
{ "key": 2, "genus": "felis", "epithet": "catus" }
{ "key": 3, "genus": "lynx", "epithet": "lynx" }
目标:写一个连接器,把静态文件逐行转换为流。在 run 方法中读取文件,并使用 next 方法把每一行的内容发送到缓冲区:
import json
import time
import pathway as pw
class FileStreamSubject(pw.io.python.ConnectorSubject):
def run(self):
with open("cats.jsonl") as file:
for line in file:
data = json.loads(line)
self.next(**data)
time.sleep(1)
这里 json.loads(line) 把每行解析为字典,再用 **data 展开为关键字参数传给 self.next()——next 接受关键字参数,每个参数对应表的一列,所有列的值会被作为单行发送到引擎。time.sleep(1) 模拟"逐条到达"的流式节奏;处理真实源时通常不需要它。
接下来,通过 pw.io.python.read 方法消费这个 subject 并执行计算。调用时需要传入你的 subject 实例,并指定数据的 schema——即声明将被提取为列的字段及其类型:
class InputSchema(pw.Schema):
key: int = pw.column_definition(primary_key=True)
genus: str
epithet: str
table = pw.io.python.read(
FileStreamSubject(),
schema=InputSchema,
)
pw.io.csv.write(table, "output.csv")
pw.run()
注意 pw.column_definition(primary_key=True) 将 key 声明为主键:后续相同主键的数据会触发更新/删除语义,而非无限追加。关于 schema 定义的完整规则,可参阅仓库中的 Schema 教程。
底层机制:连接器会在一个专门线程中启动,并持续工作直到 pw.run() 停止。对应源码 ConnectorSubject.start:start 创建并启动一个线程执行 run;若 run 抛出异常,异常会被记录到 self._exception,随后触发 on_stop 和 close。
进阶场景:使用外部 Python 库接入任意流
⚠️ 注意:Twitter 已关闭其免费层流式 API,以下 Tweepy 示例仅作模式演示,接入真实 Twitter API 需付费凭证。
第二个示例使用外部库 Tweepy 加载推文流。Tweepy 是访问 Twitter API 的 Python 库,可用 pip install tweepy 安装。
第一步:创建持有 subject 的客户端类。 创建 TwitterClient 类,继承 tweepy.StreamingClient:
import tweepy
class TwitterClient(tweepy.StreamingClient):
_subject: TwitterSubject
def __init__(self, subject: TwitterSubject) -> None:
super().__init__(BEARER_TOKEN)
self._subject = subject
def on_response(self, response) -> None:
self._subject.next(
key=response.data.id,
text=response.data.text,
)
客户端持有 subject 对象;on_response 在收到流的新响应时被调用,这正是把消息转换为目标格式并送入 subject 缓冲区的正确位置。注意在 next 方法中,每一列的值都作为独立的关键字参数传入。
第二步:定义 subject。 与之前一样,定义并连接客户端:
import pathway as pw
class TwitterSubject(pw.io.python.ConnectorSubject):
_twitter_client: TwitterClient
def __init__(self) -> None:
super().__init__()
self._twitter_client = TwitterClient(self)
def run(self) -> None:
self._twitter_client.sample()
def on_stop(self) -> None:
self._twitter_client.disconnect()
这里发生了三件事:
TwitterClient在 subject 初始化时创建。由于TwitterClient内部需要访问 subject,必须把 subject 传入其构造函数;run方法启动推文流。一旦启动,流会持续流动,直到被关闭或发生失败;on_stop在流关闭或失败时被调用,是你执行清理(如断开连接)的机会。
第三步:像之前一样调用 pw.io.python.read:
class InputSchema(pw.Schema):
key: int = pw.column_definition(primary_key=True)
text: str
table = pw.io.python.read(
TwitterSubject(),
schema=InputSchema
)
pw.io.csv.write(table, "output.csv")
pw.run()
仓库中提供了该示例的完整可运行版本:twitter_connector_example.py。完整示例脚本从环境变量 TWITTER_API_TOKEN 读取 bearer token,在 pw.io.python.read 中显式设置了 autocommit_duration_ms=1000,并用 try/except KeyboardInterrupt 优雅结束 示例入口;该目录下的 README 说明了运行步骤:安装 Pathway、pip install -r requirements.txt、export TWITTER_API_TOKEN=<BEARER_TOKEN>、python twitter_connector_example.py,最后按 CTRL+C 停止。
ConnectorSubject 接口参考
上面两个示例都实现了 ConnectorSubject。下面结合 源码实现 详解该类的接口。
需要实现的方法
| 方法 | 职责 |
|---|---|
run |
主函数,负责消费数据,并使用下文某个方法把数据喂入缓冲区;在独立线程中运行 |
on_stop |
在流关闭或失败(即 run 结束)后被调用,是执行各种清理逻辑的好位置;基类默认实现为空 定义位置 |
内置方法
next(**kwargs):以关键字参数接收值,表中所有列的值都应传入该方法,所有值作为单行发送到引擎;若某列在 schema 中定义了默认值,则可不传 实现。next_bytes(message: bytes):把bytes发送到引擎的data列;若想让字节出现在其他列或多列中,直接使用next。next_str(message: str):把字符串发送到引擎的data列;多列场景同样建议用next。commit():向引擎发送提交消息。只有数据被提交后,引擎才开始处理它。你可以用该方法手动提交,也可以依赖pw.io.python.read中设置的自动提交 实现。close():表示不会再有更多消息,发送终止哨兵;当run方法结束时会被自动调用 实现。
此外,源码中还暴露了与持久化配合的钩子方法,供进阶场景使用:
seek(state: bytes)/_seek:由 Rust 核心在启动时调用,从上次停止的位置恢复读取;on_persisted_run:在持久化运行期间被调用,用于通知状态将被持久化。
pw.io.python.read 连接器方法参考
pw.io.python.read 接受以下参数(依据当前仓库源码签名):
| 参数 | 说明 |
|---|---|
subject |
要消费的连接器 subject,即 ConnectorSubject 的实例。注意:同一个 subject 对象不能用于多个 Python 连接器,复用需创建两个独立实例 |
schema |
描述数据的 schema,即列及其类型等属性 |
format |
已弃用。原先取 "json"/"raw"/"binary","raw"/"binary" 会生成单列 data 表;现在应直接给 next 传入正确类型的值,传入该参数会触发 DeprecationWarning |
autocommit_duration_ms |
两次提交之间的最大时间(毫秒)。每隔该毫秒数,连接器收到的更新会被自动提交并推入计算图,默认 1500 毫秒 |
debug_data |
调试模式激活时替代原始数据的静态数据 |
name |
连接器的唯一名称,用于日志与监控面板;启用持久化时还会作为存储连接器进度的快照名称 |
max_backlog_size |
限制从输入源读取并在任意时刻保留在处理中的条目数量;达到上限时读取暂停,处理完成后恢复,适合对初始数据爆发的大源避免内存尖峰。注意:达到上限时 subject 的 next/next_json/next_str/next_bytes 调用会阻塞,直到队列降回限制以下 |
从实现看,read 内部会做两件事:一是校验 subject 未被重复使用(subject._already_used 检查,重复使用会抛出 ValueError),二是经由 _create_python_datasource 把 subject 的 start/seek/read/end 等回调打包成 api.PythonSubject,交给底层数据源构建——autocommit_duration_ms 正是以 commit_duration_ms 形式传入 DataSourceOptions 的。
底层原理:缓冲区、线程与 Rust 读取器
Python 侧的 ConnectorSubject 通过一个 Queue 缓冲与引擎通信:__init__ 中创建 self._buffer = Queue() 初始化位置。next 最终调用 _add_inner,把 (PythonConnectorEventType.INSERT, key, values) 放入队列;而 commit、close 则通过 _send_special_message 发送特殊字面量消息(*COMMIT*、*FINISH* 等)——这些特殊字面量常量 定义在模块顶部。
Rust 侧对应实现在 src/connectors/data_storage/python.rs:PythonReader 持有 Python subject 引用与 schema,从队列中持续读取事件并转换为引擎内部的 Value。值得注意的是它维护了 total_entries_read 与 current_external_offset 两个字段,current_offset 方法把二者组合为 OffsetValue::PythonCursor 偏移量实现,用于持久化场景下记录读取进度;seek 方法(约 L114-L120)则调用 Python subject 的 on_persisted_run 以恢复上次位置。这解释了为什么 ConnectorSubject 提供 _report_offset 等内部钩子:外部数据源(如 Kafka)可以上报自己的外部偏移量,框架据此在重启后精确恢复。
另外,源码中有一个启发式检测函数 _are_deletions_reachable:它会静态解析 subject 的 run 方法源码,判断连接器是否使用了删除 API,从而自动推断 deletions_enabled。如果你的连接器只追加不删除,显式将 deletions_enabled 设为 False 可让引擎以 append-only 模式处理,提升性能。
小结
Pathway 的自定义 Python 连接器遵循一个极简契约:继承 ConnectorSubject、实现 run、用 next 喂数据、用 on_stop 清理。该模式不绑定任何特定协议或数据源——文件中静态文件逐行转流的简单案例,到借助 Tweepy 对接外部流式 API 的进阶案例,再到 autocommit_duration_ms、max_backlog_size、name、持久化偏移量等参数的完整参考,覆盖了从入门到生产调优的全部要素。完整可运行示例见 examples/projects/custom-python-connector-twitter。
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