TypeScript 封装多交易所 WebSocket 实时代币行情

在行情看板、资产列表和交易详情页中,我们经常需要同时展示多个交易所的代币价格。最直接的做法是在每个组件里创建一个 WebSocket,但这很快会带来重复连接、格式不统一、组件卸载后仍在推送,以及断线后无法恢复等问题。

本文将封装一个可复用的 TypeScript 行情客户端。它具备以下能力:

  • 支持 Binance、OKX 和 Bybit 现货行情。
  • 将不同交易所的数据统一成同一种结构。
  • 同一交易所只建立一条连接,多个页面或组件共享。
  • 同一交易对支持多个回调,不会相互覆盖。
  • 最后一个调用方取消订阅时,才向交易所发送取消请求。
  • 断线后指数退避重连,并自动恢复已有订阅。
  • 提供心跳、错误回调和统一的资源释放方法。

本文只读取无需 API Key 的公开市场数据,不包含下单与账户操作。交易所协议可能调整,上线前应再次核对官方文档。


一、先统一对外接口

三个交易所的字段和交易对格式并不相同:Binance 和 Bybit 使用 BTCUSDT,OKX 使用 BTC-USDT。业务组件不应该知道这些差异,因此公共方法统一接收 BTC/USDT,并返回统一行情对象。

export type Exchange = "binance" | "okx" | "bybit"

export interface MarketTicker {
  exchange: Exchange
  symbol: string
  price: number
  bid?: number
  ask?: number
  change24h?: number
  timestamp: number
}

export type TickerListener = (ticker: MarketTicker) => void
export type ErrorListener = (error: Error) => void

change24h 统一使用百分比。例如 1.25 表示上涨 1.25%,而不是小数 0.0125


二、完整公共封装

可以把下面的代码保存为 lib/market-price-client.ts。它只能在浏览器环境使用;在 Next.js 中,应从带有 "use client" 的组件或自定义 Hook 中调用。

export type Exchange = "binance" | "okx" | "bybit"

export interface MarketTicker {
  exchange: Exchange
  symbol: string
  price: number
  bid?: number
  ask?: number
  change24h?: number
  timestamp: number
}

export type TickerListener = (ticker: MarketTicker) => void
export type ErrorListener = (error: Error) => void

type JsonRecord = Record<string, unknown>

interface ExchangeAdapter {
  url: string
  toExchangeSymbol(symbol: string): string
  toPublicSymbol(symbol: string): string
  subscribe(symbols: string[]): string
  unsubscribe(symbols: string[]): string
  parse(payload: unknown): MarketTicker | null
  heartbeat?: string
}

const asRecord = (value: unknown): JsonRecord | null =>
  typeof value === "object" && value !== null ? (value as JsonRecord) : null

const text = (value: unknown): string | undefined =>
  typeof value === "string" ? value : undefined

const finiteNumber = (value: unknown): number | undefined => {
  const number = typeof value === "number" ? value : Number(value)
  return Number.isFinite(number) ? number : undefined
}

const normalizePublicSymbol = (symbol: string): string => {
  const normalized = symbol.trim().toUpperCase().replace(/[-_]/g, "/")

  if (!/^[A-Z0-9]+\/[A-Z0-9]+$/.test(normalized)) {
    throw new Error(`Invalid symbol: ${symbol}. Use a value such as BTC/USDT.`)
  }

  return normalized
}

const compactSymbol = (symbol: string) =>
  normalizePublicSymbol(symbol).replace("/", "")

const dashedSymbol = (symbol: string) =>
  normalizePublicSymbol(symbol).replace("/", "-")

const binance: ExchangeAdapter = {
  url: "wss://stream.binance.com:9443/ws",
  toExchangeSymbol: (symbol) => compactSymbol(symbol).toLowerCase(),
  toPublicSymbol: (symbol) => symbol.toUpperCase(),
  subscribe: (symbols) =>
    JSON.stringify({
      method: "SUBSCRIBE",
      params: symbols.map((symbol) => `${symbol}@ticker`),
      id: Date.now(),
    }),
  unsubscribe: (symbols) =>
    JSON.stringify({
      method: "UNSUBSCRIBE",
      params: symbols.map((symbol) => `${symbol}@ticker`),
      id: Date.now(),
    }),
  parse: (payload) => {
    const data = asRecord(payload)
    if (!data || data.e !== "24hrTicker") return null

    const price = finiteNumber(data.c)
    const symbol = text(data.s)
    if (price === undefined || !symbol) return null

    return {
      exchange: "binance",
      symbol,
      price,
      bid: finiteNumber(data.b),
      ask: finiteNumber(data.a),
      change24h: finiteNumber(data.P),
      timestamp: finiteNumber(data.E) ?? Date.now(),
    }
  },
}

