Zobacz, jak wdrożyć SSE w Next.js przez Route Handlers, obsłużyć reconnect i dostarczać powiadomienia za pomocą Upstash Redis Pub/Sub.
Maciej Sala
Founder StriveLab
6 min czytaniaOpublikowano 11 kwietnia 2026 (Aktualizacja 23 lipca 2026)
Zakres: powiadomienia real-time przez SSE w Next.js
Metoda
Kierunek
Połączenie
Złożoność
Najlepsze dla
Polling
Klient → Serwer
Nowe per request
Niska
Proste dashboardy, statusy
SSE
Serwer → Klient
Persistent (HTTP)
Średnia
Powiadomienia, feed, logi
WebSockets
Dwukierunkowy
Persistent (WS)
Wysoka
Chat, gry, kolaboracja
Polling jako najprostsze podejście do powiadomień
Klient odpytuje serwer co N sekund. Zero utrzymywania połączenia, działa na każdej infrastrukturze.
Code
// hooks/use-polling.ts'use client'import { useEffect, useState } from 'react'export function usePolling<T>(url: string, intervalMs = 5000) { const [data, setData] = useState<T | null>(null) useEffect(() => { const controller = new AbortController() let active = true let timeout: ReturnType<typeof setTimeout> | undefined async function poll() { try { const response = await fetch(url, { signal: controller.signal }) if (!response.ok) throw new Error(`HTTP ${response.status}`) const result: T = await response.json() if (active) setData(result) } catch (error) { if (!controller.signal.aborted) { console.error('Polling error', error) } } finally { if (active) timeout = setTimeout(poll, intervalMs) } } void poll() return () => { active = false controller.abort() if (timeout) clearTimeout(timeout) } }, [url, intervalMs]) return data}
Code
// Sprawdzaj status zamówienia co 10 sekund.function OrderStatus({ orderId }: { orderId: string }) { const status = usePolling<{ status: string }>( `https://api.example.com/orders/${orderId}/status`, 10000, ) return <p>Status: {status?.status || 'Ładowanie...'}</p>}
Polling ma opóźnienie zależne od interwału i generuje zbędne żądania, gdy nic się nie zmieniło. Rekurencyjny setTimeout uruchamia kolejne sprawdzenie dopiero po zakończeniu poprzedniego, więc wolna odpowiedź nie powoduje nakładania żądań.
Parametr T opisuje oczekiwany wynik wyłącznie dla kompilatora. Jeśli odpowiedź pochodzi z niezaufanego API, odczytaj JSON jako unknown i zweryfikuj go parserem przed wywołaniem setData.
Server-Sent Events (SSE) dla powiadomień real-time
to jednokierunkowy stream HTTP. Serwer wysyła zdarzenia
do klienta przez otwarte połączenie. Przeglądarka obsługuje SSE natywnie
przez .
Route Handler jako endpoint SSE w Next.js
Wspólny typ transportowy powinien zawierać wyłącznie wartości serializowalne do JSON:
// app/api/notifications/stream/route.tsimport { auth } from '@/auth'import { getNotificationsAfter, getUnreadNotifications,} from '@/lib/notifications'import type { Notification } from '@/lib/notification-types'export const maxDuration = 300const encoder = new TextEncoder()function encodeNotification(notification: Notification) { return encoder.encode( [ `id: ${notification.id}`, 'event: notification', `data: ${JSON.stringify(notification)}`, '', '', ].join('\n'), )}export async function GET(request: Request) { const session = await auth() if (!session?.user) { return new Response('Unauthorized', { status: 401 }) } const userId = session.user.id let cleanup: (closeStream?: boolean) => void = () => undefined const stream = new ReadableStream<Uint8Array>({ start(controller) { let closed = false let pollTimer: ReturnType<typeof setTimeout> | undefined let heartbeat: ReturnType<typeof setInterval> | undefined let onAbort: (() => void) | undefined let cursor = request.headers.get('last-event-id') cleanup = (closeStream = false) => { if (closed) return closed = true if (pollTimer) clearTimeout(pollTimer) if (heartbeat) clearInterval(heartbeat) if (onAbort) request.signal.removeEventListener('abort', onAbort) if (closeStream) { try { controller.close() } catch { // Strumień mógł zostać zamknięty przez runtime. } } } controller.enqueue(encoder.encode('retry: 5000\n\n')) heartbeat = setInterval(() => { if (!closed) { controller.enqueue(encoder.encode(': heartbeat\n\n')) } }, 30000) async function poll() { try { const notifications: Notification[] = cursor ? await getNotificationsAfter(userId, cursor) : await getUnreadNotifications(userId, { limit: 50 }) for (const notification of notifications) { if (closed) return controller.enqueue(encodeNotification(notification)) cursor = notification.id } } catch (error) { if (!closed) console.error('SSE polling error', error) } finally { if (!closed) pollTimer = setTimeout(poll, 2000) } } void poll() onAbort = () => cleanup(true) request.signal.addEventListener('abort', onAbort, { once: true }) if (request.signal.aborted) onAbort() }, cancel() { cleanup(false) }, }) return new Response(stream, { headers: { 'Content-Type': 'text/event-stream; charset=utf-8', 'Cache-Control': 'no-cache, no-transform', 'X-Accel-Buffering': 'no', }, })}
Funkcja getUnreadNotifications powinna wybrać 50 najnowszych rekordów i zwrócić je od najstarszego do najnowszego. getNotificationsAfter zachowuje ten sam porządek i mapuje daty na string ISO. Po wysłaniu każdego rekordu kursor przesuwa się na jego trwałe id. Nie używaj czasu jako jedynego kursora, ponieważ dwa zdarzenia mogą mieć ten sam znacznik czasu.
Po długiej przerwie liczba pominiętych zdarzeń może być duża. Wtedy pobieraj je partiami zamiast ładować całą historię jednym zapytaniem. Przy strumieniach o wysokiej częstotliwości kontroluj również controller.desiredSize i ustal limit kolejki. Wolny klient nie powinien powodować nieograniczonego wzrostu pamięci procesu.
Na Vercel maxDuration określa limit tej funkcji w sekundach. Wartość 300 mieści się w aktualnym limicie planu Hobby. Projekty Pro i Enterprise z Fluid compute mogą ustawić wyższą wartość zgodnie z limitem konta. W Next.js 15 i nowszych Route Handlers GET nie są domyślnie cache'owane, więc dynamic = 'force-dynamic' nie jest tu potrzebne.
Reconnect bez gubienia powiadomień
EventSource po zerwaniu połączenia ponawia je samodzielnie, a sugerowany odstęp kontrolujesz polem retry w strumieniu. Sam reconnect nie odzyska jednak zdarzeń wysłanych podczas przerwy. Pole id ustawia identyfikator ostatniego zdarzenia po stronie przeglądarki. Podczas automatycznego wznowienia przeglądarka przesyła go w nagłówku Last-Event-ID, a endpoint doładowuje z bazy wszystkie późniejsze rekordy.
Nowy obiekt EventSource utworzony po pełnym przeładowaniu strony nie dziedziczy kursora poprzedniej instancji. Dlatego pierwsze połączenie pobiera ograniczony zestaw nieprzeczytanych powiadomień, a reconnect korzysta z Last-Event-ID. Strumień jest transportem, natomiast baza pozostaje źródłem prawdy.
Jak przetestować strumień SSE
Opcja -N wyłącza buforowanie po stronie curl, dzięki czemu każde zdarzenie pojawia się od razu. Dla endpointu chronionego sesją przekaż ciasteczko skopiowane z lokalnego środowiska:
Powinieneś od razu zobaczyć retry, a później komentarze heartbeat i bloki zakończone pustą linią. Replay sprawdzisz, zamykając połączenie, tworząc powiadomienie i uruchamiając żądanie z kursorem:
W produkcji obserwuj liczbę aktywnych strumieni, częstotliwość reconnectów, błędy subskrypcji, czas replay i zamknięcia spowodowane przez maxDuration. Sam test lokalny nie wykryje buforowania w CDN lub reverse proxy, dlatego powtórz go na adresie wdrożenia.
W trybie Strict Mode React wykonuje w development dodatkowy cykl uruchomienia i cleanupu efektu. Krótkie otwarcie dwóch połączeń w logach lokalnych jest więc możliwe, ale pierwsze z nich musi zostać natychmiast zamknięte przez eventSource.close(). W produkcji ten dodatkowy cykl nie występuje.
Zamiast odpytywać bazę co dwie sekundy możesz użyć . Po utworzeniu powiadomienia zapisujesz rekord w bazie, a następnie publikujesz go na kanale użytkownika:
Code
// lib/notifications.tsimport { Redis } from '@upstash/redis'import { db } from '@/lib/db'import type { Notification } from '@/lib/notification-types'const redis = new Redis({ url: process.env.UPSTASH_REDIS_REST_URL!, token: process.env.UPSTASH_REDIS_REST_TOKEN!,})type NewNotification = Pick<Notification, 'message' | 'type'>export async function publishNotification( userId: string, input: NewNotification,) { const record = await db.notification.create({ data: { ...input, userId }, }) const notification: Notification = { id: record.id, message: record.message, type: record.type, createdAt: record.createdAt.toISOString(), } await redis.publish(`notifications:${userId}`, notification) return notification}
Zapis następuje przed publikacją, ponieważ baza jest źródłem prawdy. W sytuacji kiedy publikacja się nie powiedzie, rekord nadal może zostać odtworzony po reconnect, a gdy system wymaga gwarantowanej publikacji bez czekania na reconnect, zastosuj transactional outbox i osobny proces publikujący.
Po stronie endpointu SSE subskrybujesz ten sam kanał. subscribe() w @upstash/redis używa Server-Sent Events nad REST API Upstash, więc nie wymaga klasycznego połączenia Redis przez TCP:
Code
// app/api/notifications/stream/route.tsimport { Redis } from '@upstash/redis'import { auth } from '@/auth'import { getNotificationsAfter, getUnreadNotifications,} from '@/lib/notifications'import type { Notification } from '@/lib/notification-types'export const maxDuration = 300const redis = Redis.fromEnv()const encoder = new TextEncoder()function encodeNotification(notification: Notification) { return encoder.encode( [ `id: ${notification.id}`, 'event: notification', `data: ${JSON.stringify(notification)}`, '', '', ].join('\n'), )}export async function GET(request: Request) { const session = await auth() if (!session?.user) { return new Response('Unauthorized', { status: 401 }) } const userId = session.user.id const channel = `notifications:${userId}` let cleanup: (closeStream?: boolean) => void = () => undefined const stream = new ReadableStream<Uint8Array>({ async start(controller) { let closed = false let replaying = true let heartbeat: ReturnType<typeof setInterval> | undefined let onAbort: (() => void) | undefined let resolveAbort: (() => void) | undefined const buffered = new Map<string, Notification>() const sentDuringReplay = new Set<string>() const subscription = redis.subscribe<Notification>(channel) const aborted = new Promise<void>((resolve) => { resolveAbort = resolve }) cleanup = (closeStream = false) => { if (closed) return closed = true if (heartbeat) clearInterval(heartbeat) if (onAbort) request.signal.removeEventListener('abort', onAbort) resolveAbort?.() void subscription .unsubscribe() .catch((error) => { console.error('Upstash unsubscribe error', error) }) .finally(() => { subscription.removeAllListeners() }) if (closeStream) { try { controller.close() } catch { // Strumień mógł zostać zamknięty przez runtime. } } } function send(notification: Notification) { if (closed || sentDuringReplay.has(notification.id)) return controller.enqueue(encodeNotification(notification)) if (replaying) sentDuringReplay.add(notification.id) } subscription.on('message', ({ message }) => { if (replaying) buffered.set(message.id, message) else send(message) }) subscription.on('error', (error) => { if (!closed) { console.error('Upstash subscribe error', error) cleanup(true) } }) const subscribed = new Promise<void>((resolve, reject) => { subscription.on('subscribe', () => resolve()) subscription.on('error', reject) }) controller.enqueue(encoder.encode('retry: 5000\n\n')) heartbeat = setInterval(() => { if (!closed) { controller.enqueue(encoder.encode(': heartbeat\n\n')) } }, 30000) onAbort = () => cleanup(true) request.signal.addEventListener('abort', onAbort, { once: true }) if (request.signal.aborted) onAbort() try { const state = await Promise.race([ subscribed.then(() => 'subscribed' as const), aborted.then(() => 'aborted' as const), ]) if (state === 'aborted' || closed) return const lastEventId = request.headers.get('last-event-id') const missed: Notification[] = lastEventId ? await getNotificationsAfter(userId, lastEventId) : await getUnreadNotifications(userId, { limit: 50 }) for (const notification of missed) send(notification) for (const notification of buffered.values()) send(notification) buffered.clear() sentDuringReplay.clear() replaying = false } catch (error) { if (!closed) { console.error('SSE initialization error', error) cleanup(true) } } }, cancel() { cleanup(false) }, }) return new Response(stream, { headers: { 'Content-Type': 'text/event-stream; charset=utf-8', 'Cache-Control': 'no-cache, no-transform', 'X-Accel-Buffering': 'no', }, })}
Kolejność inicjalizacji jest celowa. Endpoint najpierw czeka na potwierdzenie subskrypcji, następnie odtwarza historię z bazy, a zdarzenia odebrane w tym czasie przechowuje w buforze. Zbiór sentDuringReplay usuwa duplikaty występujące w obu źródłach. Po zakończeniu replay kolejne wiadomości trafiają bezpośrednio do klienta. Błąd subskrypcji zamyka zewnętrzny strumień, dzięki czemu EventSource tworzy nowe połączenie wraz z nową subskrypcją Upstash.
Redis Pub/Sub nie przechowuje historii i nie gwarantuje dostarczenia wiadomości klientowi, który był rozłączony. Jeśli nie chcesz opierać replay na bazie aplikacji, rozważ trwały log zdarzeń, taki jak Redis Streams. Niezależnie od wybranego transportu cleanup musi wywołać unsubscribe(), aby rozłączenie przeglądarki nie pozostawiało aktywnej subskrypcji Upstash.
WebSockets: kiedy SSE nie wystarczy?
SSE jest jednokierunkowe i przesyła dane od serwera do klienta. Gdy obie strony muszą wysyłać dane przez jedno trwałe połączenie, na przykład w czacie, współedycji lub grze, wybierz WebSocket.
Vercel Functions nie mogą działać jako serwer WebSocket. Możesz użyć usługi zarządzanej albo uruchomić osobny, długowieczny serwer poza Route Handlerem.
Polling, SSE czy WebSockets: co wybrać?
Polling sprawdza się przy statusie zamówienia, dostępności produktu i prostym dashboardzie z rzadkimi aktualizacjami.
SSE pasuje do powiadomień, feedu aktywności, logów i metryk przesyłanych od serwera do klienta.
WebSocket wybierz do czatu, współedycji dokumentów, gier wieloosobowych i innych przepływów dwukierunkowych.
Elastyczne i wydajne narzędzia dla biznesu, które dotrzymają kroku Twojemu rozwojowi.
Tak, ale pojedyncze połączenie nadal podlega limitowi czasu funkcji. Według stanu na 23 lipca 2026 roku limit Hobby wynosi 300 sekund, a funkcje Node.js na planach Pro i Enterprise mogą korzystać z maxDuration do 1800 sekund w becie przy włączonym Fluid compute. Po zakończeniu funkcji EventSource połączy się ponownie, dlatego endpoint musi obsługiwać Last-Event-ID i odtwarzanie pominiętych zdarzeń. Heartbeat nie omija twardego limitu czasu.
Ile jednoczesnych połączeń SSE obsłuży serwer?
Nie ma jednej bezpiecznej liczby, ponieważ pojemność zależy od runtime'u, limitów hostingu, pamięci, deskryptorów plików, proxy oraz liczby połączeń przypadających na użytkownika. Na serverless trzeba uwzględnić także czas funkcji i koszt utrzymywania strumieni. Przed wdrożeniem wykonaj test obciążeniowy z docelowym czasem połączenia, a przy dużej skali porównaj koszt własnego rozwiązania z usługą zarządzaną.
Czy SSE działa na urządzeniach mobilnych?
Tak. EventSource jest obsługiwany przez przeglądarki mobilne. Trzeba jednak pamiętać, że systemy mobilne potrafią zamykać połączenia sieciowe, gdy aplikacja albo karta przechodzi w tło, by oszczędzać baterię. To nie jest problem SSE jako takiego. Po powrocie na pierwszy plan EventSource automatycznie ponawia połączenie, więc strumień się wznawia bez dodatkowego kodu. Pominięte dane trzeba jednak odtworzyć po identyfikatorze ostatniego zdarzenia.
Czym SSE różni się od WebSockets?
SSE to jednokierunkowy strumień serwer→klient po zwykłym HTTP, z natywną obsługą w przeglądarce i automatycznym reconnectem. Sprawdza się przy powiadomieniach, feedach i logach, gdzie dane płyną tylko w jedną stronę. WebSockets dają komunikację dwukierunkową przez osobny protokół, co jest niezbędne przy czacie, współedycji czy grach, ale kosztuje więcej złożoności i nie jest natywnie wspierane w Route Handlerach Next.js. Jeśli klient tylko odbiera zdarzenia, zwykle zacznij od SSE.
Po co stosować Redis Pub/Sub zamiast odpytywać bazę w endpoincie SSE?
Bo odpytywanie bazy co kilka sekund w każdym otwartym strumieniu skaluje się słabo. Przy wielu klientach mnożysz zapytania niezależnie od tego, czy są nowe dane. Z Redis Pub/Sub publikujesz event dokładnie wtedy, gdy powiadomienie powstaje, a endpoint SSE subskrybuje kanał i przekazuje je dalej. Upstash realizuje subskrypcję przez SSE nad swoim REST API, więc działa to także w środowiskach serverless bez trwałych połączeń TCP. Pub/Sub jest jednak ulotny, dlatego baza danych nadal pozostaje źródłem prawdy potrzebnym do odtwarzania zdarzeń.
Czy EventSource może wysłać nagłówek Authorization?
Natywny konstruktor EventSource nie pozwala dodawać własnych nagłówków. Najprostszy wariant używa endpointu tego samego originu i ciasteczka sesji. Przy innym originie potrzebujesz withCredentials oraz poprawnej konfiguracji CORS. Nie wkładaj długowiecznego tokenu do query string, ponieważ adres może trafić do logów i historii.
O autorze
Maciej Sala
Maciej Sala — Product Manager i Frontend Developer z bogatym doświadczeniem w marketingu internetowym oraz SEO. Na co dzień pracuje z Reactem, Next.js i TypeScriptem, a ostatnio także z Astro i narzędziami do automatyzacji procesów AI. Sprawnie łączy perspektywę produktową z praktycznym podejściem do kodu. Przez kilka lat był związany z branżą gier wideo jako project manager i game designer. Absolwent historii na Uniwersytecie Jagiellońskim oraz studiów podyplomowych z marketingu internetowego na AGH w Krakowie. Po godzinach trenuje na siłowni, maluje figurki i rozwijam własne projekty.
Pierwsza część serii uporządkowała fundamenty: serwer, bazę danych, API i CORS. Teraz przechodzimy do obszarów, które zwykle pojawiają się chwilę później, gdy aplikacja przestaje być prostym CRUD-em: komunikacja w czasie rzeczywistym, webhooki, integracje zewnętrzne oraz autentykacja .
Maciej Sala
Founder StriveLab
To trzecia część mojej małej serii „Backend dla frontendowca” i po fundamentach API oraz tematach real-time, webhooków i uwierzytelniania zostaje warstwa, która decyduje o tym, czy aplikacja wytrzyma prawdziwe użytkowanie: cache , kolejki, pliki, deployment, monitoring i bezpieczeństwo .
Maciej Sala
Founder StriveLab
App Router daje Ci dwa sposoby na uruchomienie kodu po stronie serwera, poprzez Route Handlers i Server Actions . Może wyglądają podobnie, ponieważ oba działają na serwerze, sięgają do bazy i do zmiennych środowiskowych, ale każdy z nich rozwiązuje zupełnie inny problem. Użycie jednego tam, gdzie pasuje drugi, nie jest odpowiednim rozwiązaniem, ponieważ zły wybór odbija się potem na architekturze.