gs-quant 批量风控低延迟请求:WebSocket 订阅、msgpack 解码与断线重连的 3 个关键机制
gs-quant 批量风控低延迟请求WebSocket 订阅、msgpack 解码与断线重连的 3 个关键机制【免费下载链接】gs-quantPython toolkit for quantitative finance项目地址: https://gitcode.com/GitHub_Trending/gs/gs-quantgs-quant 是一个 Python 量化金融工具包。当你对几百个持仓批量计算风险指标时请求与结果回传的通道往往先于计算本身成为瓶颈。下面拆开它低延迟链路上靠 WebSocket 跑通且不丢结果的 3 个机制。场景几百个持仓批量风控串行请求为什么撑不住假设你有 500 个持仓、要输出 10 个日期的风险指标最直觉的写法是 for 循环加同步 HTTP 请求。每发一个都要发出—等待—收回一个来回总耗时是所有来回的累加几百个请求排下来干等的时间远超计算本身。更麻烦的是服务端各请求的计算时长不一结果乱序返回连接中途一断在途请求的结果就全丢了。gs-quant 把这条链路当成生产级流水线来设计请求派发与结果回传拆成两条通道回传走长连接序列化用更紧凑的二进制格式断线还有降级路径。机制拆解 双通道设计批量 HTTP 派发 WebSocket 订阅回传派发端由GsRiskApi.calc_multi()负责把整批请求用一次批量 POST 打到/risk/calculate/bulk服务端给每个请求发回一个 reportId。回传端不走轮询而是建立一条 WebSocket 长连接客户端与服务端可随时互发消息的持久连接订阅地址是/risk/calculate/results/subscribe。每产生一批新 reportId客户端就用ws.send()把 id 列表推过去服务端算完就立刻推帧回来。连接本身由 gs_quant/session.py 里connect_websocket()异步上下文管理器建立协商子协议为msgpack-binary。队列工具drain_queue_async()与shutdown_queue_listener()在 gs_quant/api/gs/risk.py 同包的基类中实现。 一帧结果怎么读一个状态字符区分四种载荷服务端推回的每帧都有固定形状reportId;状态符载荷。按第一个分号切开左半是请求 id右半的首字符决定解码方式E 是错误字符串R 是 JSON 字符串M 是 base64 包装的 msgpackB 是原始 msgpack 二进制。msgpack 是二进制序列化格式比 JSON 更紧凑、解码更快可理解为 JSON 的轻量表亲。走 B 分支时省掉一次 base64 编解码。同一套application/x-msgpack内容类型在数据接口 gs_quant/api/gs/data.py 的行情查询里也在用属于全库统一的编解码约定。️ 断线容错指数退避重连再不行就轮询兜底连接断开时先查关闭码。若属于 1000、1001、1006正常关闭、离开、异常关闭客户端按 1s、2s、4s、8s 翻倍退避重建连接并把未完成的 reportId 列表重新订阅一遍日志里会打Re-subscribing N requests最多重试 5 次。若域名解析失败gaierror直接抛WebsocketUnavailableget_results()捕获后切到__get_results_poll()轮询路径把 reportId 一次性 POST 到/risk/calculate/results/bulk拉结果。也就是说没有 WebSocket 也能把任务跑完只是慢一些——降级不是失败只是变慢。代码走读两处看清链路要害第一处在__get_results_ws()的帧解析段位于 gs_quant/api/gs/risk.py# 消息形状: REQUEST_ID;STATUS_CHARDATA raw_res result_listener.result() separator b; if isinstance(raw_res, bytes) else ; # 在第一个分号处切分: 左侧请求 id, 右侧结果体 request_id_raw, _, result_data_raw raw_res.partition(separator) status, risk_data result_data_raw[0], result_data_raw[1:] # E错误 / RJSON / Mbase64包msgpack / B原始msgpack二进制 result ( msgpack.unpackb(risk_data) if status B else msgpack.unpackb(base64.b64decode(risk_data)) if status M else json.loads(risk_data) if status R else RuntimeError(risk_data) )第二处是重连主循环注意退避用math.pow(2, attempts - 1)关闭确认只等 50ms源码注释提到实际观察到过最长约 1000ms 的等待所以干脆不等attempts, max_attempts 0, 5 while attempts max_attempts: if attempts 0: await asyncio.sleep(math.pow(2, attempts - 1)) # 1s,2s,4s,8s ws_url f/{api_version}/risk/calculate/results/subscribe async with risk_session.async_.connect_websocket( ws_url, subprotocols[msgpack-binary] if cls.USE_MSGPACK else None, close_timeout0.05, # 不阻塞等待对端关闭确认 ) as ws: error await handle_websocket()⚙️ 数据与效果三个值得记的工程口径场景指标说明批量派发一次 POST 携带整批请求/risk/calculate/bulk一趟返回全部 reportId免去逐请求往返结果匹配乱序帧按 reportId 对号入座客户端持有 pending 字典帧到即弹出不依赖到达顺序断线容错退避 1s/2s/4s/8s最多 5 次重连失败或域名不可达时降级为轮询/risk/calculate/results/bulk另有两处硬编码值得留意单条派发 POST 的超时是 181 秒_exec()里timeout181订阅连接的发送超时是 30 秒。 落地建议三步启用 msgpack 批量路径把请求列表整体交给calc_multi()或RiskApi.run()派发不要单条请求另起一趟保持GsRiskApi.USE_MSGPACK True默认值session 会自动带Content-Type: application/x-msgpack头注意只有批量请求才走 msgpack 编码单条请求仍用 JSON连接建立后确认协商到的子协议是msgpack-binary否则帧会落到 Mbase64分支多一次编解码开销。如何验证 WebSocket 订阅在你的环境生效重连场景观察日志是否出现Re-subscribing N requests说明重新订阅链路在走若捕获到WebsocketUnavailable说明域名解析或网络策略不通已进入轮询降级先查网络白名单频繁撞上 181 秒 POST 超时时先查并发与批量大小再怀疑网络。如何控制请求在途量不失控gs_quant/api/risk.py 的run_async()按持仓数 × 日期数折算每个请求的权重攒够一个 chunk 才派发每收回一份结果就放行等量新请求让在途量保持大致恒定。自己调用时按机器内存和服务端容量调max_concurrent避免一次性全量压入。用一句话收尾批量风控要快关键不在算得多快而在派发与回传是否解耦、编码是否紧凑、断线是否兜得住——gs-quant 把这三件事分别交给了双通道、msgpack 和退避重连。风控 API 批量派发与订阅实现gs_quant/api/gs/risk.py请求队列、结果组装与轮询基类gs_quant/api/risk.pyWebSocket 连接与 msgpack 序列化gs_quant/session.py【免费下载链接】gs-quantPython toolkit for quantitative finance项目地址: https://gitcode.com/GitHub_Trending/gs/gs-quant创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
上一篇/下一篇内容由系统自动关联
返回资讯列表 →