如何实现Polygon WebSocket经Node.js中继至多React前端客户端
Node.js WebSocket 中继服务实现方案
核心逻辑:服务端维持与Polygon的单条持久WebSocket连接,统一管理所有前端客户端的订阅关系,收到行情推送后按需转发给对应前端,API密钥仅存于服务端,无泄露风险。
前置依赖
服务端使用轻量WebSocket库ws实现,先安装依赖:
npm install ws
建议把Polygon API密钥存到环境变量中,不要硬编码在代码里,启动服务时可通过POLYGON_API_KEY=你的密钥 node server.js的方式传入。
服务端完整代码
const { WebSocketServer, WebSocket } = require('ws'); const PORT = 8080; // 服务端监听端口,可自行修改 const POLYGON_WS_URL = 'wss://delayed.polygon.io/stocks'; const POLYGON_API_KEY = process.env.POLYGON_API_KEY; // 全局状态管理 let polygonSocket = null; // 存储所有前端客户端连接 const connectedClients = new Set(); // 存储订阅关系:key为订阅频道(比如A.AAPL),value为订阅该频道的客户端Set const channelSubscribers = new Map(); // 记录当前已经向Polygon成功订阅的频道,避免重复发送订阅请求 const activePolygonSubscriptions = new Set(); // 初始化与Polygon的WebSocket连接 function initPolygonConnection() { polygonSocket = new WebSocket(POLYGON_WS_URL); polygonSocket.on('open', () => { console.log('已连接到Polygon WebSocket服务'); // 连接建立后自动完成鉴权 polygonSocket.send(JSON.stringify({ action: 'auth', params: POLYGON_API_KEY })); }); polygonSocket.on('message', (rawData) => { const messages = JSON.parse(rawData.toString()); for (const msg of messages) { // 鉴权成功后,恢复之前所有有效订阅 if (msg.status === 'auth_success') { console.log('Polygon鉴权成功'); if (activePolygonSubscriptions.size > 0) { polygonSocket.send(JSON.stringify({ action: 'subscribe', params: Array.from(activePolygonSubscriptions).join(',') })); } continue; } // 行情消息,转发给对应订阅的客户端 if (msg.ev === 'A') { const channel = `A.${msg.sym}`; const subscribers = channelSubscribers.get(channel); if (subscribers) { const sendData = JSON.stringify([msg]); subscribers.forEach(client => { if (client.readyState === WebSocket.OPEN) { client.send(sendData); } }); } } } }); // 断线自动重连 polygonSocket.on('close', () => { console.log('Polygon连接断开,5秒后重连'); setTimeout(initPolygonConnection, 5000); }); polygonSocket.on('error', (err) => { console.error('Polygon连接错误:', err); polygonSocket.close(); }); } // 启动面向前端的WebSocket服务 const wss = new WebSocketServer({ port: PORT }); console.log(`WebSocket中继服务已启动,监听端口${PORT}`); wss.on('connection', (clientSocket) => { connectedClients.add(clientSocket); // 存储当前客户端订阅的频道,断开时用于清理 const clientSubscriptions = new Set(); clientSocket.on('message', (rawData) => { try { const msg = JSON.parse(rawData.toString()); // 处理前端订阅请求 if (msg.action === 'subscribe') { const channels = msg.params.split(','); channels.forEach(channel => { if (clientSubscriptions.has(channel)) return; clientSubscriptions.add(channel); // 更新频道订阅者列表 if (!channelSubscribers.has(channel)) { channelSubscribers.set(channel, new Set()); } channelSubscribers.get(channel).add(clientSocket); // 若该频道未向Polygon订阅,发起订阅 if (!activePolygonSubscriptions.has(channel) && polygonSocket.readyState === WebSocket.OPEN) { activePolygonSubscriptions.add(channel); polygonSocket.send(JSON.stringify({ action: 'subscribe', params: channel })); } }); } // 处理前端取消订阅请求 if (msg.action === 'unsubscribe') { const channels = msg.params.split(','); channels.forEach(channel => { if (!clientSubscriptions.has(channel)) return; clientSubscriptions.delete(channel); const subscribers = channelSubscribers.get(channel); if (subscribers) { subscribers.delete(clientSocket); // 频道无订阅者时,向Polygon取消订阅节省流量 if (subscribers.size === 0) { channelSubscribers.delete(channel); if (activePolygonSubscriptions.has(channel) && polygonSocket.readyState === WebSocket.OPEN) { activePolygonSubscriptions.delete(channel); polygonSocket.send(JSON.stringify({ action: 'unsubscribe', params: channel })); } } } }); } } catch (e) { console.error('解析前端消息错误:', e); } }); // 客户端断开时自动清理订阅关系 clientSocket.on('close', () => { connectedClients.delete(clientSocket); clientSubscriptions.forEach(channel => { const subscribers = channelSubscribers.get(channel); if (subscribers) { subscribers.delete(clientSocket); if (subscribers.size === 0) { channelSubscribers.delete(channel); if (activePolygonSubscriptions.has(channel) && polygonSocket.readyState === WebSocket.OPEN) { activePolygonSubscriptions.delete(channel); polygonSocket.send(JSON.stringify({ action: 'unsubscribe', params: channel })); } } } }); }); }); // 服务启动时初始化Polygon连接 initPolygonConnection();
React 前端改造
将原来直连Polygon的地址替换为自建中继服务地址即可,无需再传输API密钥,鉴权逻辑完全由服务端处理:
let socket; useEffect(() => { if (marketStatus() !== 'closed') { // 替换为自己的中继服务地址,本地测试为ws://localhost:8080 socket = new WebSocket('ws://你的服务地址:8080'); socket.onopen = () => { // 无需发送鉴权请求,直接发送订阅指令即可 socket.send(JSON.stringify({ action: 'subscribe', params: `A.${props.ticker}` })); }; socket.onmessage = (event) => { const data = JSON.parse(event.data); if (data[0].ev === 'A') { console.log(event.data) // 保留原有行情处理逻辑即可 } }; return () => { // 断开时可主动发送取消订阅请求,服务端也会自动清理 if (socket.readyState === WebSocket.OPEN) { socket.send(JSON.stringify({ action: 'unsubscribe', params: `A.${props.ticker}` })); } socket.close(); }; } }, [props.ticker]);
实现说明
- 服务端与Polygon全程仅维持1条连接,完全符合平台连接数限制
- Polygon API密钥仅存储在服务端环境变量,不会泄露到前端
- 自动维护订阅关系,无客户端订阅的标的会自动向Polygon取消订阅,减少无效流量消耗
- 内置Polygon连接断线重连逻辑,重连后自动恢复所有有效订阅,无需人工干预
- 支持任意数量前端同时连接,仅向前端推送其订阅的标的数据,无冗余消息传输
内容的提问来源于stack exchange,提问作者Steven Talafous
相关产品推荐
相关产品推荐

