尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

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

Pathway WebSocket 自定义连接器实战:用 ConnectorSubject 与 aiohttp 消费实时数据流 Pathway WebSocket 自定义连接器实战用 ConnectorSubject 与 aiohttp 消费实时数据流【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway本文基于 PathwayPython 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基类封装“建连接、收消息、缓冲写入”的公共逻辑针对具体 APIPolygon.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_asciiFalse).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(fUnhandled 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(urlURL, api_keyAPI_KEY, symbols.*) table pw.io.python.read(subject, schemaStockAggregates)结合 read 的源码实现有几点与 WebSocket 长连接场景直接相关参数默认值说明subject必填连接器主体实例。源码中 read 会检查_already_used同一个 subject 对象只能用于一个连接器需要复用请创建新实例schema按 format 推断描述输出表的列与类型本例为StockAggregatesformatjson已废弃。源码提示应改为直接通过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_typepython的GenericDataSource把subject.start/subject.seek/subject._read/subject.end绑定到引擎侧引擎在独立线程中调start启动你的run之后不断调_read从缓冲队列取事件run结束或异常时走on_stopclose收尾。第五步订阅表变化并运行流水线教程使用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变更后的行字段名到值的 dicttime变更的处理时间单位微秒可理解为 minibatch IDis_additionTrue表示插入False表示删除/更新中的删除部分——一次更新在同一批内表现为“删旧 插新”两个操作。subscribe还支持on_end流结束时回调、on_time_end每个处理时间关闭时回调、name用于日志与监控和sort_by批内按列排序输出参数可按需扩展。工程要点小结线程与事件循环的分工Pathway 引擎在专用线程里跑run你在run内部自由地建 asyncio 事件循环跑 aiohttp 协程两者的衔接点就是缓冲队列。run返回 连接器结束引擎不再等待新消息。用next还是next_jsonnext_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 源码。【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表