首页
/ Pathway 自定义 Python 连接器(Custom Python Connector)实战:把 Twitter 等任意数据流接入实时表

Pathway 自定义 Python 连接器(Custom Python Connector)实战:把 Twitter 等任意数据流接入实时表

2026-09-07 19:56:52作者:魏侃纯Zoe

本篇以仓库内 Custom Python Connector(Twitter)示例为核心,讲解如何在 Pathway 中通过 pw.io.python.ConnectorSubject 编写自定义 Python 输入连接器:把一段持续产生的数据流(示例为 Twitter Streaming API 返回的推文)注入 Pathway 实时计算表,再交给下游算子处理并输出到 output.csv。读完本文,你将掌握 ConnectorSubject 的生命周期约定、pw.io.python.read 各参数的真实作用,以及把这套模式复用到任意数据源(Kafka 之外的私有协议、WebSocket、轮询接口等)的完整套路。

这个示例要解决什么问题

在 Pathway 中接入数据通常有现成连接器(Kafka、CSV、PostgreSQL 等),但当数据源是私有协议、专有 SDK 或官方尚未提供连接器的流式接口时,就需要“自定义 Python 连接器”。关联文档 custom-python-connector-twitter.md 明确指出其目的:

The purpose of this example is to show how to implement custom python connector.

示例本体是一个完整可运行的 Pipeline:用 Twitter 官方 Python SDK(tweepy)订阅实时推文流,把每条推文转换为 Pathway 表的一行,再把表的变更流写入 output.csv。正如文档所强调的:

Twitter API was used as an example, but any data stream can be fed into Pathway using this method.

也就是说,Twitter 只是"饵",这套「编写一个 Python 生产者类 → 用 read 接入 → 交给引擎」的方法论可以套到任何数据流上。

⚠️ 需要特别留意文档开头的警告:Twitter 已关闭其免费档的 Streaming API。因此直接复现"拉真实推文"可能需要付费档凭证或替代方案;但这不影响它作为学习自定义连接器写法的绝佳样板。

示例目录与文件

仓库中该示例由 4 个文件组成:

文件 作用
README.md 与关联文档同内容的示例说明
requirements.txt 额外依赖,仅一行:tweepy==4.13.0
twitter_connector_example.py 核心示例代码
output.csv 运行后由 Pathway 自动生成的结果文件(不在仓库中)

tweepy 是唯一的第三方依赖——Pathway 本身已随主包安装完毕,其余全部逻辑只依赖 pathway 与标准库。

五分钟跑通示例

关联文档给出了 5 步启动流程,这里逐条展开:

1. 安装 Pathway

安装 Pathway 包(示例代码通过 import pathway as pw 使用其全部能力),安装方式按项目官方安装指引操作即可。

2. 安装额外依赖

pip install -r requirements.txt

等价于直接安装 tweepy==4.13.0——示例使用该版本的 tweepy.StreamingClient 订阅 Twitter 样本流。

3. 配置访问令牌

export TWITTER_API_TOKEN=<BEARER_TOKEN>

Bearer Token 需要从 Twitter 开发者门户申请。示例代码通过 os.environ["TWITTER_API_TOKEN"] 读取该环境变量(见 twitter_connector_example.py),因此这一步缺失会直接抛 KeyError

4. 运行示例

python twitter_connector_example.py

脚本会启动 Pathway 引擎,实时接收推文并把每条推文的 idtext 追加写入当前目录的 output.csv

5. 停止

CTRL+C 触发 KeyboardInterrupt,脚本捕获后打印 Done. 并退出(对应代码中 try: pw.run() except KeyboardInterrupt: 分支)。

逐行拆解示例源码

完整代码位于 twitter_connector_example.py,其设计分为三块:schema 定义、连接器 Subject、组装与运行。这里去掉版权注释后逐步解读。

Schema:告诉 Pathway 一行数据长什么样

class InputSchema(pw.Schema):
    key: int = pw.column_definition(primary_key=True)
    text: str

