Подписки на события

RTK Query изначально проектировался как инструмент для работы с HTTP-запросами, но его внутренняя архитектура опирается на единый стор Redux, что позволяет расширять модель взаимодействия с сервером до более сложных сценариев. Одним из таких сценариев являются подписки на события — механизм, позволяющий получать данные не только через запросы, но и через потоковые обновления: WebSocket, Server-Sent Events или кастомные event-driven каналы.

Подписки в RTK Query не являются отдельной сущностью уровня API, как query или mutation. Они реализуются через комбинацию cache lifecycle, updateQueryData, onCacheEntryAdded и интеграцию с внешними источниками событий. Это делает систему гибкой, но требует понимания внутреннего поведения кэша.


Базовая концепция подписок через жизненный цикл кэша

Каждый запрос в RTK Query создаёт запись в кэше. Эта запись существует, пока на неё есть активные подписчики (components, hooks или manual subscriptions). Когда последний подписчик исчезает, RTK Query может удалить кэш или оставить его в зависимости от настроек keepUnusedDataFor.

Ключевой момент:

подписка на данные = подписка на кэш

const { data } = useGetMessagesQuery(roomId);

В этом примере компонент автоматически подписывается на cache entry getMessages(roomId). Пока компонент активен — данные считаются «живыми».

Однако этот механизм работает только для HTTP-запросов. Для событийных потоков требуется расширение через onCacheEntryAdded.


onCacheEntryAdded как основа event-based подписок

onCacheEntryAdded — это lifecycle-хук, который вызывается при создании первой подписки на конкретный cache key. Он позволяет подключать внешние источники данных и синхронизировать их с кэшем.

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

export const chatApi = createApi({
  reducerPath: 'chatApi',
  baseQuery: fetchBaseQuery({ baseUrl: '/api' }),
  endpoints: (builder) => ({
    getMessages: builder.query({
      query: (roomId) => `/rooms/${roomId}/messages`,

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

        const ws = new WebSocket(`wss://example.com/rooms/${roomId}`);

        const listener = (event) => {
          const message = JSON.parse(event.data);

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

        ws.addEventListener('message', listener);

        await cacheEntryRemoved;

        ws.removeEventListener('message', listener);
        ws.close();
      }
    })
  })
});

Механизм работы onCacheEntryAdded

Внутри RTK Query происходит следующая последовательность:

  1. При первом useQuery создаётся cache entry.
  2. Вызывается onCacheEntryAdded.
  3. Выполняется await cacheDataLoaded — ожидание первичного запроса.
  4. Подключается внешний источник событий.
  5. При каждом событии обновляется кэш через updateCachedData.
  6. Когда последний подписчик исчезает — срабатывает cacheEntryRemoved.
  7. Выполняется cleanup логика (закрытие WebSocket, отписка от SSE и т.д.).

Ключевой принцип:

жизненный цикл подписки полностью привязан к жизненному циклу cache entry


updateCachedData как механизм реактивного обновления

updateCachedData предоставляет доступ к текущим данным кэша через Immer-прокси. Это позволяет изменять структуру данных иммутабельно, но писать код как будто происходит мутация.

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

Особенности:

  • изменения немедленно отражаются во всех подписанных компонентах;
  • не требуется dispatch отдельных actions;
  • используется тот же механизм, что и для optimistic updates.

Подписки через WebSocket

Наиболее частый сценарий использования event-based подписок — чаты, торговые системы, live-дашборды.

getPriceUpdates: builder.query({
  query: (symbol) => `/prices/${symbol}`,

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

    const socket = new WebSocket(`wss://stream.example.com/prices/${symbol}`);

    socket.addEventListener('message', (event) => {
      const data = JSON.parse(event.data);

      updateCachedData((draft) => {
        draft.price = data.price;
        draft.timestamp = data.timestamp;
      });
    });

    await cacheEntryRemoved;

    socket.close();
  }
});

