You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何查询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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.24 23:33:00