key 是推文 ID,被声明为主键(primary_key=True);text 是推文正文。主键用于后续变更流(新增/删除)的行定位,是"把流变成可维护状态表"的关键。注意推文 ID 在 Twitter API 中返回的是数值类型,因此这里声明为 int

TwitterSubject:连接器的"心脏"

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()

这是整套机制的核心:任何自定义 Python 连接器都必须继承 pw.io.python.ConnectorSubject 并实现 run() 方法。run() 内部调用 tweepy 的 sample() 阻塞式订阅样本流——数据进入 Pathway 并非靠返回值,而是靠内部回调把消息"推"进引擎

TwitterClient:SDK 回调与 Pathway 的桥

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,
        )
    ...

TwitterClient 继承 tweepy.StreamingClient,在 tweepy 的 on_response 回调里拿到每一条推文,然后调用 self._subject.next(key=..., text=...)——next()ConnectorSubject 提供的"发消息给引擎"的方法,参数必须与传入 read 的 schema 字段一一对应(详见下文原理章节)。

注意两个类互相持有引用(Subject 持 Client、Client 持 Subject),这是典型的回调型连接器写法:SDK 的事件循环在自己的线程里跑,每来一条数据就通过 next() 塞进 Pathway 内部队列

组装 Pipeline

input = pw.io.python.read(
    TwitterSubject(),
    schema=InputSchema,
    autocommit_duration_ms=1000,
)

pw.io.csv.write(input, "output.csv")

try:
    pw.run()
except KeyboardInterrupt:
    print("Done.")
  • pw.io.python.read(...):用自定义 Subject 建一个实时输入表;
  • pw.io.csv.write(input, "output.csv"):把表的变更流持续写进 CSV;
  • pw.run():启动引擎并开始处理,阻塞直到被中断。

关于代码顶部的一行,示例注释说明:使用 Pathway 企业版(Scale)时需要设置自己的 license key;使用社区版(Community)则注释掉该行即可(见 twitter_connector_example.py)。读者在本地复现时应按自己所用版本处理这一行。

原理纵深一:ConnectorSubject 的生命周期约定

ConnectorSubject 定义在 python/pathway/io/python/init.py,其文档字符串把整条契约说得很清楚,这里结合示例逐一对应:

run() 在独立线程中被引擎启动

"This function will be started by pathway engine in a separate thread."

这就是为什么示例中 run() 可以放一个无限阻塞的 sample():引擎不会等它返回,而是并行消费它不断 next() 出来的数据。这也是把任何数据源(WebSocket、串口、私有的轮询循环)接入 Pathway 的统一姿势——把"取数循环"写进 run(),每取到一条就 next() 一次

run() 返回 = 连接器结束

"When the run function terminates, the connector will be considered finished and pathway won't wait for new messages from it."

对有限的批式数据源,run() 里 for 循环结束后连接器自然收尾,引擎继续处理已接收数据即可。

on_stop()run() 结束后被调用

"Called after the end of the run function."

示例在这里调用 disconnect() 优雅关闭 tweepy 的连接,避免泄漏流式 socket。

next(**kwargs) 发送一行消息

"The arguments should be compatible with the schema passed to pw.io.python.read."

keytext 必须对应 schema 字段;若有默认值的字段可以省略。ConnectorSubject 还提供了 next_json(dict)next_str(str)next_bytes(bytes) 三种便捷发送方式(见 python/pathway/io/python/init.py),分别对应 read 在 json/raw/binary 格式下的使用。

内部缓冲与背压 类内部维护一个 _buffer: Queue。文档特别注明:如果设置了 max_backlog_size,当待处理事件数达到上限时 next() 等调用会阻塞,直到队列回落到限制以下——这正是"背压"机制的实现方式,避免生产过快打爆引擎。

性能提示:deletions_enabled

"If the subject won't delete records, set the class property deletions_enabled to False as it may help to improve the performance."

