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 — это 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();
}
})
})
});
Внутри RTK Query происходит следующая последовательность:
useQuery создаётся cache entry.onCacheEntryAdded.await cacheDataLoaded — ожидание первичного
запроса.updateCachedData.cacheEntryRemoved.Ключевой принцип:
жизненный цикл подписки полностью привязан к жизненному циклу cache entry
updateCachedData предоставляет доступ к текущим данным
кэша через Immer-прокси. Это позволяет изменять структуру данных
иммутабельно, но писать код как будто происходит мутация.
updateCachedData((draft) => {
draft.push(newMessage);
});
Особенности:
Наиболее частый сценарий использования 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-данных. Кэш становится единственным источником истины.
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 важно понимать:
Это предотвращает избыточные WebSocket-соединения.
Поведение подписок зависит от:
createApi({
keepUnusedDataFor: 60
});
Параметр определяет, сколько секунд кэш остаётся после исчезновения последнего подписчика.
При event-based подписках это критично:
Подписки и события не заменяют классическую модель инвалидации.
builder.mutation({
query: (message) => ({
url: '/messages',
method: 'POST',
body: message
}),
invalidatesTags: ['Messages']
});
Комбинация:
создаёт устойчивую гибридную систему.
RTK Query использует tag-based invalidation для управления зависимостями данных.
Подписки через onCacheEntryAdded работают поверх этого
механизма:
Если не обработать cacheEntryRemoved, соединение может
остаться активным:
await cacheEntryRemoved;
socket.close();
События могут приходить быстрее, чем завершится initial fetch. Именно поэтому важно:
await cacheDataLoaded;
При WebSocket иногда приходят повторные события. RTK Query не решает
это автоматически — требуется дедупликация на уровне
updateCachedData.
Наиболее стабильная архитектура:
await cacheDataLoaded; // REST initial state
socket.onmess age = (event) => {
updateCachedData((draft) => {
mergeOrReplace(draft, JSON.parse(event.data));
});
};
RTK Query передаёт signal, который можно использовать
для отмены операций:
async onCacheEntryAdded(
arg,
{ signal }
) {
const controller = new AbortController();
signal.addEventListener('abort', () => {
controller.abort();
});
}
Это важно при интеграции с fetch-stream API или кастомными транспортами.
Подписки через RTK Query позволяют:
Архитектура становится особенно эффективной в системах:
Несмотря на гибкость, существуют ограничения:
Это делает RTK Query хорошим транспортным слоем, но не полноценным event streaming framework.