Как распределять данные из одного источника по WebSocket

Короткий ответ: если данные приходят из одного источника, не нужно открывать отдельное подключение к источнику для каждого WebSocket-клиента. Сделайте один центральный обработчик: он получает входящие данные, нормализует их, смотрит подписки клиентов и рассылает только тем, кому событие действительно нужно

Рабочая модель называется по-разному: hub, broker, dispatcher, fan-out. Смысл один: один входящий поток, много исходящих WebSocket-подключений

Как должна выглядеть схема

Пример:

Источник данных
      |
      v
Сервер WebSocket
      |
      +--> клиент A подписан на prices
      +--> клиент B подписан на news
      +--> клиент C подписан на prices и news

Если источник прислал событие prices, его получают A и C. Клиент B не получает лишние данные

Формат клиентской подписки

Клиент может отправить серверу сообщение:

socket.send(JSON.stringify({
  type: "subscribe",
  topic: "prices",
}));

И отписаться:

socket.send(JSON.stringify({
  type: "unsubscribe",
  topic: "prices",
}));

Так сервер понимает, кому что отправлять

Пример сервера

Установите ws:

npm init -y
npm install ws

server.js:

const { WebSocketServer } = require("ws");

const wss = new WebSocketServer({ port: 8080 });
const subscriptions = new Map();

function getTopics(socket) {
  if (!subscriptions.has(socket)) {
    subscriptions.set(socket, new Set());
  }

  return subscriptions.get(socket);
}

wss.on("connection", (socket) => {
  getTopics(socket);

  socket.on("message", (message) => {
    const data = JSON.parse(message.toString());
    const topics = getTopics(socket);

    if (data.type === "subscribe") {
      topics.add(data.topic);
      socket.send(JSON.stringify({ type: "subscribed", topic: data.topic }));
    }

    if (data.type === "unsubscribe") {
      topics.delete(data.topic);
      socket.send(JSON.stringify({ type: "unsubscribed", topic: data.topic }));
    }
  });

  socket.on("close", () => {
    subscriptions.delete(socket);
  });
});

function broadcast(topic, payload) {
  for (const client of wss.clients) {
    const topics = subscriptions.get(client);

    if (client.readyState === client.OPEN && topics?.has(topic)) {
      client.send(JSON.stringify({ type: "event", topic, payload }));
    }
  }
}

setInterval(() => {
  broadcast("prices", {
    symbol: "BTC",
    value: Math.round(60000 + Math.random() * 1000),
  });
}, 1000);

console.log("WebSocket hub: ws://localhost:8080");

Здесь setInterval имитирует один источник данных. В реальности на этом месте может быть API, очередь сообщений, база, биржевой поток или внутренний сервис

Почему нельзя рассылать все всем

Рассылка всех событий всем клиентам быстро становится дорогой. Клиент получает лишние данные, браузер тратит CPU, сеть забивается, а сервер держит большие буферы отправки

Грамотная раздача строится на трех правилах:

  • у каждого события есть topic или тип;
  • у каждого клиента есть список подписок;
  • сервер проверяет readyState и не шлет в закрытые соединения.

Что делать с быстрым источником

Если источник присылает 1000 событий в секунду, не всегда надо отправлять каждое событие каждому клиенту. Возможные приемы:

  • throttle: отправлять не чаще раза в N миллисекунд;
  • snapshot: отправлять последнее состояние, а не каждое изменение;
  • batching: собирать несколько событий в один пакет;
  • фильтр по подписке: отправлять только нужную категорию.

Для дашборда часто достаточно отправлять последний снимок раз в секунду. Для торгового терминала требования могут быть жестче

Частые ошибки

Первая ошибка — создать новое подключение к источнику на каждого WebSocket-клиента. Так источник получает лишнюю нагрузку, а сервер теряет контроль над очередями

Вторая ошибка — забыть удалить клиента из подписок при close. Тогда в памяти останутся старые Set и Map-записи

Третья ошибка — отправлять данные без проверки readyState. Закрытые или закрывающиеся соединения могут ломать отправку

Четвертая ошибка — не ограничивать частоту. WebSocket не делает магического backpressure для вашей бизнес-логики: если клиент не успевает обрабатывать поток, приложение начнет тормозить

Самопроверка

Откройте два клиента. Первый подпишите на prices, второй не подписывайте. Если первый получает события, а второй молчит, распределение работает

Потом подпишите второй клиент на prices и убедитесь, что он тоже начал получать события. Закройте вкладку и проверьте, что сервер не продолжает хранить подписку закрытого клиента

Что почитать дальше по WebSocket

Если нужен общий маршрут по теме, откройте рубрику WebSocket. Для соседних задач пригодятся эти разборы:

Оцените статью
0 0 голоса
Рейтинг статьи
Подписаться
Уведомить о
guest

0 комментариев
Старые
Новые Популярные
Межтекстовые Отзывы
Посмотреть все комментарии
0
Оставьте комментарий! Напишите, что думаете по поводу статьи.x