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

资讯详情

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

WebSocket 连接池生产级实现:实时行情高可用与负载均衡的 TaoToken 配置骨架

WebSocket 连接池生产级实现:实时行情高可用与负载均衡的 TaoToken 配置骨架 1. 从轮询到连接池实时行情接入的必经之路如果你正在做量化行情、盯盘工具或者交易信号推送大概率经历过这个阶段先用requests写个循环每隔一秒拉一次价格跑起来看着还行标的加到几十个之后就开始出问题——延迟忽高忽低、接口限频、CPU 占用飙升。这不是代码写得不好而是轮询这种模式本身有天花板。WebSocket 把这个问题解决了一大半服务端主动推送延迟从秒级降到毫秒级一个连接就能订阅多个标的。但单连接同样有上限订阅数一多消息积压、单点故障、服务端订阅软限制都会冒出来。真正上生产环境你需要的是连接池多条 WebSocket 连接分摊订阅压力配合健康检查、故障切换和负载均衡让业务层对底层连接状态完全无感知。这篇文章聚焦实时行情场景从统一 Key/API 通道的角度切入给出一套可复制的config.toml与settings.json配置骨架并附上连接池健康检查与故障切换的验证动作。适合已经跑通单连接、准备把行情接入做到生产级可用的开发者。文中涉及的统一接入通道以 TaoToken 为例官网地址是 https://taotoken.net/?utm_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewriteutm_content API 入口为 https://taotoken.net/api 。2. 前置准备统一 Key 与 API 通道2.1 为什么需要统一通道多市场行情接入最烦的不是 WebSocket 本身而是每个市场一套心跳协议、一套重连策略、一套消息格式。美股一个连接、港股一个连接、加密货币再来一个代码里的连接管理器很快就变成一团乱麻。工程上的做法是找一个能跨市场的统一网关把心跳、鉴权、订阅协议标准化客户端只维护一套连接池逻辑。TaoToken 在这里扮演的就是统一通道的角色一个 API Key 覆盖多市场订阅心跳协议统一为固定间隔 ping消息格式一致。这样连接池的负载均衡器不需要关心“这条连接是美股还是港股”只需要按订阅数分配即可。2.2 获取 API Key登录 TaoToken 控制台在 API Keys 页面创建一个新 Key。建议按环境区分开发环境一个 Key生产环境一个 Key方便出问题时快速定位和吊销。创建后立即复制保存页面刷新后不再显示完整 Key。控制台入口https://taotoken.net/console?utm_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewriteutm_content API Keys 管理https://taotoken.net/api-keys?utm_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewriteutm_content注意API Key 严禁硬编码到代码里也不要提交到 Git 仓库。生产环境统一从环境变量或密钥管理服务读取。2.3 环境依赖Python 侧需要websockets和tomliPython 3.11 以下读取 TOML 用如果追求更低延迟可以额外装uvlooppip install websockets tomli uvloopNode.js 侧如果做前端行情展示可以用ws库配置读取用内置的fs加JSON.parse即可不需要额外依赖。3. 可复制配置骨架config.toml 与 settings.json3.1 config.toml连接池核心参数下面这份config.toml是连接池的骨架配置覆盖了连接数、订阅上限、心跳、重连、负载均衡策略等关键项。你可以直接复制到项目根目录按实际订阅规模调整数值。# config.toml - WebSocket 连接池生产级配置骨架 [api] # 统一接入通道从环境变量读取禁止硬编码 base_url wss://taotoken.net/api/realtime api_key_env TAOTOKEN_API_KEY # 请求超时秒 connect_timeout 10 [pool] # 连接池大小订阅数 / 单连接上限 1 个热备 pool_size 4 # 单连接最大订阅标的数实测超过 30 后 P99 延迟陡升 max_symbols_per_conn 20 # 热备连接数用于故障时毫秒级接管 hot_spare_count 1 # 启动时是否预热所有连接 prewarm true [heartbeat] # 心跳间隔秒与服务端协议保持一致 interval 1 # 僵死判定连续多少次心跳无响应判定连接死亡 pong_timeout 5 [reconnect] # 初始重连延迟秒 base_delay 1 # 最大重连延迟秒 max_delay 60 # 最大重试次数超过后告警人工介入 max_retries 10 # 抖动比例避免重连风暴 jitter_ratio 0.1 [load_balance] # 负载均衡策略round_robin / least_subscriptions / weighted_random strategy least_subscriptions # 加权随机策略下延迟采样窗口秒 latency_window 30 [health_check] # 健康检查间隔秒 check_interval 1 # 连续失败多少次触发故障恢复 failure_threshold 3 # 是否启用消息 ID 去重 dedup_enabled true # 去重表大小LRU dedup_cache_size 100003.2 settings.json业务层与运行时配置config.toml管连接池本身settings.json管业务层怎么用这个池子。两者分离的好处是连接池参数相对稳定业务订阅列表可能频繁变动改一个文件不影响另一个。{ runtime: { event_loop: uvloop, log_level: INFO, metrics_enabled: true, metrics_port: 9090 }, subscriptions: { markets: [US, HK, CRYPTO], symbols: [ AAPL.US, TSLA.US, NVDA.US, MSFT.US, GOOGL.US, 700.HK, 9988.HK, 3690.HK, BTCUSDT, ETHUSDT ], channels: [ticker, depth], batch_size: 20 }, dispatch: { queue_maxsize: 10000, backpressure_policy: drop_oldest, consumer_count: 2 }, alerting: { on_conn_dead: true, on_reconnect_failed: true, webhook_url_env: ALERT_WEBHOOK } }3.3 配置加载代码把两份配置读进来注入连接池构造函数。下面这段代码可以直接用import os import json import tomli from pathlib import Path def load_config(config_path: str config.toml, settings_path: str settings.json) - dict: with open(config_path, rb) as f: config tomli.load(f) with open(settings_path, r, encodingutf-8) as f: settings json.load(f) # API Key 从环境变量注入不落盘 api_key_env config[api][api_key_env] api_key os.environ.get(api_key_env) if not api_key: raise RuntimeError(f环境变量 {api_key_env} 未设置) config[api][api_key] api_key return {config: config, settings: settings}提示config.toml里的api_key_env只存变量名真实 Key 通过环境变量注入。这样配置文件可以安全地提交到仓库Key 不会泄露。4. 连接池核心实现与验证4.1 连接池骨架下面这段代码是连接池的核心骨架重点展示负载均衡分配、心跳保活、故障恢复三个环节。为保持骨架清晰消息分发和订阅迁移部分用注释标注了位置你可以按业务填充。import asyncio import json import random import websockets from dataclasses import dataclass, field from enum import Enum from typing import List, Dict, Optional class ConnState(Enum): IDLE idle ACTIVE active DEAD dead dataclass class WSConnection: conn_id: str state: ConnState ConnState.IDLE ws: Optional[websockets.WebSocketClientProtocol] None symbols: List[str] field(default_factorylist) last_pong: float 0.0 avg_latency: float 0.0 class ConnectionPool: def __init__(self, config: dict, settings: dict): self.cfg config self.settings settings self.api_key config[api][api_key] self.pool_size config[pool][pool_size] self.max_per_conn config[pool][max_symbols_per_conn] self.connections: List[WSConnection] [] self._hot_spare: List[WSConnection] [] self._lock asyncio.Lock() async def start(self): 预热连接池预留热备 for i in range(self.pool_size): conn WSConnection(conn_idfconn-{i}) await self._connect_with_backoff(conn) asyncio.create_task(self._heartbeat_loop(conn)) asyncio.create_task(self._message_loop(conn)) self.connections.append(conn) # 预留热备 spare_count self.cfg[pool][hot_spare_count] for _ in range(spare_count): if self.connections: spare self.connections.pop() spare.state ConnState.IDLE self._hot_spare.append(spare) async def _connect_with_backoff(self, conn: WSConnection): url f{self.cfg[api][base_url]}?api_key{self.api_key} retry, base, cap 0, self.cfg[reconnect][base_delay], \ self.cfg[reconnect][max_delay] while retry self.cfg[reconnect][max_retries]: try: conn.ws await websockets.connect( url, open_timeoutself.cfg[api][connect_timeout]) conn.state ConnState.ACTIVE conn.last_pong asyncio.get_event_loop().time() return except Exception: delay min(base * (2 ** retry), cap) jitter random.uniform(0, delay * self.cfg[reconnect][jitter_ratio]) await asyncio.sleep(delay jitter) retry 1 raise RuntimeError(f{conn.conn_id} 重连失败需人工介入) async def _heartbeat_loop(self, conn: WSConnection): interval self.cfg[heartbeat][interval] while conn.state ! ConnState.DEAD: try: if conn.ws and conn.state ConnState.ACTIVE: await conn.ws.send(json.dumps({cmd: ping})) except Exception: pass await asyncio.sleep(interval) async def _message_loop(self, conn: WSConnection): timeout self.cfg[heartbeat][pong_timeout] while conn.state ! ConnState.DEAD: try: msg await asyncio.wait_for(conn.ws.recv(), timeouttimeout) data json.loads(msg) if data.get(cmd) pong: conn.last_pong asyncio.get_event_loop().time() else: # 生产环境在此处调用消息分发await self._dispatch(conn, data) pass except asyncio.TimeoutError: # 超时触发健康检查判定 pass except Exception: conn.state ConnState.DEAD asyncio.create_task(self._recover(conn)) break async def _recover(self, dead_conn: WSConnection): 故障恢复热备接管 指数退避重连 if self._hot_spare: hot self._hot_spare.pop() hot.state ConnState.ACTIVE # 生产环境需实现订阅迁移await self._migrate_subscriptions(dead_conn, hot) self.connections.append(hot) try: await self._connect_with_backoff(dead_conn) # 生产环境需恢复订阅await self._restore_subscriptions(dead_conn) except RuntimeError: # 触发告警 pass async def subscribe(self, symbols: List[str]): 外部订阅接口自动负载均衡 async with self._lock: strategy self.cfg[load_balance][strategy] if strategy least_subscriptions: target min(self.connections, keylambda c: len(c.symbols)) elif strategy round_robin: target self.connections[0] self.connections.append(self.connections.pop(0)) else: weights [1.0 / (c.avg_latency 0.001) for c in self.connections] target random.choices(self.connections, weightsweights)[0] if len(target.symbols) len(symbols) self.max_per_conn: if self._hot_spare: target self._hot_spare.pop() target.state ConnState.ACTIVE self.connections.append(target) # 生产环境需实现实际订阅await self._do_subscribe(target, symbols) target.symbols.extend(symbols)4.2 验证请求确认连接池正常工作写完骨架后用下面这段验证脚本确认连接池能正常启动、订阅、接收消息。脚本会打印每个连接的订阅数和收到的第一条行情消息。import asyncio from pool import ConnectionPool, load_config async def verify(): cfg load_config() pool ConnectionPool(cfg[config], cfg[settings]) await pool.start() symbols cfg[settings][subscriptions][symbols] await pool.subscribe(symbols) # 打印连接池状态 for conn in pool.connections: print(f{conn.conn_id} state{conn.state.value} fsymbols{len(conn.symbols)}) # 等待 10 秒观察消息接收 await asyncio.sleep(10) # 模拟故障手动关闭一条连接观察热备接管 if pool.connections: victim pool.connections[0] print(f模拟故障关闭 {victim.conn_id}) await victim.ws.close() await asyncio.sleep(5) print(f故障后连接数{len(pool.connections)} f热备剩余{len(pool._hot_spare)}) if __name__ __main__: asyncio.run(verify())预期输出类似conn-0 stateactive symbols20 conn-1 stateactive symbols20 conn-2 stateactive symbols20 模拟故障关闭 conn-0 故障后连接数3热备剩余0如果看到连接数在故障后保持不变、热备被激活说明故障切换链路是通的。4.3 成功结果判定一次成功的连接池验证应该满足以下三个条件第一启动后所有连接状态为active订阅数按least_subscriptions策略均匀分布偏差不超过 1 个标的。第二模拟关闭一条连接后5 秒内热备接管业务层收到的行情消息不中断。第三连续运行 1 小时心跳日志中pong响应间隔稳定在 1 秒左右无僵死判定误报。5. 本篇常见错排查5.1 连接建立后立即断开现象websockets.connect返回成功但_message_loop第一次recv就抛异常连接状态变为DEAD。排查方向先确认 API Key 是否正确注入。常见错误是config.toml里写了api_key_env TAOTOKEN_API_KEY但环境变量没导出代码里读到空字符串。用echo $TAOTOKEN_API_KEY确认。其次检查base_url是否拼错wss://前缀不能少。5.2 心跳正常但收不到行情现象pong响应正常last_pong持续更新但业务层收不到 ticker 消息。排查方向检查订阅命令是否真正发出。骨架代码里_do_subscribe是注释状态需要你填充实际订阅逻辑。另外确认settings.json里的channels字段与服务端支持的频道名一致ticker和depth是常见命名但不同服务商可能有差异。5.3 重连风暴导致服务端限流现象网络抖动后所有连接同时进入重连服务端返回 429 或直接拒绝连接。排查方向确认jitter_ratio配置生效。骨架代码里重连延迟是delay random.uniform(0, delay * jitter_ratio)如果jitter_ratio设为 0所有连接会在同一时刻重连。建议保持 0.1 以上。另外检查max_retries超过后应该触发告警而不是无限重试。5.4 消息重复导致策略误判现象重连后收到重复的 ticker 消息策略层重复下单。排查方向启用dedup_enabled在消息分发前用消息 ID 做 LRU 去重。骨架代码里dedup_cache_size设为 10000按每秒 100 条消息估算可以覆盖 100 秒的重传窗口。如果消息量更大适当调大这个值。5.5 热备连接未激活现象模拟故障后连接数减少热备没有被激活。排查方向检查_recover方法里self._hot_spare.pop()是否执行。常见错误是hot_spare_count设为 0或者start方法里预留热备的逻辑被跳过。另外确认_migrate_subscriptions是否实现如果没实现热备接管后订阅列表是空的业务层依然收不到消息。6. 接入文档与后续动作连接池跑通之后下一步通常是把它接入到实际的策略或展示层。如果你在排障过程中需要确认 API 的请求格式、心跳协议细节或者错误码含义可以查阅接入文档接入文档https://taotoken.net/doc?utm_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewriteutm_content如果只是想先验证模型对话或行情数据格式不急着写连接池可以直接在模型对话页面发一条测试请求确认 Key 和通道都正常模型对话https://taotoken.net/chat?utm_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewriteutm_content对于需要长期跑编码任务或者 Agent 场景的读者连接池只是基础设施的一部分更完整的方案可以参考 Coding PlanCoding Planhttps://taotoken.net/coding-plan?utm_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewriteutm_content最后提醒一句连接池的pool_size和max_symbols_per_conn不是拍脑袋定的。先用单连接压测出你实际场景下的 P99 延迟拐点再按拐点反推单连接上限最后用订阅总数除以单连接上限加一得到池子大小。这个顺序反过来做很容易配出一个看起来能用、实际一压就垮的池子。
返回列表