如果连接器只增不删(如本示例的推文流),把 deletions_enabled 关掉能提升性能。从源码看(python/pathway/io/python/init.py),_deletions_enabled 默认通过启发式扫描(_are_deletions_reachable)自动判断;手动关闭时类内重写 _deletions_enabled 返回 False 即可。反之,若试图在关闭删除的场景下调 remove 或对同一主键再次插入(UPSERT),会抛出 ValueError: Trying to delete a row ... but deletions_enabled is set to False.(见 python/pathway/io/python/init.py)。

原理纵深二:pw.io.python.read 的参数含义

read 的底层实现 _create_python_datasourcepython/pathway/io/python/init.py)把 Subject 包装成一个可被 Rust 核心调用的 PythonSubject 数据源。常用参数:

参数 默认值 作用
schema 必填 规定表结构,字段与 next() 的关键字参数对应
autocommit_duration_ms 1500 两次提交之间的最大时间间隔(毫秒)。每过该时长,连接器收到的更新会被提交并推入 Pathway 计算图。示例设为 1000,即约 1 秒一个批次刷新到下游
name None 连接器的唯一名称,供状态/日志标识
max_backlog_size None 待处理积压上限,超过后 next() 调用阻塞,实现背压

注意示例里写 output.csvpw.io.csv.write 默认就带有"变更流"语义——Pathway 写入的不是一次性快照,而是持续追加后续到达的行;这正是 autocommit_duration_ms 控制的提交节奏在起作用。

从 Twitter 到任意数据流:三种改写模板

基于上文原理,把 Twitter 换成你自己的数据源只需改写 Subject.run() 内的取数逻辑。这里给出与示例对应的通用心智模型:

模板一:回调/推送型(本示例) 外部 SDK 或库在事件循环里回调我们 → 在回调里调 self.next(...)。改 Twitter 为 WebSocket 客户端、消息推送 SDK 等都走这条路。

模板二:主动轮询/拉取型 run()while True: 定时拉取(如每隔 N 秒查一次 REST 接口或数据库),拉到就 next()time.sleep 控制频率。适用于没有推送能力的数据源。

模板三:有限批式 run() 里遍历本地文件/列表后自然结束,引擎处理完即收尾(对应 _is_finite 返回 True 的语义)。

一个值得参考的进阶样本是仓库中更完整的 Twitter 实时情感分析项目 examples/projects/twitter 及教程文档 2.twitter.md——它在同一主题上叠加了流式表 join、情感聚合等更多 Pathway 能力,可作为学会自定义连接器后的下一站练习。

常见问题与注意事项

  1. Twitter 免费档 Streaming API 已关闭:文档开头即有警告,示例中 sample() 可能无法拉到数据。学习阶段可以把 TwitterClient 换成任意测试数据生成器(例如模板二的轮询循环),机制完全一致。
  2. 环境变量缺失os.environ["TWITTER_API_TOKEN"] 是硬读取,忘记 export 会直接报 KeyError
  3. 主键语义:schema 中 primary_key=True 的字段决定了 Pathway 如何识别"同一行"。若你的数据源没有天然 ID,需要自行构造(如时间戳、递增序号)以保证变更流正确。
  4. License 行:示例顶部 pw.set_license_key(...) 一行仅在企业版需要;使用社区版时按注释说明删除即可,不影响连接器机制本身。
  5. 只增不删的流记得关 deletions:能显著减少不必要的删除检测开销(见上文 deletions_enabled 说明)。

总结

自定义 Python 连接器的本质只有三句话:继承 ConnectorSubject 并实现 run() 把数据取进来;用 next() 把每一行按 schema 推进引擎;用 pw.io.python.read 组装并交给 pw.run() 驱动。Twitter 示例的价值在于把这三句话落成了一个可运行的完整样板——即便 Twitter API 已不再免费,这套"任意数据流接入 Pathway 实时表"的模式依旧可直接迁移到你的下一个数据源上。

相关参考文件:

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

项目优选

收起
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++
916
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