const okx: ExchangeAdapter = {
  url: "wss://ws.okx.com:8443/ws/v5/public",
  toExchangeSymbol: dashedSymbol,
  toPublicSymbol: (symbol) => symbol.replace("-", "/"),
  subscribe: (symbols) =>
    JSON.stringify({
      op: "subscribe",
      args: symbols.map((instId) => ({ channel: "tickers", instId })),
    }),
  unsubscribe: (symbols) =>
    JSON.stringify({
      op: "unsubscribe",
      args: symbols.map((instId) => ({ channel: "tickers", instId })),
    }),
  heartbeat: "ping",
  parse: (payload) => {
    const message = asRecord(payload)
    const rows = Array.isArray(message?.data) ? message.data : []
    const data = asRecord(rows[0])
    if (!data) return null

    const price = finiteNumber(data.last)
    const symbol = text(data.instId)
    if (price === undefined || !symbol) return null

    const open = finiteNumber(data.open24h)
    const change24h = open && open !== 0 ? ((price - open) / open) * 100 : undefined

    return {
      exchange: "okx",
      symbol: symbol.replace("-", "/"),
      price,
      bid: finiteNumber(data.bidPx),
      ask: finiteNumber(data.askPx),
      change24h,
      timestamp: finiteNumber(data.ts) ?? Date.now(),
    }
  },
}

const bybit: ExchangeAdapter = {
  url: "wss://stream.bybit.com/v5/public/spot",
  toExchangeSymbol: compactSymbol,
  toPublicSymbol: (symbol) => symbol,
  subscribe: (symbols) =>
    JSON.stringify({
      op: "subscribe",
      args: symbols.map((symbol) => `tickers.${symbol}`),
    }),
  unsubscribe: (symbols) =>
    JSON.stringify({
      op: "unsubscribe",
      args: symbols.map((symbol) => `tickers.${symbol}`),
    }),
  heartbeat: JSON.stringify({ op: "ping" }),
  parse: (payload) => {
    const message = asRecord(payload)
    if (!text(message?.topic)?.startsWith("tickers.")) return null

    const data = asRecord(message?.data)
    const price = finiteNumber(data?.lastPrice)
    const symbol = text(data?.symbol)
    if (price === undefined || !symbol) return null

    const ratio = finiteNumber(data?.price24hPcnt)

    return {
      exchange: "bybit",
      symbol,
      price,
      bid: finiteNumber(data?.bid1Price),
      ask: finiteNumber(data?.ask1Price),
      change24h: ratio === undefined ? undefined : ratio * 100,
      timestamp: finiteNumber(message?.ts) ?? Date.now(),
    }
  },
}

const adapters: Record<Exchange, ExchangeAdapter> = {
  binance,
  okx,
  bybit,
}

class ExchangeConnection {
  private socket: WebSocket | null = null
  private reconnectTimer: ReturnType<typeof setTimeout> | null = null
  private heartbeatTimer: ReturnType<typeof setInterval> | null = null
  private reconnectAttempts = 0
  private manuallyClosed = false
  private readonly listeners = new Map<string, Set<TickerListener>>()

  constructor(
    private readonly exchange: Exchange,
    private readonly adapter: ExchangeAdapter,
    private readonly onError?: ErrorListener,
  ) {}

  add(symbol: string, listener: TickerListener): () => void {
    const publicSymbol = normalizePublicSymbol(symbol)
    const exchangeSymbol = this.adapter.toExchangeSymbol(publicSymbol)
    const listeners = this.listeners.get(exchangeSymbol) ?? new Set<TickerListener>()
    const isFirstListener = listeners.size === 0

    listeners.add(listener)
    this.listeners.set(exchangeSymbol, listeners)
    this.manuallyClosed = false

    if (!this.socket || this.socket.readyState === WebSocket.CLOSED) {
      this.connect()
    } else if (isFirstListener && this.socket.readyState === WebSocket.OPEN) {
      this.socket.send(this.adapter.subscribe([exchangeSymbol]))
    }

    let active = true
    return () => {
      if (!active) return
      active = false
      this.remove(exchangeSymbol, listener)
    }
  }

