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

使用Binance API订阅组合数据流遇限:请求377个仅成功320个

币安WebSocket组合流订阅数量不足问题排查与修复

问题核心

尝试订阅币安现货377个BUSD交易对的1分钟K线组合流,仅成功建立320个数据流,查阅官方文档无明确数量限制,确认是实现层面问题。

代码问题分析

  1. 单WebSocket连接订阅超限:币安单WebSocket连接存在隐性订阅数量上限(实际测试单连接最多支持约200-300个流订阅),一次性提交377个订阅请求会被服务器部分拒绝,导致部分流无法建立。
  2. 同步文件读写阻塞事件循环:消息处理逻辑中使用fs.readFileSync和fs.writeFileSync同步读写文件,会阻塞Node.js事件循环,导致WebSocket消息接收不及时,部分订阅确认或数据流消息丢失,影响计数准确性。
  3. 字符串排序逻辑错误:contentParsed.sort((a, b) => a.s - b.s)中,交易对符号是字符串类型,直接相减会得到NaN,排序逻辑失效,但不影响订阅数量。

解决方案

1. 拆分WebSocket连接

将订阅列表分成多个批次,每个批次的订阅数控制在200以内,创建多个WebSocket连接分别订阅不同批次的流。

2. 替换同步文件操作为异步

使用fs.promises.readFile和fs.promises.writeFile异步读写文件,避免阻塞事件循环,确保消息处理流畅。

3. 修复排序逻辑

使用字符串的localeCompare方法实现正确的字典序排序。

修改后的代码示例

import { busdAllSymbols } from './symbols'
var cors = require('cors')
import fs from 'fs/promises'
import path from 'path'

app.use(cors())

app.get('/exchangeInfo', async (req: Request, res: Response) => {
    const stableCoinName = req.query.stableCoinName as string
    const exchangeInfo = await axios.get('https://api.binance.com/api/v3/exchangeInfo')
    const symbols = exchangeInfo.data.symbols.map((el: any) => el.symbol.toLowerCase())
    const symbolsFiltered = symbols.filter((name: string) => {
        return name.substring(name.length - stableCoinName.length) === stableCoinName.toLowerCase()
    })
    res.send(symbolsFiltered)
})

const server = app.listen(8081, () => {
    console.log('server running')
})

const wss = new WebSocketServer({ server })
const clients = new Set<WebSocket.WebSocket>()

wss.on('connection', (client) => {
    console.log('client is connected')
    clients.add(client)

    client.on('message', () => console.log('message received'))

    client.on('close', () => {
        clients.delete(client)
    })
})

// 拆分订阅批次,每200个流一个连接
const BATCH_SIZE = 200
const batches = []
for (let i = 0; i < busdAllSymbols.length; i += BATCH_SIZE) {
    batches.push(busdAllSymbols.slice(i, i + BATCH_SIZE))
}

// 为每个批次创建WebSocket连接
batches.forEach((batch, index) => {
    const streamBinance = new WebSocket('wss://stream.binance.com:9443/ws')
    streamBinance.on('open', () => {
        console.log(`stream ${index} opened`)
        const subs = {
            method: 'SUBSCRIBE',
            params: batch.map((symbol) => `${symbol}@kline_1m`),
            id: index + 1, // 每个连接使用唯一ID
        }
        streamBinance.send(JSON.stringify(subs))
    })

    streamBinance.on('message', async (message) => {
        const candle = JSON.parse(message.toString())
        // 异步读写文件
        try {
            const contentString = await fs.readFile(path.join(__dirname, 'data.txt'), 'utf8')
            let contentParsed = JSON.parse(contentString) as any[]
            
            const existingIndex = contentParsed.findIndex(el => el.s === candle.s)
            if (existingIndex !== -1) {
                contentParsed[existingIndex] = candle
            } else {
                contentParsed.push(candle)
            }
            // 修复排序逻辑
            contentParsed.sort((a, b) => a.s.localeCompare(b.s))
            await fs.writeFile(path.join(__dirname, 'data.txt'), JSON.stringify(contentParsed))

            // 推送给客户端
            clients.forEach((client) => {
                if (client.readyState === WebSocket.OPEN) {
                    client.send(JSON.stringify(candle))
                }
            })
        } catch (err) {
            console.error('File operation error:', err)
            // 首次文件不存在时初始化
            if ((err as NodeJS.ErrnoException).code === 'ENOENT') {
                await fs.writeFile(path.join(__dirname, 'data.txt'), JSON.stringify([candle]))
            }
        }
    })

    streamBinance.on('error', (err) => {
        console.error(`stream ${index} error:`, err)
    })
})

验证步骤

  1. 启动服务后,查看控制台输出的多个stream X opened日志,确认所有批次连接成功建立。
  2. 检查币安返回的订阅确认消息(每个连接会返回包含result字段的确认),统计成功订阅的流数量是否与377一致。
  3. 观察data.txt中的数据条目数量,确认最终能收集到所有377个交易对的K线数据。

内容的提问来源于stack exchange,提问作者jpxcar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 23:07:10