RTK Query изначально ориентирован на работу с HTTP-запросами, где основная модель взаимодействия строится вокруг query/mutation и кэширования результатов. WebSocket меняет характер обмена данными: соединение становится постоянным, а данные поступают асинхронным потоком событий. Это требует иной архитектуры обновления кэша и управления подписками.
В RTK Query нет встроенного транспортного уровня для WebSocket, но предусмотрены расширения жизненного цикла запросов, позволяющие интегрировать потоковые данные без выхода за пределы архитектуры Redux Toolkit.
Ключевая точка интеграции — onCacheEntryAdded,
позволяющая подключить WebSocket при появлении первой подписки и
корректно отключить его при отсутствии активных слушателей.
Основная схема строится вокруг 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 вызывается при первой подписке на
endpointupdateCachedData модифицирует кэш напрямую через
Immer-слойcacheEntryRemoved завершает жизненный цикл
соединенияЭто обеспечивает автоматическое управление ресурсами без ручного контроля подписок.
WebSocket не должен открываться на каждый компонент отдельно. RTK
Query гарантирует дедупликацию запросов, но соединение внутри
onCacheEntryAdded требует дополнительного контроля, если
endpoint активно используется в нескольких местах.
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 передаёт разные типы событий.
Простейшая схема — использование поля 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 в реальных условиях требует устойчивости к разрывам соединения.
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()
};
}
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, после чего включается поток обновлений.
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 может выступать триггером для полной или частичной инвалидации 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 не возвращает коллекцию данных, а работает как поток событий без начальной выборки. В таком случае 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-менеджер.
В более сложных архитектурах 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 необходимо синхронизировать с состоянием RTK Query.
ws.oner ror = () => {
updateCachedData((draft) => {
draft.wsError = true;
});
};
или через отдельное поле состояния:
updateCachedData((draft) => {
draft.status = 'disconnected';
});
Это позволяет UI реагировать на состояние соединения через тот же механизм, что и данные.
На практике наиболее устойчивая архитектура включает три уровня:
Синхронизация обеспечивается через:
updateCachedData для потоковых измененийinvalidateTags для пересборки данныхonCacheEntryAdded как точка жизненного цикла
соединенияТакая модель устраняет необходимость ручного state management и сохраняет консистентность между сервером и клиентом в реальном времени