使用Binance API订阅组合数据流遇限:请求377个仅成功320个
币安WebSocket组合流订阅数量不足问题排查与修复
问题核心
尝试订阅币安现货377个BUSD交易对的1分钟K线组合流,仅成功建立320个数据流,查阅官方文档无明确数量限制,确认是实现层面问题。
代码问题分析
- 单WebSocket连接订阅超限:币安单WebSocket连接存在隐性订阅数量上限(实际测试单连接最多支持约200-300个流订阅),一次性提交377个订阅请求会被服务器部分拒绝,导致部分流无法建立。
- 同步文件读写阻塞事件循环:消息处理逻辑中使用
fs.readFileSync和fs.writeFileSync同步读写文件,会阻塞Node.js事件循环,导致WebSocket消息接收不及时,部分订阅确认或数据流消息丢失,影响计数准确性。 - 字符串排序逻辑错误:
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) }) })
验证步骤
- 启动服务后,查看控制台输出的多个
stream X opened日志,确认所有批次连接成功建立。 - 检查币安返回的订阅确认消息(每个连接会返回包含
result字段的确认),统计成功订阅的流数量是否与377一致。 - 观察
data.txt中的数据条目数量,确认最终能收集到所有377个交易对的K线数据。
内容的提问来源于stack exchange,提问作者jpxcar
相关产品推荐
相关产品推荐

