Pathway WebSocket 自定义连接器实战:用 ConnectorSubject 与 aiohttp 消费实时数据流
2026/9/7 18:36:44 网站建设 项目流程

Pathway WebSocket 自定义连接器实战:用 ConnectorSubject 与 aiohttp 消费实时数据流

【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway

本文基于 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 引擎。

教程的完整目标链条是:

  1. 抽象一个通用AIOHttpWebsocketSubject基类,封装“建连接、收消息、缓冲写入”的公共逻辑;
  2. 针对具体 API(Polygon.io)实现消息处理与握手流程;
  3. 定义pw.Schema描述输出表结构;
  4. 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方法是同步入口,内部驱动 asyncioconsume协程在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_messagejson.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
formatjson已废弃。源码提示应改为直接通过next传入正确类型的值;使用next_json时默认按 json 格式处理
autocommit_duration_ms1500两次 commit 之间的最大间隔。每经过该时长,连接器收到的更新会被自动提交并推进入计算图。对持续推流的 WebSocket 场景,这个自动提交机制保证数据以有界延迟流入下游
nameNone连接器唯一名称,用于日志与监控面板;启用持久化时也作为进度快照的名称
max_backlog_sizeNone处理中事件数的上限。达到上限时,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_additionTrue表示插入,False表示删除/更新中的删除部分——一次更新在同一批内表现为“删旧 + 插新”两个操作。

subscribe还支持on_end(流结束时回调)、on_time_end(每个处理时间关闭时回调)、name(用于日志与监控)和sort_by(批内按列排序输出)参数,可按需扩展。

工程要点小结

  1. 线程与事件循环的分工:Pathway 引擎在专用线程里跑run,你在run内部自由地建 asyncio 事件循环跑 aiohttp 协程;两者的衔接点就是缓冲队列。run返回 = 连接器结束,引擎不再等待新消息。
  2. next还是next_jsonnext_json把 dict 序列化为 JSON 后按 schema 解析,适合消息本身接近 JSON 对象的场景(如本例);如果需要把消息拆分到多个字段、或使用与 schema 类型不直接对应的 Python 值,可直接用next传关键字参数,并显式匹配 schema 类型。
  3. 背压与提交:对 WebSocket 这类持续推流源,可关注autocommit_duration_ms(默认 1500ms)决定数据可见延迟;流量大时用max_backlog_size引入背压,防止缓冲无限增长。
  4. 失败语义:在on_ws_messageraise会让连接器线程捕获异常并在end时重抛,整个pw.run()会以错误终止——这是把远端 API 的error状态显式暴露给运行时的正确做法。
  5. 清理钩子:若连接资源需要在停止时显式释放(如调用服务端断开),可覆写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 源码。

【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询