NestJS对接币安WebSocket API请求挂起、前端仅收首条数据问题求助
问题解决方案
1 初始HTTP请求一直pending的原因
你最初在HTTP接口中对接币安WebSocket时,因为HTTP是单次请求-响应模型,你没有主动调用响应对象的send/json方法结束请求,同时币安WebSocket是持续推送的长连接,没有主动断开的逻辑,导致请求一直处于挂起状态。你后续改用Socket.IO网关做双向推送的方案符合这类实时数据推送的场景,是正确的技术选型。
2 前端仅能收到第一条推送的核心原因
- Socket.IO的
emit确认回调(即你emit第三个参数传入的函数)是单次响应机制,只能接收服务器返回的第一次数据,后续Observable推送的新值不会再触发这个回调。 - NestJS默认对于
@SubscribeMessage装饰的方法返回的Observable,只会将第一个next值通过ack回调返回给客户端,后续流数据没有主动推送的逻辑。 - 客户端没有监听服务器主动推送的自定义事件,后续的新数据没有对应的接收入口。
- 你当前的代码每次收到客户端
events事件都会新建一个币安WebSocket连接,多客户端访问时会触发币安的频率限制,也会造成资源浪费。
修复代码
服务端修改(Coin.gateway.ts)
将主动推送逻辑替换默认的Observable返回逻辑,复用同一个币安连接减少资源消耗:
import { MessageBody, SubscribeMessage, WebSocketGateway, WebSocketServer, ConnectedSocket } from '@nestjs/websockets'; import { Server, Socket } from 'socket.io'; import { map } from 'rxjs'; import { Coin } from './classes/coin'; import * as coinlist from './list/coins.json' @WebSocketGateway(811, {transports: ['websocket', 'polling'], cors: true}) export class CoinGateway { @WebSocketServer() server: Server; // 全局复用单个币安连接实例 private coinInstance: Coin; @SubscribeMessage('events') handleMessage(@MessageBody() data: any, @ConnectedSocket() client: Socket) { console.log('data',data) // 仅首次初始化币安连接 if (!this.coinInstance) { this.coinInstance = new Coin(coinlist, 'usdt', 'miniTicker') // 订阅币安数据流,有新数据就主动推送给所有连接的客户端 this.coinInstance.getCryptoData().pipe(map((c) => c)).subscribe((data) => { this.server.emit('binanceData', data) }) } // 给emit的ack回调返回订阅成功状态 return { status: 'success', msg: '已成功订阅币安实时数据' } } }
客户端修改(useEffect钩子)
新增监听服务端推送的binanceData事件,新增组件卸载时的断开连接逻辑避免内存泄漏:
useEffect(() => { const socket = io('ws://localhost:811', {transports: ['websocket']}) socket.on('connect', () => { console.log('Connection established from client') socket.emit('events', '', (res: any) => { console.log('订阅结果:', res) }) // 监听服务端主动推送的币安数据,每次有新数据都会触发 socket.on('binanceData', (data) => { console.log('收到币安实时数据:', data) // 此处编写业务逻辑更新UI即可 }) const engine = socket.io.engine; console.log(engine.transport.name); // in most cases, prints "polling" engine.once("upgrade", () => { // called when the transport is upgraded (i.e. from HTTP long-polling to WebSocket) console.log(engine.transport.name); // in most cases, prints "websocket" }); engine.on("packetCreate", ({ type, data }) => { // called for each packet sent console.log('Stype', type) console.log('Sdata', data) }); }) // 组件卸载时主动断开连接 return () => { socket.disconnect() } }, [])
可选优化点
- 可以在
Coin类的Observable中添加重连逻辑,币安WebSocket异常断开时自动重连,无需手动重启服务。 - 如果需要按客户端订阅不同的交易对,可以按客户端ID维护订阅列表,单独给对应客户端推送数据,不需要全量广播。
内容的提问来源于stack exchange,提问作者Timur
相关产品推荐
相关产品推荐

