Tâches longues en Node.js : notifier les utilisateurs de l'avancée des tâches via SSE
- Les Server-Sent Events (SSE) sont un mécanisme simple et léger, basé sur HTTP, qui envoie des mises à jour en temps réel du serveur vers le navigateur.
- L'EventEmitter de NestJS fait le pont entre les workers de la queue et le flux SSE, selon une approche Event-Driven.
- Cette architecture s'applique à de nombreux cas concrets : analyse IA, OCR, génération de documents, import de données volumineuses.
Il arrive que l'on traite des tâches lourdes côté API, ce qui complique la gestion des requêtes. Un OCR, une analyse IA ou toute autre tâche de plusieurs secondes dégrade l'expérience utilisateur dès qu'elle bloque la requête.
Dans un précédent article, nous avions vu comment un système de queue règle les problèmes liés au couplage fort de ce genre de tâches : timeouts des requêtes, scalabilité, etc.
Il est donc temps de s'attaquer à notre nouveau problème : Quel est le sens de la vie ? Comment informer les différents clients qui consomment l'API de l'avancée des tâches traitées en arrière-plan ?
Si vous n'avez pas lu notre article sur le traitement de tâches lourdes en arrière-plan, nous vous conseillons de remédier à cette hérésie et d'aller le consulter juste ici.
Le système de queues c'est bien, mais...
Le principe de Jobs et de Workers délègue le traitement à un processus en arrière-plan et répond immédiatement au client. Une fois la requête HTTP classique honorée, le client perd toute visibilité sur ce qui se passe côté serveur.
Rassurez-vous, nous avons la solution, et elle se passe de polling.
Notifier les clients : l'approche Event-Driven
L'architecture Event-Driven repose sur un principe simple. Quand quelque chose se passe, on émet un événement, et les composants intéressés s'y abonnent pour réagir.
Dans notre cas, l'idée est d'utiliser ce principe à deux niveaux :
- Le worker émet un événement vers le serveur quand le job progresse.
- Le serveur pousse l'information vers les navigateurs connectés.
Nous empruntons quelques concepts à cette architecture pour envoyer des événements côté client, quitte à la suivre seulement en partie.
L'EventEmitter de NestJS
NestJS fournit un package dédié (@nestjs/event-emitter) qui implémente un système d'événements interne à l'application. Un worker peut émettre un événement, et un service peut s'y abonner, sans que les deux se connaissent directement.
flowchart LR
A@{ shape: processes, label: "Jobs" } --> B@{ "shape": "database", label: "Queue" }
B --> C@{ shape: lin-rect, label: "Worker 1" }
B --> D@{ shape: lin-rect, label: "Worker 2" }
B --> E@{ shape: lin-rect, label: "Worker 3" }
C -.-> | événement | F{EventEmitter}
D -.-> | événement | F{EventEmitter}
E -.-> | événement | F{EventEmitter}
F --> G@{ shape: paper-tape, label: "Listener"}
F --> H@{ shape: paper-tape, label: "Listener"}
Ce découplage est essentiel. Le worker se contente d'émettre un événement avec les données de progression, peu importe qui l'écoute.
Server-Sent Events (SSE) : pousser des mises à jour vers le navigateur
Qu'est-ce que le SSE ?
On entend souvent parler de polling ou de WebSockets. D'autres manières existent pour faire remonter des informations de l'API vers le client.
Les Server-Sent Events (ou SSE) sont un mécanisme standard du web (API EventSource) qui permet au serveur d'envoyer des événements au client via une connexion HTTP persistante.
sequenceDiagram
Client->>Serveur: Ouvre un flux d'événements (GET /events)
Serveur-->>Client: Événement 1
Serveur-->>Client: Événement 2
Serveur-->>Client: Événement 3
Client->>Serveur: Ferme la connexion
Le client ouvre une connexion sur un endpoint dédié. Le serveur garde cette connexion ouverte et envoie des événements au format text/event-stream chaque fois qu'il a quelque chose à communiquer.
SSE ou WebSockets : comment choisir
Il est de notre devoir de vous faire un comparatif :
| WebSockets | Server-Sent Events | |
|---|---|---|
| Direction du flux | Bidirectionnel (serveur ↔ client) | Unidirectionnel (serveur → client) |
| Protocole | WebSocket (upgrade depuis HTTP) | HTTP/HTTPS standard |
| Complexité | Plus complexe (gestion d'un protocole dédié, reconnexion manuelle) | Simple (HTTP natif, reconnexion automatique par le navigateur) |
| Performance | Surcharge à l'ouverture, mais très efficace pour du trafic bidirectionnel à haute fréquence | Surcharge minimale à l'initialisation, idéal pour du push occasionnel |
| Cas d'usage | Chat en temps réel, jeux multijoueurs, collaboration simultanée | Notifications, suivi de progression, flux d'actualités |
Si l'objectif est simplement de prévenir les clients d'une mise à jour côté serveur de manière ponctuelle, le SSE est la solution la plus adaptée. Il est plus simple à mettre en place et repose sur HTTP standard. Le navigateur gère automatiquement la reconnexion en cas de coupure.
Si vos besoins sont plus complexes, notamment avec de la communication bidirectionnelle entre le client et le serveur, le protocole WebSocket sera sûrement un meilleur choix.
Gardez en tête que d'autres technologies font communiquer votre serveur et vos clients dans des cas d'usage différents, certaines plus complexes comme le WebRTC. Voyez ce comparatif comme un point de départ et amusez-vous à fouiller par vous-même pour trouver la perle rare qu'il vous faut !
Limitations à connaître
Sur HTTP/1.1, le SSE plafonne à 6 connexions simultanées par domaine et par navigateur. Cette limite disparaît avec HTTP/2, aujourd'hui largement supporté. Pensez-y si vous devez ouvrir plusieurs flux SSE différents dans la même application.
L'architecture complète : de la queue au navigateur
En combinant queue, EventEmitter et SSE, on obtient le pipeline suivant :
flowchart LR
subgraph Queue system
A@{ shape: processes, label: "Jobs" } --> B@{ "shape": "database", label: "Queue" }
B --> C@{ shape: lin-rect, label: "Worker 1" }
B --> D@{ shape: lin-rect, label: "Worker 2" }
B --> E@{ shape: lin-rect, label: "Worker 3" }
end
C -.-> | événement | F{EventEmitter}
D -.-> | événement | F{EventEmitter}
E -.-> | événement | F{EventEmitter}
subgraph Events-Driven
F --> G@{ shape: paper-tape, label: "Listener"}
end
subgraph "SSE"
G --> H["Observable (RxJS Subject)"]
H --> | événement | I["Client 1"]
H --> | événement | J["Client 2"]
H --> | événement | K["Client 3"]
end
Le flux se déroule ainsi :
- Le client envoie une requête
POST /analyzepour lancer une analyse. - Le serveur ajoute un job dans la queue BullMQ et répond immédiatement.
- Un worker récupère le job et commence le traitement.
- À chaque étape, le worker émet un événement via l'EventEmitter.
- Un listener capte cet événement et le pousse dans un RxJS
Subject. - Le
Subjectenvoie l'événement à tous les clients connectés au flux SSE.
Place à l'exemple : l'analyse du contenu d'un fichier
Vous vous souvenez de l'exemple de l'article sur l'utilisation de BullMQ ? Eh bien il est temps de le compléter.
Pour rappel, nous allons prendre un cas que nous avons rencontré chez Lonestone : l'analyse du contenu d'un fichier pour en retirer des informations importantes.
L'architecture du système
Dans le précédent exemple, nous avions déjà vu comment envoyer une requête POST vers l'API pour lancer une analyse. Cette requête ajoute un job dans la queue et retourne une réponse au client.
Avec l'EventEmitter et le SSE, nous ajoutons maintenant la communication en temps réel. Le client va donc se connecter à un flux SSE pour écouter les événements côté API.
Voici le diagramme de séquence mis à jour avec notre système de SSE :
sequenceDiagram
participant Client
participant API
participant Queue
participant Worker
Client->>API: Ouvre /events (SSE)
API-->>Client: Flux SSE ouvert
Client->>API: POST /analyze
API->>Queue: Ajoute un job
API-->>Client: 200 OK
Queue->>Worker: Traite le job
Worker->>Worker: Étape 1 (extraction)
Worker-->>API: Émet événement "extraction"
API-->>Client: SSE → extraction
Worker->>Worker: Étape 2 (analyse)
Worker-->>API: Émet événement "analyse"
API-->>Client: SSE → analyse
Worker->>Worker: Terminé
Worker-->>API: Émet événement "completed"
API-->>Client: SSE → completed
Un peu de code : le SSE avec NestJS
Deux grands concepts vont nous être utiles : l'EventEmitter pour les événements internes à l'API, et RxJS pour exposer le SSE.
Configuration du module
Notre système de queue étant en place, on se rend dans AnalysisModule pour y déclarer notre AnalysisEventsService :
import { BullModule } from '@nestjs/bullmq'import { Module } from '@nestjs/common'import { AnalysisEventsService } from './analysis-events.service'import { AnalysisController } from './analysis.controller'import { ANALYSIS_QUEUE_NAME, AnalysisProcessor } from './analysis.processor'import { AnalysisService } from './analysis.service'
@Module({ imports: [ BullModule.registerQueue({ name: ANALYSIS_QUEUE_NAME }), ], controllers: [AnalysisController], providers: [AnalysisService, AnalysisEventsService, AnalysisProcessor],})export class AnalysisModule {}Le service d'événements : le pont entre EventEmitter et SSE
Ce service écoute les événements internes et les transforme en un flux Observable (via un RxJS Subject) que NestJS peut exposer en SSE :
import { Injectable, OnModuleDestroy } from '@nestjs/common'import { EventEmitter2 } from '@nestjs/event-emitter'import { Subject } from 'rxjs'
export const ANALYSIS_UPDATED_EVENT = 'analysis-updated-event'
@Injectable()export class AnalysisEventsService implements OnModuleDestroy { private readonly eventSubject = new Subject<MessageEvent>()
constructor(private readonly eventEmitter: EventEmitter2) { this.eventEmitter.on(ANALYSIS_UPDATED_EVENT, (event: AnalysisEvent) => { this.eventSubject.next({ data: event } as MessageEvent) }) }
onUpdated() { return this.eventSubject.asObservable() }
onModuleDestroy() { this.eventSubject.complete() }}Le Subject RxJS fait le pont :
- En entrée, il reçoit les événements via
next()à chaque fois que l'EventEmitter déclencheanalysis-updated-event. - En sortie, il expose un
Observableque NestJS utilise pour alimenter le flux SSE. - À la destruction du module, le
Subjectest complété proprement pour fermer les connexions.
Le processor (worker) : émettre les événements
Une fois l'AnalysisEventsService en place, nous pouvons mettre à jour notre processor, écrit dans l'article précédent.
Pour rappel, c'est ici que nous traitons les tâches lourdes.
Nous devons donc émettre un événement à chaque étape via l'EventEmitter2 de NestJS, comme ceci :
this.eventEmitter.emit(ANALYSIS_UPDATED_EVENT, { id: job.data.analysisId, step: ANALYSIS_STEPS.EXTRACTION,})Dans le code, ça donne :
import { Processor, WorkerHost } from '@nestjs/bullmq'import { EventEmitter2 } from '@nestjs/event-emitter'import { Job } from 'bullmq'
export const ANALYSIS_QUEUE_NAME = 'analysis_queue'export const ANALYSIS_JOB_NAME = 'analysis_job'export const ANALYSIS_JOBS_CONCURRENCY = 10
@Processor(ANALYSIS_QUEUE_NAME, { concurrency: ANALYSIS_JOBS_CONCURRENCY, removeOnComplete: { age: 3600, count: 1000 }, removeOnFail: { age: 24 * 3600 },})export class AnalysisProcessor extends WorkerHost { constructor(private readonly eventEmitter: EventEmitter2) { super() }
async process(job: Job<AnalysisJobData>) { this.eventEmitter.emit(ANALYSIS_UPDATED_EVENT, { id: job.data.analysisId, step: ANALYSIS_STEPS.EXTRACTION, })
await performExtraction(job.data) // Tâche longue
this.eventEmitter.emit(ANALYSIS_UPDATED_EVENT, { id: job.data.analysisId, step: ANALYSIS_STEPS.ANALYSIS_PART_ONE, })
await performAnalysis(job.data) // Tâche longue
this.eventEmitter.emit(ANALYSIS_UPDATED_EVENT, { id: job.data.analysisId, step: ANALYSIS_STEPS.COMPLETED, }) }}Chaque emit() envoie un événement interne que l'AnalysisEventsService capte.
Le contrôleur : exposer le flux SSE
Il nous reste à ouvrir la connexion côté API (je sais, vous ne vous y attendiez pas du tout). Nous déclarons l'endpoint dans l'AnalysisController :
import { Controller, HttpCode, HttpStatus, Param, Post, Sse } from '@nestjs/common'
@Controller('analysis')export class AnalysisController { constructor( private readonly analysisEvents: AnalysisEventsService, ) {}
// [...]
@Sse('events') getEvents() { return this.analysisEvents.onUpdated() }}Le décorateur @Sse('events') de NestJS fait tout le travail : il ouvre une connexion HTTP persistante avec le header Content-Type: text/event-stream et envoie chaque valeur émise par l'Observable au format SSE standard.
Implémentation côté front avec React
Côté client, l'objectif est de consommer le flux SSE et de mettre à jour l'interface en temps réel. L'exemple utilise React avec TanStack Query (React Query).
Écouter le flux SSE
On déclare le hook useAnalysisEvents, qui s'abonne au flux SSE via la fonction expérimentale streamedQuery de TanStack Query. À chaque événement reçu, il met à jour le cache des analyses :
import { queryOptions, experimental_streamedQuery as streamedQuery, useQuery, useQueryClient,} from '@tanstack/react-query'import { useEffect, useRef } from 'react'
export function useAnalysisEvents() { const queryClient = useQueryClient() const lastProcessedLengthRef = useRef(0)
const query = queryOptions({ queryKey: ['analysis-events'], queryFn: streamedQuery({ streamFn: async () => { const { stream } = await analysisControllerGetEvents() return stream }, }), })
const { data: streamedData } = useQuery(query)
useEffect(() => { const events = Array.isArray(streamedData) ? streamedData : [] const from = lastProcessedLengthRef.current if (from >= events.length) return
const toProcess = events.slice(from)
for (const raw of toProcess) { const eventData = zAnalysisEventsSchema.parse(raw)
queryClient.setQueryData(['analyses'], (previous) => { return previous.map((analysis) => analysis.id === eventData.id ? { ...analysis, ...eventData } : analysis, ) }) }
lastProcessedLengthRef.current = events.length }, [streamedData, queryClient])}Quelques détails sur le fonctionnement :
streamedQuerygère nativement les flux (ReadableStream) et accumule les événements dans un tableau. Si vous préférez passer par une méthode plus conventionnelle, vous pouvez utiliser l'interfaceEventSourcedirectement (en ajoutant uneventListenerdessus).lastProcessedLengthRefretient le nombre d'événements déjà traités, si bien que chaque rendu reprend uniquement les nouveaux.setQueryDatamet à jour le cache local de façon optimiste. L'UI se rafraîchit instantanément sans avoir à relancer une requête. Nous procédons ainsi car dans notre exemple, la donnée vit uniquement dans le cache du front. Libre à vous de gérer la synchronisation autrement (invalidation du cache, refetch après chaque événement, etc.).
Afficher la progression
Le dashboard affiche déjà une carte par analyse avec son statut, mais ce statut reste figé. Il est temps d'y remédier en appelant le hook useAnalysisEvents.
export default function DashboardPage() { useAnalysisEvents() const { data: analyses } = useAnalyses() const { mutate } = useStartAnalysis()
return ( <main className="container mx-auto py-8 px-4 space-y-6"> <h1 className="text-3xl font-bold">Analyses</h1> {analyses?.map((analysis) => ( <Card key={analysis.id}> <CardHeader> <CardTitle>Analysis</CardTitle> <AnalysisBadge step={analysis.step} /> </CardHeader> <CardFooter> <Button onClick={() => mutate({ id: analysis.id })} disabled={ analysis.step !== 'completed' && analysis.step !== 'failed' } > Lancer l'analyse </Button> </CardFooter> </Card> ))} </main> )}Le composant AnalysisBadge affiche un badge coloré en fonction de l'étape en cours (extraction, analyse, terminé, échoué). Le cache TanStack Query pilote tout l'affichage : quand le hook SSE y écrit les nouvelles données, React refait automatiquement le rendu des composants concernés.
Et voilà, normalement vous avez un exemple fonctionnel de bout en bout, queue et SSE compris !
Pour conclure
Ces deux articles font le tour du traitement des tâches lourdes dans une application Node.js, avec NestJS côté API et React côté client.
Ensemble, BullMQ, l'EventEmitter et le SSE donnent une architecture simple à maintenir :
- BullMQ gère le traitement en arrière-plan via Redis, avec concurrence, retries et nettoyage automatique.
- L'EventEmitter fait circuler les événements de progression à l'intérieur de l'application.
- Le SSE pousse ces événements vers les navigateurs connectés, en s'appuyant sur HTTP standard.
Vous pouvez ajuster cette architecture à vos besoins, en ajoutant des étapes au worker ou en persistant les données côté serveur !
Si vous avez ce genre de traitement lourd à sortir de vos requêtes HTTP, on peut le faire avec vous chez Lonestone. Quand l'application tourne déjà et que les temps de réponse posent problème, un Audit identifie ce qui mérite de passer en asynchrone et dans quel ordre. Pour la mise en œuvre, Build prend le produit en charge de bout en bout, et Scale renforce une équipe interne avec des développeurs rodés à ces architectures.
Le code complet de l'exemple, qui combine BullMQ et le SSE, est sur GitHub.