Интеграция с WebSocket соединениями

RTK Query изначально ориентирован на работу с HTTP-запросами, где основная модель взаимодействия строится вокруг query/mutation и кэширования результатов. WebSocket меняет характер обмена данными: соединение становится постоянным, а данные поступают асинхронным потоком событий. Это требует иной архитектуры обновления кэша и управления подписками.

В RTK Query нет встроенного транспортного уровня для WebSocket, но предусмотрены расширения жизненного цикла запросов, позволяющие интегрировать потоковые данные без выхода за пределы архитектуры Redux Toolkit.

Ключевая точка интеграции — onCacheEntryAdded, позволяющая подключить WebSocket при появлении первой подписки и корректно отключить его при отсутствии активных слушателей.


Базовая архитектура WebSocket внутри createApi

Основная схема строится вокруг createApi, где endpoint не выполняет HTTP-запрос, а становится точкой подписки на поток данных.

import { createApi, fetchBaseQuery } from '@reduxjs/toolkit/query/react';

export const wsApi = createApi({
  reducerPath: 'wsApi',
  baseQuery: fetchBaseQuery({ baseUrl: '/api' }),
  endpoints: (builder) => ({
    getMessages: builder.query({
      query: () => 'messages',
      async onCacheEntryAdded(
        arg,
        { updateCachedData, cacheEntryRemoved }
      ) {
        const ws = new WebSocket('wss://example.com/ws');

        try {
          await new Promise((resolve) => {
            ws.ono pen = resolve;
          });

          ws.onmess age = (event) => {
            const data = JSON.parse(event.data);

            updateCachedData((draft) => {
              draft.push(data);
            });
          };

          await cacheEntryRemoved;
        } finally {
          ws.close();
        }
      }
    })
  })
});

Поведение механизма

  • onCacheEntryAdded вызывается при первой подписке на endpoint
  • updateCachedData модифицирует кэш напрямую через Immer-слой
  • cacheEntryRemoved завершает жизненный цикл соединения
  • WebSocket существует строго в рамках активного кэша

Это обеспечивает автоматическое управление ресурсами без ручного контроля подписок.


Управление жизненным циклом соединения

WebSocket не должен открываться на каждый компонент отдельно. RTK Query гарантирует дедупликацию запросов, но соединение внутри onCacheEntryAdded требует дополнительного контроля, если endpoint активно используется в нескольких местах.

Подход с singleton-соединением

let socket;

function getSocket() {
  if (!socket) {
    socket = new WebSocket('wss://example.com/ws');
  }
  return socket;
}

Однако в контексте RTK Query такой подход часто избыточен, поскольку жизненный цикл уже синхронизирован с кэшем.


Потоковая синхронизация данных с кэшем

Основная задача WebSocket в RTK Query — не просто получение сообщений, а синхронизация серверного потока с локальным кэшем.

Добавление новых элементов

ws.onmess age = (event) => {
  const message = JSON.parse(event.data);

  updateCachedData((draft) => {
    const exists = draft.find((m) => m.id === message.id);
    if (!exists) {
      draft.unshift(message);
    }
  });
};

Обновление существующих записей

updateCachedData((draft) => {
  const index = draft.findIndex((m) => m.id === message.id);
  if (index !== -1) {
    draft[index] = {
      ...draft[index],
      ...message
    };
  }
});

Удаление элементов

updateCachedData((draft) => {
  return draft.filter((m) => m.id !== message.id);
});

RTK Query позволяет работать с кэшем как с локальным mutable state, сохраняя при этом иммутабельность через Immer.


Разделение событий WebSocket по типам

В реальных приложениях WebSocket передаёт разные типы событий. Простейшая схема — использование поля type.

ws.onmess age = (event) => {
  const payload = JSON.parse(event.data);

  updateCachedData((draft) => {
    switch (payload.type) {
      case 'message_added':
        draft.push(payload.data);
        break;

      case 'message_updated':
        const idx = draft.findIndex(m => m.id === payload.data.id);
        if (idx !== -1) {
          draft[idx] = payload.data;
        }
        break;

      case 'message_deleted':
        return draft.filter(m => m.id !== payload.data.id);
    }
  });
};

Такой подход позволяет превращать WebSocket в событийный слой поверх RTK Query кэша.


Управление повторным подключением

WebSocket в реальных условиях требует устойчивости к разрывам соединения.

Базовый retry-механизм

function createSocket(onMessage) {
  let ws;
  let retries = 0;

  const connect = () => {
    ws = new WebSocket('wss://example.com/ws');

    ws.onmess age = onMessage;

    ws.oncl ose = () => {
      const timeout = Math.min(1000 * 2 ** retries, 30000);
      retries++;

      setTimeout(connect, timeout);
    };
  };

  connect();

  return {
    close: () => ws?.close()
  };
}

Интеграция в RTK Query

