如何实现两个独立WebSocket协同?Node.js对接Kraken与Angular
问题分析与解决方案
你的代码核心问题有两个:
- 重复初始化Socket.IO服务器:在Kraken的
message事件中每次都创建新的Socket.IO实例,这会导致大量冗余服务器对象堆积,直接引发内存泄漏,而且新实例无法和已连接的客户端通信。 - 事件触发逻辑错误:
io.on('connection')只有在客户端首次连接时才会触发,你在Kraken消息回调里绑定这个事件,只会在收到第一条消息后处理后续的客户端连接,无法把当前消息推送给已在线的客户端。
正确实现方案
步骤1:全局初始化Socket.IO服务器
Socket.IO服务器应该在Node.js服务启动时只初始化一次,绑定到HTTP服务器实例上,确保所有客户端共用同一个服务实例。
步骤2:独立处理Kraken WebSocket连接
单独建立Kraken的WebSocket连接,在收到消息时直接通过已初始化的Socket.IO实例广播数据给所有客户端。
完整示例代码
const express = require('express'); const http = require('http'); const { Server } = require('socket.io'); const WebSocket = require('ws'); const app = express(); const server = http.createServer(app); // 全局初始化Socket.IO服务器(仅执行一次) const io = new Server(server, { cors: { origin: "http://localhost:4200", methods: ["GET", "POST"], allowedHeaders: ["*"], credentials: true, }, }); // 监听客户端连接/断开事件 io.on('connection', (socket) => { console.log('用户已连接'); socket.on('disconnect', () => { console.log('用户已断开连接'); }); }); // 初始化Kraken WebSocket连接 const krakenWsUrl = 'wss://ws.kraken.com'; const krakenWs = new WebSocket(krakenWsUrl); krakenWs.on('open', () => { console.log('已连接到Kraken WebSocket'); // 发送订阅请求(示例为订阅BTC/USD交易数据) const subscription = JSON.stringify({ event: 'subscribe', pair: ['XBT/USD'], subscription: { name: 'trade' } }); krakenWs.send(subscription); }); krakenWs.on('message', (wsMsg) => { try { const data = JSON.parse(wsMsg); // 过滤Kraken的订阅确认消息,只处理交易数据 if (Array.isArray(data) && data[1]) { const changes = parseTrades(data); // 你的数据处理函数 // 广播数据给所有已连接的前端客户端 io.emit('kraken-trade-update', changes); console.log('已转发交易数据:', changes); } } catch (err) { console.error('解析Kraken消息失败:', err); } }); krakenWs.on('close', () => { console.log('Kraken WebSocket连接关闭'); // 可选:添加重连逻辑 }); krakenWs.on('error', (err) => { console.error('Kraken WebSocket错误:', err); }); // 启动Node.js服务器 const PORT = 3000; server.listen(PORT, () => { console.log(`服务器运行在 http://localhost:${PORT}`); });
Angular前端接收示例
在Angular组件中监听Socket.IO事件:
import { Component, OnInit, OnDestroy } from '@angular/core'; import { io, Socket } from 'socket.io-client'; @Component({ selector: 'app-trade-monitor', templateUrl: './trade-monitor.component.html' }) export class TradeMonitorComponent implements OnInit, OnDestroy { private socket: Socket; tradeData: any[] = []; ngOnInit(): void { // 连接到Node.js的Socket.IO服务 this.socket = io('http://localhost:3000'); // 监听交易数据更新事件 this.socket.on('kraken-trade-update', (data) => { this.tradeData = [...this.tradeData, ...data]; // 根据业务需求渲染数据 }); } ngOnDestroy(): void { // 组件销毁时断开Socket连接 this.socket.disconnect(); } }
关键注意事项
- 避免重复初始化:Socket.IO和Kraken WebSocket都应该只初始化一次,不要在事件回调中重复创建实例。
- 数据过滤:Kraken会返回订阅确认、心跳等非交易消息,需要过滤后再转发。
- 错误处理:添加try-catch捕获JSON解析错误,避免单个异常导致服务崩溃。
- 重连机制:可以给Kraken WebSocket添加自动重连逻辑,防止连接中断后数据停止推送。
内容的提问来源于stack exchange,提问作者user6680
相关产品推荐
相关产品推荐

