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

如何实现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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 18:06:27