async onCacheEntryAdded(arg, { updateCachedData, cacheEntryRemoved }) {
  const socket = createSocket((event) => {
    const data = JSON.parse(event.data);

    updateCachedData((draft) => {
      draft.push(data);
    });
  });

  await cacheEntryRemoved;
  socket.close();
}

Синхронизация WebSocket и HTTP bootstrap

Часто WebSocket не является единственным источником данных. Первичная загрузка выполняется через HTTP, после чего включается поток обновлений.

getMessages: builder.query({
  query: () => 'messages',

  async onCacheEntryAdded(
    arg,
    { updateCachedData, cacheDataLoaded, cacheEntryRemoved }
  ) {
    await cacheDataLoaded;

    const ws = new WebSocket('wss://example.com/ws');

    ws.onmess age = (event) => {
      const msg = JSON.parse(event.data);

      updateCachedData((draft) => {
        draft.push(msg);
      });
    };

    await cacheEntryRemoved;
    ws.close();
  }
});

Ключевой момент — ожидание cacheDataLoaded, чтобы WebSocket не конфликтовал с начальным состоянием.


Инвалидация кэша через WebSocket события

WebSocket может выступать триггером для полной или частичной инвалидации RTK Query кэша.

ws.onmess age = (event) => {
  const data = JSON.parse(event.data);

  if (data.type === 'force_refresh') {
    dispatch(api.util.invalidateTags(['Messages']));
  }
};

При использовании providesTags и invalidatesTags появляется гибридная модель:

  • WebSocket сигнализирует об изменении
  • RTK Query пересобирает данные через HTTP
  • Кэш остаётся согласованным

Использование endpoint-level подписок вместо query

Иногда WebSocket не возвращает коллекцию данных, а работает как поток событий без начальной выборки. В таком случае endpoint можно моделировать как пустой query.

listenEvents: builder.query({
  queryFn: () => ({ data: [] }),

  async onCacheEntryAdded(
    arg,
    { cacheEntryRemoved }
  ) {
    const ws = new WebSocket('wss://example.com/events');

    ws.onmess age = (event) => {
      const eventData = JSON.parse(event.data);
      console.log(eventData);
    };

    await cacheEntryRemoved;
    ws.close();
  }
});

Здесь RTK Query используется исключительно как lifecycle-менеджер.


Совмещение с Redux Listener Middleware

В более сложных архитектурах WebSocket может быть вынесен из RTK Query и управляться через listenerMiddleware.

import { createListenerMiddleware } from '@reduxjs/toolkit';

const listenerMiddleware = createListenerMiddleware();

listenerMiddleware.startListening({
  actionCreator: wsApi.endpoints.getMessages.matchFulfilled,
  effect: async (action, listenerApi) => {
    const ws = new WebSocket('wss://example.com/ws');

    ws.onmess age = (event) => {
      listenerApi.dispatch(
        wsApi.util.updateQueryData(
          'getMessages',
          undefined,
          (draft) => {
            draft.push(JSON.parse(event.data));
          }
        )
      );
    };

    await listenerApi.delay(0);
  }
});

Такой подход отделяет транспортный слой от RTK Query lifecycle.


Оптимизация частоты обновлений

WebSocket может генерировать высокочастотный поток данных. Прямое обновление кэша на каждое событие приводит к лишним рендерам.

Батчинг обновлений

let buffer = [];

setInterval(() => {
  if (buffer.length === 0) return;

  updateCachedData((draft) => {
    draft.push(...buffer);
  });

  buffer = [];
}, 200);

ws.onmess age = (event) => {
  buffer.push(JSON.parse(event.data));
};

Это снижает нагрузку на React-рендеринг и Immer-процессинг.


Контроль конкурентных подписок

RTK Query может иметь несколько активных подписок на один endpoint. WebSocket должен учитывать, что данные обновляют общий кэш.

Критический принцип:

  • WebSocket привязывается к cache entry, а не к компоненту
  • один cache entry = одно соединение
  • несколько компонентов = общая подписка

Обработка ошибок соединения

Ошибки WebSocket необходимо синхронизировать с состоянием RTK Query.

ws.oner ror = () => {
  updateCachedData((draft) => {
    draft.wsError = true;
  });
};

или через отдельное поле состояния:

updateCachedData((draft) => {
  draft.status = 'disconnected';
});

Это позволяет UI реагировать на состояние соединения через тот же механизм, что и данные.


Гибридная модель RTK Query + WebSocket + HTTP

На практике наиболее устойчивая архитектура включает три уровня:

  • HTTP для начальной загрузки
  • WebSocket для инкрементальных обновлений
  • RTK Query кэш как единый источник состояния

Синхронизация обеспечивается через:

  • updateCachedData для потоковых изменений
  • invalidateTags для пересборки данных
  • onCacheEntryAdded как точка жизненного цикла соединения

Такая модель устраняет необходимость ручного state management и сохраняет консистентность между сервером и клиентом в реальном времени