Pathway 自定义 Python 连接器(Custom Python Connector)实战:把 Twitter 等任意数据流接入实时表
本篇以仓库内 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 引擎,实时接收推文并把每条推文的 id 与 text 追加写入当前目录的 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."
即 key、text 必须对应 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_datasource(python/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.csv 的 pw.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 能力,可作为学会自定义连接器后的下一站练习。
常见问题与注意事项
- Twitter 免费档 Streaming API 已关闭:文档开头即有警告,示例中
sample()可能无法拉到数据。学习阶段可以把TwitterClient换成任意测试数据生成器(例如模板二的轮询循环),机制完全一致。 - 环境变量缺失:
os.environ["TWITTER_API_TOKEN"]是硬读取,忘记export会直接报KeyError。 - 主键语义:schema 中
primary_key=True的字段决定了 Pathway 如何识别"同一行"。若你的数据源没有天然 ID,需要自行构造(如时间戳、递增序号)以保证变更流正确。 - License 行:示例顶部
pw.set_license_key(...)一行仅在企业版需要;使用社区版时按注释说明删除即可,不影响连接器机制本身。 - 只增不删的流记得关 deletions:能显著减少不必要的删除检测开销(见上文
deletions_enabled说明)。
总结
自定义 Python 连接器的本质只有三句话:继承 ConnectorSubject 并实现 run() 把数据取进来;用 next() 把每一行按 schema 推进引擎;用 pw.io.python.read 组装并交给 pw.run() 驱动。Twitter 示例的价值在于把这三句话落成了一个可运行的完整样板——即便 Twitter API 已不再免费,这套"任意数据流接入 Pathway 实时表"的模式依旧可直接迁移到你的下一个数据源上。
相关参考文件:
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 StartedRust0629
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