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

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

2026-09-06 09:15:20作者:胡唯隽

本文基于 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
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_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 源码

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