如何查询Binance开放的Websocket流状态及实现断连重连?
WebSocket流管理、状态查询与自动重连实现方案
1. 多流管理与活跃状态查询
原生WebSocket没有内置全局统计接口,你需要自行封装实例存储逻辑,通过映射表统一管理所有连接,同时可以直接读取WebSocket实例自带的readyState属性判断状态:
0:正在连接中1:连接已开启、活跃可用2:正在关闭3:已关闭/连接失败
封装流创建与管理逻辑
// 全局流管理映射表,key为交易对名称,value存储流的全部元信息 const streamMap = new Map() // 封装创建Binance数据流的公共方法 function createStream(symbol, interval = '1h') { const address = `wss://stream.binance.com:9443/ws/${symbol.toLowerCase()}@kline_${interval}` const ws = new WebSocket(address) const streamInfo = { symbol, interval, address, ws, isActive: false, reconnectCount: 0, heartbeatTimer: null } streamMap.set(symbol, streamInfo) // 连接成功时标记为活跃 ws.onopen = () => { streamInfo.isActive = true streamInfo.reconnectCount = 0 startHeartbeat(streamInfo) // 心跳逻辑见下文 } // 你原本的消息处理逻辑可以在这里统一绑定 ws.onmessage = (event) => { const data = JSON.parse(event.data) // 处理K线数据的业务逻辑 } return streamInfo }
活跃流查询方法
function getActiveStreams() { const activeSymbols = [] streamMap.forEach((info, symbol) => { // 同时用自定义标记和原生状态双重校验,结果更准确 if (info.isActive && info.ws.readyState === WebSocket.OPEN) { activeSymbols.push(symbol) } }) return activeSymbols }
2. 断流检测与自动重连
断流分为两种情况:显式断连(触发onerror/onclose事件)和隐式断连(网络闪断但底层未触发关闭事件),需要结合心跳检测+重连退避机制实现稳定重连。
心跳检测实现
const HEARTBEAT_INTERVAL = 30000 // 每30秒发一次心跳 const HEARTBEAT_TIMEOUT = 10000 // 10秒未收到响应判定为断连 function startHeartbeat(streamInfo) { clearInterval(streamInfo.heartbeatTimer) let pongReceived = true streamInfo.heartbeatTimer = setInterval(() => { if (streamInfo.ws.readyState !== WebSocket.OPEN) { clearInterval(streamInfo.heartbeatTimer) triggerReconnect(streamInfo) return } if (!pongReceived) { clearInterval(streamInfo.heartbeatTimer) streamInfo.ws.close() triggerReconnect(streamInfo) return } pongReceived = false // Node.js环境下ws库直接调用ping方法;浏览器环境下可替换为发送请求:ws.send(JSON.stringify({method:"ping"})) streamInfo.ws.ping() setTimeout(() => { if (!pongReceived) { clearInterval(streamInfo.heartbeatTimer) streamInfo.ws.close() triggerReconnect(streamInfo) } }, HEARTBEAT_TIMEOUT) }, HEARTBEAT_INTERVAL) // 监听服务端返回的pong响应 streamInfo.ws.onpong = () => { pongReceived = true } }
带指数退避的重连逻辑
const MAX_RECONNECT_COUNT = 10 // 最大重连次数,避免无限重试 const BASE_RECONNECT_DELAY = 1000 // 基础重连延迟1秒 function triggerReconnect(streamInfo) { if (streamInfo.reconnectCount >= MAX_RECONNECT_COUNT) { console.log(`${streamInfo.symbol} 重连次数超过上限,已停止重连`) streamMap.delete(streamInfo.symbol) return } streamInfo.isActive = false streamInfo.reconnectCount += 1 // 指数退避计算延迟,避免短时间频繁请求触发限流 const delay = BASE_RECONNECT_DELAY * Math.pow(2, streamInfo.reconnectCount - 1) setTimeout(() => { streamInfo.ws = null const newWs = new WebSocket(streamInfo.address) streamInfo.ws = newWs newWs.onopen = () => { streamInfo.isActive = true startHeartbeat(streamInfo) } newWs.onerror = () => triggerReconnect(streamInfo) newWs.onclose = () => triggerReconnect(streamInfo) newWs.onpong = () => pongReceived = true newWs.onmessage = (event) => { const data = JSON.parse(event.data) // 复用原本的K线数据处理逻辑 } }, delay) }
使用示例
// 创建BTC、ETH的小时K线流 createStream('btcusdt') createStream('ethusdt') // 随时查询当前活跃流 console.log('当前活跃交易对流:', getActiveStreams())
内容的提问来源于stack exchange,提问作者James Butler
相关产品推荐
相关产品推荐