Особенность данного подхода заключается в том, что RTK Query не требует отдельного slice для WebSocket-данных. Кэш становится единственным источником истины.


Server-Sent Events как альтернатива WebSocket

SSE проще в реализации, но менее гибок. RTK Query одинаково хорошо поддерживает оба подхода.

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

  const eventSource = new EventSource(`/events/stream/${arg}`);

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

    updateCachedData((draft) => {
      draft.events.unshift(payload);
    });
  };

  cacheEntryRemoved.then(() => {
    eventSource.close();
  });
}

SSE особенно эффективен для:

  • уведомлений
  • логов
  • событий аудита
  • систем мониторинга

Проблема множественных подписок

RTK Query автоматически дедуплицирует запросы. Однако при использовании event streams важно понимать:

  • один cache entry = одна подписка на событие
  • несколько компонентов используют одну и ту же подписку
  • повторное подключение не происходит при каждом mount

Это предотвращает избыточные WebSocket-соединения.


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

Поведение подписок зависит от:

keepUnusedDataFor

createApi({
  keepUnusedDataFor: 60
});

Параметр определяет, сколько секунд кэш остаётся после исчезновения последнего подписчика.

При event-based подписках это критично:

  • слишком маленькое значение приведёт к частым reconnection
  • слишком большое увеличит расход памяти

Принудительное обновление через invalidation

Подписки и события не заменяют классическую модель инвалидации.

builder.mutation({
  query: (message) => ({
    url: '/messages',
    method: 'POST',
    body: message
  }),
  invalidatesTags: ['Messages']
});

Комбинация:

  • WebSocket для live обновлений
  • invalidation для синхронизации после мутаций

создаёт устойчивую гибридную систему.


Связь с tag-based системой кеширования

RTK Query использует tag-based invalidation для управления зависимостями данных.

Подписки через onCacheEntryAdded работают поверх этого механизма:

  • теги определяют когда данные пересчитываются
  • события определяют как данные обновляются в реальном времени

Потенциальные проблемы архитектуры

Утечки соединений

Если не обработать cacheEntryRemoved, соединение может остаться активным:

await cacheEntryRemoved;
socket.close();

Несинхронизированное состояние

События могут приходить быстрее, чем завершится initial fetch. Именно поэтому важно:

await cacheDataLoaded;

Дублирование сообщений

При WebSocket иногда приходят повторные события. RTK Query не решает это автоматически — требуется дедупликация на уровне updateCachedData.


Паттерн объединения REST и событий

Наиболее стабильная архитектура:

  1. initial query через HTTP
  2. дальнейшие обновления через WebSocket/SSE
  3. синхронизация через invalidation при мутациях
await cacheDataLoaded; // REST initial state

socket.onmess age = (event) => {
  updateCachedData((draft) => {
    mergeOrReplace(draft, JSON.parse(event.data));
  });
};

Использование abort-сигналов

RTK Query передаёт signal, который можно использовать для отмены операций:

async onCacheEntryAdded(
  arg,
  { signal }
) {
  const controller = new AbortController();

  signal.addEventListener('abort', () => {
    controller.abort();
  });
}

Это важно при интеграции с fetch-stream API или кастомными транспортами.


Роль подписок в масштабируемых приложениях

Подписки через RTK Query позволяют:

  • централизовать управление состоянием
  • избежать отдельного слоя state management для streaming данных
  • использовать Redux как реактивный кэш
  • уменьшить количество boilerplate-кода

Архитектура становится особенно эффективной в системах:

  • real-time аналитики
  • финансовых панелей
  • чат-приложений
  • IoT мониторинга

Ограничения модели

Несмотря на гибкость, существуют ограничения:

  • нет встроенного reconnect/backoff механизма
  • нет стандартизированной очереди событий
  • сложные топологии событий требуют внешнего менеджера
  • логика потоков находится внутри endpoint, а не в отдельном слое

Это делает RTK Query хорошим транспортным слоем, но не полноценным event streaming framework.