  close(): void {
    this.manuallyClosed = true
    this.clearTimers()
    this.listeners.clear()
    this.socket?.close(1000, "Client disposed")
    this.socket = null
  }

  private connect(): void {
    if (typeof window === "undefined") {
      this.report(new Error("MarketPriceClient can only run in the browser."))
      return
    }

    if (
      this.socket?.readyState === WebSocket.OPEN ||
      this.socket?.readyState === WebSocket.CONNECTING
    ) {
      return
    }

    const socket = new WebSocket(this.adapter.url)
    this.socket = socket

    socket.onopen = () => {
      this.reconnectAttempts = 0
      const symbols = [...this.listeners.keys()]
      if (symbols.length > 0) socket.send(this.adapter.subscribe(symbols))
      this.startHeartbeat()
    }

    socket.onmessage = (event) => {
      if (event.data === "pong") return

      try {
        const ticker = this.adapter.parse(JSON.parse(String(event.data)))
        if (!ticker) return

        const exchangeSymbol = this.adapter.toExchangeSymbol(ticker.symbol)
        const normalizedTicker = {
          ...ticker,
          symbol: normalizePublicSymbol(
            this.adapter.toPublicSymbol(exchangeSymbol),
          ),
        }

        this.listeners
          .get(exchangeSymbol)
          ?.forEach((listener) => listener(normalizedTicker))
      } catch (error) {
        this.report(error)
      }
    }

    socket.onerror = () => {
      this.report(new Error(`${this.exchange} WebSocket error.`))
    }

    socket.onclose = () => {
      this.stopHeartbeat()
      this.socket = null

      if (!this.manuallyClosed && this.listeners.size > 0) {
        this.scheduleReconnect()
      }
    }
  }

  private remove(symbol: string, listener: TickerListener): void {
    const listeners = this.listeners.get(symbol)
    if (!listeners) return

    listeners.delete(listener)
    if (listeners.size > 0) return

    this.listeners.delete(symbol)
    if (this.socket?.readyState === WebSocket.OPEN) {
      this.socket.send(this.adapter.unsubscribe([symbol]))
    }
  }

  private scheduleReconnect(): void {
    if (this.reconnectTimer) return

    const delay = Math.min(1000 * 2 ** this.reconnectAttempts, 30_000)
    this.reconnectAttempts += 1
    this.reconnectTimer = setTimeout(() => {
      this.reconnectTimer = null
      this.connect()
    }, delay)
  }

  private startHeartbeat(): void {
    this.stopHeartbeat()
    if (!this.adapter.heartbeat) return

    this.heartbeatTimer = setInterval(() => {
      if (this.socket?.readyState === WebSocket.OPEN) {
        this.socket.send(this.adapter.heartbeat)
      }
    }, 20_000)
  }

  private stopHeartbeat(): void {
    if (this.heartbeatTimer) clearInterval(this.heartbeatTimer)
    this.heartbeatTimer = null
  }

  private clearTimers(): void {
    if (this.reconnectTimer) clearTimeout(this.reconnectTimer)
    this.reconnectTimer = null
    this.stopHeartbeat()
  }

  private report(error: unknown): void {
    this.onError?.(
      error instanceof Error ? error : new Error("Unknown WebSocket error."),
    )
  }
}

export class MarketPriceClient {
  private readonly connections = new Map<Exchange, ExchangeConnection>()

  constructor(private readonly onError?: ErrorListener) {}

  subscribe(
    exchange: Exchange,
    symbol: string,
    listener: TickerListener,
  ): () => void {
    let connection = this.connections.get(exchange)

    if (!connection) {
      connection = new ExchangeConnection(
        exchange,
        adapters[exchange],
        this.onError,
      )
      this.connections.set(exchange, connection)
    }

    return connection.add(symbol, listener)
  }

  dispose(): void {
    this.connections.forEach((connection) => connection.close())
    this.connections.clear()
  }
}

export const marketPriceClient = new MarketPriceClient((error) => {
  console.error("[market-price]", error)
})

三、为什么可以在多个地方调用

marketPriceClient 是模块级单例。同一个浏览器标签页无论有多少组件导入它,都使用同一实例。

内部通过两层结构管理订阅:

MarketPriceClient
  ├─ Binance connection
  │    ├─ BTCUSDT -> listener A, listener B
  │    └─ ETHUSDT -> listener C
  ├─ OKX connection
  └─ Bybit connection

调用 subscribe() 会返回一个取消函数。组件 A 取消订阅时只移除 A 的回调;只要组件 B 仍在监听同一交易对,底层订阅就会保留。


四、在 React / Next.js 中使用

可以再封装一个 Hook,供任意客户端组件调用:

"use client"

import { useEffect, useState } from "react"
import {
  marketPriceClient,
  type Exchange,
  type MarketTicker,
} from "@/lib/market-price-client"

export function useMarketPrice(exchange: Exchange, symbol: string) {
  const [ticker, setTicker] = useState<MarketTicker | null>(null)

  useEffect(() => {
    return marketPriceClient.subscribe(exchange, symbol, setTicker)
  }, [exchange, symbol])

  return ticker
}

在价格卡片中调用:

"use client"

import { useMarketPrice } from "@/hooks/use-market-price"

export function BitcoinPrice() {
  const ticker = useMarketPrice("binance", "BTC/USDT")

  if (!ticker) return <span>Loading...</span>

  return (
    <div>
      <strong>{ticker.exchange}</strong>
      <span>{ticker.symbol}: {ticker.price}</span>
      <span>24h: {ticker.change24h?.toFixed(2)}%</span>
    </div>
  )
}

同一页面也可以同时比较多个交易所:

const binanceTicker = useMarketPrice("binance", "BTC/USDT")
const okxTicker = useMarketPrice("okx", "BTC/USDT")
const bybitTicker = useMarketPrice("bybit", "BTC/USDT")

在 React 之外使用也很简单:

const unsubscribe = marketPriceClient.subscribe(
  "okx",
  "ETH/USDT",
  (ticker) => console.log(ticker.price),
)

// 不再需要行情时调用
unsubscribe()

通常不要在普通组件卸载时调用 marketPriceClient.dispose(),因为它会关闭所有页面共享的连接。只有在应用退出、用户登出且确定不再需要任何行情,或测试结束时才调用它。


五、生产环境还需要考虑什么

1. 行情推送不是持久化状态

WebSocket 断线期间会丢失消息。价格展示一般可以在重连后等待下一次推送;如果业务必须立即有初始值,可以先请求交易所 REST ticker,再用 WebSocket 增量更新。

2. 浏览器后台标签页会被节流

浏览器可能降低后台页面中定时器的执行频率。高可靠行情系统更适合由服务端维持交易所连接,再通过自己的 WebSocket 或 SSE 分发给前端。

3. 不要无限制订阅交易对

交易所有连接、消息频率和单连接订阅数量限制。大规模行情列表应该分批订阅、限制重订阅频率,并根据官方限制拆分连接。

4. 数字精度

示例为了方便展示使用 number。涉及资产结算、下单或精确金额计算时,不应直接使用浮点数,应保留交易所返回的字符串并使用 decimal.jsbig.js 等高精度方案。

5. 符号映射不能永远依靠字符串拼接

BTC/USDT 这类常见现货对可以直接转换,但生产项目最好定期读取各交易所的 instruments / exchange info 接口,建立明确的交易对映射,并过滤已下架或暂停交易的品种。

6. 行情不等于可成交价格

最新成交价、最优买价和最优卖价含义不同。大额交易还会受到盘口深度、滑点和手续费影响,因此行情展示结果不应表述为投资建议或成交保证。


六、总结

一个可复用的多交易所 WebSocket 封装,关键并不只是“成功收到价格”,而是要同时解决协议适配、连接复用、订阅生命周期、自动重连和资源释放。

本文的实现将 Binance、OKX 和 Bybit 的现货 ticker 转换为统一的 MarketTicker;用每个交易所一条连接承载多个交易对;用监听器集合让多个组件共享同一订阅;并通过取消函数、心跳和指数退避重连保证长期运行的稳定性。

业务层最终只需要记住一个入口:

const unsubscribe = marketPriceClient.subscribe(
  "binance",
  "BTC/USDT",
  (ticker) => console.log(ticker.price),
)

当需要增加新交易所时,只需实现新的 ExchangeAdapter,业务组件和统一数据结构都不需要跟着重写。这正是把交易所协议与页面展示解耦的核心价值。