在行情看板、资产列表和交易详情页中,我们经常需要同时展示多个交易所的代币价格。最直接的做法是在每个组件里创建一个 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.js、big.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,业务组件和统一数据结构都不需要跟着重写。这正是把交易所协议与页面展示解耦的核心价值。