首页
/ Pathway 自定义 Python 连接器:实现任意数据源接入流式处理

Pathway 自定义 Python 连接器:实现任意数据源接入流式处理

2026-09-05 19:18:51作者:咎竹峻Karen

本篇指南介绍如何在 Pathway 实时数据框架中创建自定义 Python 连接器(ConnectorSubject),将静态文件、第三方库等任意数据源转换为流式输入,并深入解析 pw.io.python.read 的完整参数与底层线程/缓冲区机制。读完本篇,你可以掌握:从数据源读数据并通过 next 推入引擎缓冲区的完整模式、外部库(如 Tweepy)与 Pathway 的对接方式,以及 autocommit_duration_msmax_backlog_size 等关键参数的作用与默认值。

核心概念:ConnectorSubject 是数据源与引擎之间的桥梁

创建自定义连接器的核心步骤是:继承 pw.io.python.ConnectorSubject 抽象类,并实现其中的 run 方法——该方法负责从数据源读取数据并将数据喂入缓冲区。

ConnectorSubject 充当数据源与 Pathway 引擎之间的桥梁,它内置了 next 方法,允许你将数据推入缓冲区。从 ConnectorSubject 的类定义 可以看到几个关键事实:

  • run 是唯一的抽象方法(@abstractmethod),必须实现;
  • 该函数由 Pathway 引擎在独立线程中启动;
  • run 函数终止后,连接器即被视为结束,Pathway 不会再等待它的新消息;
  • 如果连接器不删除记录,可设置类属性 deletions_enabledFalse,有助于提升性能;
  • 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_stopclose

进阶场景:使用外部 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()

这里发生了三件事:

  1. TwitterClient 在 subject 初始化时创建。由于 TwitterClient 内部需要访问 subject,必须把 subject 传入其构造函数;
  2. run 方法启动推文流。一旦启动,流会持续流动,直到被关闭或发生失败;
  3. 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.txtexport 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) 放入队列;而 commitclose 则通过 _send_special_message 发送特殊字面量消息(*COMMIT**FINISH* 等)——这些特殊字面量常量 定义在模块顶部

Rust 侧对应实现在 src/connectors/data_storage/python.rs:PythonReader 持有 Python subject 引用与 schema,从队列中持续读取事件并转换为引擎内部的 Value。值得注意的是它维护了 total_entries_readcurrent_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_msmax_backlog_sizename、持久化偏移量等参数的完整参考,覆盖了从入门到生产调优的全部要素。完整可运行示例见 examples/projects/custom-python-connector-twitter

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

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.13 K
2.75 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
857
1.35 K
docsdocs
暂无描述
Markdown
897
5.8 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
529
593
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
915
1.83 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.58 K
1.01 K
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.35 K
1.46 K
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
1.01 K
515
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
547
388