Pathway WebSocket 自定义连接器实战:用 ConnectorSubject 与 aiohttp 消费实时数据流
本文基于 Pathway(Python ETL 框架,用于流处理、实时分析、LLM 流水线与 RAG)官方教程文档,讲解如何创建一个自定义 WebSocket 连接器:先抽象出一个通用的 aiohttp WebSocket 消费基类,再以 Polygon.io Stocks API 为例,演示“连接 → 认证 → 订阅”的多步消息握手过程,最终把 WebSocket 实时数据流接入 Pathway 计算图。读完本篇,你可以掌握 pw.io.python.ConnectorSubject 的接口约束、pw.io.python.read 的关键参数,并能将该模式改接到任意 WebSocket API 上。
为什么需要自定义 WebSocket 连接器
WebSockets 协议的特点是:每个 API 的通信流程都可能不同——有的连接即可推流,有的需要先鉴权,有的还需要显式订阅主题。Pathway 没有为每一种 WebSocket API 内置连接器,而是提供了一套通用的 Python 连接器扩展机制,允许你用任意第三方库(本文使用 aiohttp)编写消费逻辑,再把数据喂入 Pathway 引擎。
教程的完整目标链条是:
- 抽象一个通用
AIOHttpWebsocketSubject基类,封装“建连接、收消息、缓冲写入”的公共逻辑; - 针对具体 API(Polygon.io)实现消息处理与握手流程;
- 定义
pw.Schema描述输出表结构; - 用
pw.io.python.read生成输入表,用pw.io.subscribe观察变化,用pw.run运行流水线。
第一步:抽象通用 WebSocket 消费基类
自定义连接器的入口是继承 pw.io.python.ConnectorSubject(源码位于 python/pathway/io/python/init.py),并实现唯一的抽象方法 run。教程给出的通用基类如下:
import pathway as pw
import asyncio
import aiohttp
from aiohttp.client_ws import ClientWebSocketResponse
class AIOHttpWebsocketSubject(pw.io.python.ConnectorSubject):
_url: str
def __init__(self, url: str):
super().__init__()
self._url = url
def run(self):
async def consume():
async with aiohttp.ClientSession() as session:
async with session.ws_connect(self._url) as ws:
async for msg in ws:
if msg.type == aiohttp.WSMsgType.CLOSE:
break
else:
result = await self.on_ws_message(msg, ws)
for row in result:
self.next_json(row)
asyncio.new_event_loop().run_until_complete(consume())
async def on_ws_message(self, msg, ws: ClientWebSocketResponse) -> list[dict]:
...
这段代码的设计要点:
run方法是同步入口,内部驱动 asyncio。consume协程在run中通过asyncio.new_event_loop().run_until_complete(consume())执行,即运行在一个独立的 asyncio 事件循环里。这与引擎的线程模型一致——从 ConnectorSubject.start 的源码可以看到,Pathway 会把run放在一个专用的threading.Thread中启动,run返回即表示连接器结束(close哨兵消息随后发出)。因此“无限循环消费 + 收到 CLOSE 时break”是保持连接常驻的正确写法。- 消息处理委托给抽象方法
on_ws_message。基类只负责“收消息 → 调用子类处理 → 写缓冲”的骨架,子类决定每条消息如何转换成行(可能一条消息拆出多行,也可能某些消息不产生任何行)。 - 结果通过
self.next_json(row)写入缓冲。next_json接收一个 dict,内部执行json.dumps(message, ensure_ascii=False).encode("utf-8")后压入缓冲队列(见 ConnectorSubject.next_json)。这意味着每行数据会以 JSON 编码进入引擎,再按 schema 解析成列——所以 dict 的键必须与 schema 字段名对应。
第二步:实现真实场景——Polygon.io Stocks API
教程以 Polygon.io Stocks API 为例,该连接器订阅所选股票的 1 秒级聚合(A 事件)。Polygon 的关键约束是:连接建立后不会直接推数据,必须先发送认证消息,收到 auth_success 后才能发送订阅消息。这个“多步消息交换”正是 WebSocket 连接器最典型的形态,on_ws_message 用状态机式的路由来处理它:
import json
class PolygonSubject(AIOHttpWebsocketSubject):
_api_key: str
_symbols: str
def __init__(self, url: str, api_key: str, symbols: str):
super().__init__(url)
self._api_key = api_key
self._symbols = symbols
async def on_ws_message(
self, msg: aiohttp.WSMessage, ws: ClientWebSocketResponse
) -> list[dict]:
if msg.type == aiohttp.WSMsgType.TEXT:
result = []
payload = json.loads(msg.data)
for object in payload:
match object:
case {"ev": "status", "status": "connected"}:
# make authorization request if connected successfully
await self._authorize(ws)
case {"ev": "status", "status": "auth_success"}:
# request a stream, once authenticated
await self._subscribe(ws)
case {"ev": "A"}:
# append data object to results list
result.append(object)
case {"ev": "status", "status": "error"}:
raise RuntimeError(object["message"])
case _:
raise RuntimeError(f"Unhandled payload: {object}")
return result
else:
return []
async def _authorize(self, ws: ClientWebSocketResponse):
await ws.send_json({"action": "auth", "params": self._api_key})
async def _subscribe(self, ws: ClientWebSocketResponse):
await ws.send_json({"action": "subscribe", "params": self._symbols})
对照源码理解这段握手的几个细节:
- 一条 payload 是一个 JSON 数组,逐个对象路由。Polygon 每条文本消息序列化后包含一个对象列表,所以
on_ws_message先json.loads,再对列表内每个对象做match分支:connected→ 发起认证;auth_success→ 发起订阅;A→ 追加为结果行;error→ 抛出RuntimeError使流水线失败(异常会被 ConnectorSubject 的线程包装捕获 并在end时重新抛出);未知对象同样抛错,避免静默丢数据。 - 非 TEXT 消息返回空列表。二进制帧、ping 等不产生数据行,直接返回
[]即可,基类循环会继续等待下一条消息。 - 握手是“事件驱动”的,而不是主动轮询。认证与订阅都在收到对应状态消息时才发送,顺序由 API 的状态消息自然驱动,这是处理多步 WebSocket 握手的推荐方式。
第三步:定义输出表的 Schema
定义一个 pw.Schema 来描述结果表的结构。由于连接器不对入站 payload 做任何修改,schema 字段与 API 返回的对象一一对应:
class StockAggregates(pw.Schema):
sym: str # stock symbol
o: float # opening price
v: int # tick volume
s: int # starting tick timestamp
e: int # ending tick timestamp
...
需要注意:next_json 传入的 dict 会被整体序列化,引擎按 schema 声明的列提取值;未声明的字段会被忽略,声明了但消息里没有的字段需要 schema 提供默认值,否则会解析失败。
第四步:用 pw.io.python.read 创建输入表
把 subject 交给 pw.io.python.read 即可得到输入表:
URL = "wss://delayed.polygon.io/stocks"
API_KEY = "your-api-key"
subject = PolygonSubject(url=URL, api_key=API_KEY, symbols=".*")
table = pw.io.python.read(subject, schema=StockAggregates)
结合 read 的源码实现,有几点与 WebSocket 长连接场景直接相关:
| 参数 | 默认值 | 说明 |
|---|---|---|
subject |
必填 | 连接器主体实例。源码中 read 会检查 _already_used:同一个 subject 对象只能用于一个连接器,需要复用请创建新实例 |
schema |
按 format 推断 | 描述输出表的列与类型;本例为 StockAggregates |
format |
json |
已废弃。源码提示应改为直接通过 next 传入正确类型的值;使用 next_json 时默认按 json 格式处理 |
autocommit_duration_ms |
1500 |
两次 commit 之间的最大间隔。每经过该时长,连接器收到的更新会被自动提交并推进入计算图。对持续推流的 WebSocket 场景,这个自动提交机制保证数据以有界延迟流入下游 |
name |
None |
连接器唯一名称,用于日志与监控面板;启用持久化时也作为进度快照的名称 |
max_backlog_size |
None |
处理中事件数的上限。达到上限时,subject 的 next / next_json 等调用会阻塞,直到队列回落——从 Queue(max_backlog_size) 的实现可见它把无界队列换成有界队列。对突发流量大的数据源,这是避免内存尖峰的背压手段 |
从源码结构看,read 最终通过 _create_python_datasource 构建一个 storage_type="python" 的 GenericDataSource,把 subject.start / subject.seek / subject._read / subject.end 绑定到引擎侧:引擎在独立线程中调 start 启动你的 run,之后不断调 _read 从缓冲队列取事件,run 结束或异常时走 on_stop + close 收尾。
第五步:订阅表变化并运行流水线
教程使用 pw.io.subscribe 观察表内变化:
import logging
def on_change(
key: pw.Pointer,
row: dict,
time: int,
is_addition: bool,
):
logging.info(f"{time}: {row}")
pw.io.subscribe(table, on_change)
再运行流水线:
pw.run()
on_change 回调签名的四个参数语义(见 subscribe 文档字符串):
key:变更行的指针;row:变更后的行,字段名到值的 dict;time:变更的处理时间,单位微秒(可理解为 minibatch ID);is_addition:True表示插入,False表示删除/更新中的删除部分——一次更新在同一批内表现为“删旧 + 插新”两个操作。
subscribe 还支持 on_end(流结束时回调)、on_time_end(每个处理时间关闭时回调)、name(用于日志与监控)和 sort_by(批内按列排序输出)参数,可按需扩展。
工程要点小结
- 线程与事件循环的分工:Pathway 引擎在专用线程里跑
run,你在run内部自由地建 asyncio 事件循环跑 aiohttp 协程;两者的衔接点就是缓冲队列。run返回 = 连接器结束,引擎不再等待新消息。 - 用
next还是next_json:next_json把 dict 序列化为 JSON 后按 schema 解析,适合消息本身接近 JSON 对象的场景(如本例);如果需要把消息拆分到多个字段、或使用与 schema 类型不直接对应的 Python 值,可直接用next传关键字参数,并显式匹配 schema 类型。 - 背压与提交:对 WebSocket 这类持续推流源,可关注
autocommit_duration_ms(默认 1500ms)决定数据可见延迟;流量大时用max_backlog_size引入背压,防止缓冲无限增长。 - 失败语义:在
on_ws_message中raise会让连接器线程捕获异常并在end时重抛,整个pw.run()会以错误终止——这是把远端 API 的error状态显式暴露给运行时的正确做法。 - 清理钩子:若连接资源需要在停止时显式释放(如调用服务端断开),可覆写
on_stop方法(在 源码中run结束或异常后、close之前被调用)。
该模式(通用 aiohttp 基类 + 子类状态机 + schema + pw.io.python.read)可以不改骨架地迁移到其他 WebSocket API:只需替换 _authorize / _subscribe 中的握手报文和 on_ws_message 中的消息路由,即可接入任意需要多步消息交换的 WebSocket 数据源。
参考文档:WebSockets connectors 教程、Custom Python connectors 教程、Python connector 源码。
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