6.9 KiB
6.9 KiB
WS 单通道推送设计
版本: 3.1 (2026-08-25) 状态: ✅ 已完成(2026-08-26 实现并真实验证通过) 关联: docs/设计/整体设计方案_v3.md 第五节(订阅与推送) 本次迭代: 通用接口定义 + 全量 tick 订阅(whole)
实现记录:
- 新增
src/bridge_subscription.py(订阅管理:snowflake sub_id、订阅表、QMT 回调→WS 推送)bridge_http_server.py新增POST /data/subscribe、POST /data/unsubscribe、GET /data/tickbridge_ws_server.py新增broadcast_allbridge_main.pystop 清理订阅;qmt_strategy_entry.py热重载模块列表更新- 桥端口迁移至 8610
- 真实验证通过:订阅 whole → WS 收到真实增量推送(单批最多 26673 代码)、tick 快照、退订正常
一、通道定义(已确认)
只保留一个 WS 通道,不再区分 data/trade 双通道。
- WS 端点:
ws://host:8610/ws(单通道) - 所有推送(行情数据 + 交易事件)走同一条连接
- 消息用
type字段区分类型 - 客户端按
type解析分发到不同处理逻辑
变更记录:v1 曾设计双通道(/data/ws + /trade/ws),讨论后改为单通道。
二、订阅流程(已确认)
客户端通过 HTTP 提交订阅请求,订阅成功后建立 WS 连接接收推送。
1. 客户端 → HTTP 提交订阅请求(含订阅类型 + 内容)
2. 服务端订阅成功,记录:订阅ID + 订阅类型
3. 服务端 → 客户端:返回 订阅ID
4. 客户端 建立 WS 连接 (ws://host:8610/ws)
5. 开始接收数据推送(按 type 分发)
6. 以后:客户端用 订阅ID 取消某类型的订阅
已确认要点:
- 订阅/退订统一走 HTTP(请求-响应,返回明确结果)
- 服务端记录:订阅ID + 订阅类型(+ 订阅内容,如 codes)
- 订阅响应返回:订阅ID
- 去掉 client_id(单客户端场景,客户端不需要自报身份)
- 统一按单客户端处理:不搞多客户端归属路由
- WS 连接:无 client_id,推送发到唯一 WS 连接
- 按订阅过滤推送:用户订阅了什么类型,才推送什么类型; 没订阅的数据类型不推送;退订某类型后该类型停止推送
- 不做 prime 推送:WS 只推增量(subscribe_whole_quote 回调数据); 客户端需要当前快照时,主动调用透传 get_full_tick 的 HTTP 接口(另设)
- 心跳机制:沿用现有(服务端每 30s ping,客户端回 pong,已有实现)
- 退订用订阅ID(HTTP 退订接口)
三、订阅接口 URI(已确认风格 B)
一个订阅接口,type 在 body 里:
| 方法 | URI | body | 返回 |
|---|---|---|---|
| POST | /data/subscribe |
{"type":"whole","codes":["SH","SZ"]} |
{"ok":true,"sub_id":123} |
| POST | /data/unsubscribe |
{"sub_id":123} |
{"ok":true} |
- 数据订阅类型未来扩展(whole/stock/kline)只需加 type,不改 URI
- sub_id 生成:snowflake 类唯一 ID 算法(时间戳 + 序列号,全局唯一、趋势递增)
- 错误格式:沿用项目规范 3.6——
{"detail":"..."}+ 4xx/5xx (如 400 unsupported type / codes required、404 sub_id not found、503 订阅失败)
四、两类推送(已确认架构)
| 类别 | 是否需要订阅 | 推送方式 |
|---|---|---|
| 数据订阅(whole 全量tick 等) | 需要订阅,按订阅过滤 | 订阅了什么才推什么 |
| 账号交易通知(成交/订单/持仓/资金) | 不需要订阅,统一回传 | 账号级事件,连上 WS 即收 |
四之二、get_full_tick 透传接口(已确认新增)
WS 只推增量(subscribe_whole_quote 回调数据),不做 prime 快照推送; 客户端需要当前 tick 快照时,主动调用本接口:
| 方法 | URI | 参数 | 返回 |
|---|---|---|---|
| GET | /data/tick |
codes=600000.SH,000001.SZ(逗号分隔) |
{"ok":true,"data":{code: tick_dict}} |
- 底层:透传
ContextInfo.get_full_tick(codes)(现有bridge_data_adapter.get_full_tick) - 支持批量代码;与单票
/data/quote并存(/data/tick 批量,/data/quote 保留兼容)
五、本次迭代范围(聚焦)
本次实现:
- 通用接口定义:订阅接口(HTTP)、退订接口(HTTP)、WS 单通道、心跳
- 数据订阅:全量 tick(type=whole)——订阅/退订/增量推送
- get_full_tick 透传接口(
/data/tick,客户端按需拉快照)
本次不做(预留,保留命名,不实现):
- 数据订阅其他类型:单票行情(预留)
- 账号交易通知:成交/订单/持仓/资金回调推送
六、type 枚举(初步定义)
数据订阅类(按订阅推送)
| type | 含义 | 本次 |
|---|---|---|
whole |
全量 tick(增量推送,只含变化的品种) | ✅ 实现 |
账号交易通知类(统一回传,本次不实现)
| type | 含义 | 触发 |
|---|---|---|
trade_result |
成交回报 | deal_callback |
order_update |
订单状态 | order_callback |
position_update |
持仓变更 | position_callback |
asset_update |
资金变更 | account_callback |
通用
| type | 含义 |
|---|---|
pong |
心跳响应(现有) |
已确认:不加
subscribed/unsubscribed(HTTP 响应已确认订阅结果); 不加error(订阅期错误由 HTTP 4xx/5xx 返回,运行期错误服务端内部消化,客户端无需感知)
七、whole 推送消息结构(已确认)
与 get_full_tick / subscribe_whole_quote 回调数据结构一致(官方文档确认):
{"type":"whole","data":{"600000.SH":{"timetag":"20231106 15:00:04","lastPrice":2.533,"open":2.528,"high":2.538,"low":2.521,"lastClose":2.513,"amount":1442588037.0,"volume":5701929,"pvolume":5701929,"stockStatus":5}}}
data:{code: tick_dict},code 为完整代码(600000.SH)- 增量推送只含变化的品种
- tick_dict 字段透传 QMT 原始字段(timetag/lastPrice/open/high/low/lastClose/amount/volume/pvolume/stockStatus...)
八、改造点(已确认)
| 文件 | 改动 |
|---|---|
bridge_http_server.py |
新增路由:POST /data/subscribe、POST /data/unsubscribe、GET /data/tick |
bridge_ws_server.py |
新增 broadcast_all(推送 whole 到唯一连接);现有心跳保留 |
bridge_subscription.py(新) |
订阅管理:订阅表(sub_id → type/codes)、snowflake sub_id、订阅/退订逻辑、QMT 订阅回调 → 推送 |
bridge_main.py |
初始化订阅管理器;adjust 里驱动 |
bridge_data_adapter.py |
get_full_tick 已存在(批量支持) |
九、参考资料(仅背景,非设计结论)
- docs/设计/整体设计方案_v3.md 第五节:订阅与推送的总体规划
- reference/xtquant_big_convert 的 quote_subscription_manager.py:订阅管理的参考实现