Angular

RxJS

Observables, opérateurs, subjects, interopérabilité avec les signals.

Vérifié en septembre 2026 · Angular 22.1 · environ 16 min

RxJS est la bibliothèque de flux sur laquelle Angular s'est construit : HttpClient, le routeur et les formulaires réactifs exposent encore leurs résultats et leurs événements en observables. Depuis que les signals portent l'état, le partage des rôles est net. Un signal est une valeur présente, toujours lisible. Un observable est une suite de valeurs dans le temps, qui peut se terminer, échouer ou être abandonnée en route. RxJS garde donc ce que les signals ne savent pas faire : attendre, espacer, annuler, combiner des événements. Le cours suppose connu le cours Signals et finit sur le passage de l'un à l'autre. RxJS appartient à la même famille ReactiveX que Rx.NET : un observable y tient le rôle d'IObservable<T>, un flux poussé, là où IEnumerable<T> est tiré. Les exemples utilisent RxJS 7.8.2, la version que ce projet installe et la dernière stable au 26 septembre 2026. Une 9.0 est en bêta depuis août 2026, hors de la plage que déclare @angular/core 22.1 : ^6.5.3 || ^7.4.0.

L'Observable, un flux paresseux

Un observable est une fonction qui attend un abonné. Tant que personne n'appelle subscribe, rien ne s'exécute : ni requête, ni minuteur. Chaque abonnement exécute le producteur une nouvelle fois, pour lui seul. L'abonné reçoit trois sortes de notifications : next, autant de fois qu'il y a de valeurs, puis au plus une fin, complete ou error, après laquelle plus rien ne passe.

import { Observable } from 'rxjs';

// Creer l'observable n'execute rien : la fonction passee au constructeur est
// le producteur, et elle sera rappelee pour chaque abonnement.
const tirages = new Observable<number>((abonne) => {
  console.log('producteur demarre');
  abonne.next(1);
  abonne.next(2);

  const minuteur = setTimeout(() => {
    abonne.next(3);
    abonne.complete(); // apres complete, plus rien ne passe
  }, 100);

  // La fonction rendue est le nettoyage : appelee apres complete ou error,
  // ou quand l'abonne se desabonne.
  return () => {
    clearTimeout(minuteur);
    console.log('nettoyage');
  };
});

console.log('avant subscribe');
tirages.subscribe({
  next: (valeur) => console.log('A recoit', valeur),
  error: (erreur) => console.log('A echoue', erreur),
  complete: () => console.log('A termine'),
});
console.log('apres subscribe');

const b = tirages.subscribe((valeur) => console.log('B recoit', valeur));
b.unsubscribe();

// avant subscribe
// producteur demarre     <- l'abonnement lance le producteur...
// A recoit 1             <- ...et les deux premieres valeurs arrivent
// A recoit 2                pendant l'appel a subscribe, de facon synchrone
// apres subscribe
// producteur demarre     <- B relance le producteur : une seconde execution
// B recoit 1
// B recoit 2
// nettoyage              <- B se desabonne, son minuteur est annule
// A recoit 3             <- 100 ms plus tard
// A termine
// nettoyage

Deux propriétés surprennent quand on vient des promesses. Un observable n'est pas asynchrone par nature : les deux premières valeurs arrivent pendant l'appel à subscribe, avant la ligne suivante. Et il s'annule : unsubscribe appelle la fonction de nettoyage du producteur, qui défait ce qu'il avait lancé. Une promesse commencée va à son terme ; un observable dont on se désabonne ne produit plus rien. switchMap, takeUntilDestroyed et l'annulation des requêtes de HttpClient reposent sur cette propriété.

Les opérateurs

Un opérateur est une fonction qui prend un observable et en rend un autre. pipe les enchaîne dans l'ordre d'écriture ; la source n'est jamais modifiée, et le résultat est aussi paresseux qu'elle. Les plus courants ont un équivalent LINQ direct ; la différence tient au moment où ils s'appliquent : non quand un foreach réclame la valeur suivante, mais quand la source la pousse.

import { filter, from, map, scan, tap } from 'rxjs';

interface Mouvement {
  readonly libelle: string;
  readonly montant: number;
}

const releve: readonly Mouvement[] = [
  { libelle: 'salaire', montant: 2400 },
  { libelle: 'loyer', montant: -900 },
  { libelle: 'annulation', montant: 0 },
  { libelle: 'courses', montant: -130 },
];

// pipe ne modifie pas la source : il rend un nouvel observable, qui ne fait
// rien tant que personne ne s'y abonne. Chaque operateur s'applique a chaque
// valeur, l'une apres l'autre, comme un Where et un Select de LINQ, mais sur
// des valeurs poussees par la source au lieu d'etre tirees par un foreach.
const soldes = from(releve).pipe(
  filter((mouvement) => mouvement.montant !== 0), // Where
  tap((mouvement) => console.log('passe :', mouvement.libelle)), // effet de bord, valeur intacte
  map((mouvement) => mouvement.montant), // Select
  scan((solde, montant) => solde + montant, 0), // un Aggregate qui emet chaque etape
);

soldes.subscribe((solde) => console.log('solde :', solde));

// passe : salaire
// solde : 2400
// passe : loyer
// solde : 1500
// passe : courses
// solde : 1370

tap observe sans modifier : un journal, un compteur. scan émet chaque état intermédiaire de l'accumulation, là où reduce attend la fin du flux pour émettre une seule valeur. Les opérateurs de temps n'ont pas d'équivalent LINQ : debounceTime n'émet qu'après un silence, distinctUntilChanged écarte une valeur égale à la précédente, take(n) termine le flux après n valeurs. La recherche de la dernière section se sert des deux premiers.

Froid et chaud

Un observable froid crée son producteur à chaque abonnement : c'est le cas de ceux qu'on construit avec new Observable, of ou timer, et de ceux que rend HttpClient. Un observable chaud diffuse un producteur qui existe sans lui et que tous ses abonnés partagent : router.events, les valueChanges d'un formulaire, un Subject. Qui s'abonne tard à un flux chaud a manqué ce qui est déjà passé ; qui s'abonne deux fois à un flux froid déclenche deux exécutions. C'est ce second cas qui coûte, quand il n'est pas voulu : deux pipes async sur le même http.get envoient deux requêtes.

import { AsyncPipe } from '@angular/common';
import { HttpClient } from '@angular/common/http';
import { ChangeDetectionStrategy, Component, inject } from '@angular/core';

interface Profil {
  readonly nom: string;
  readonly commandes: number;
}

@Component({
  selector: 'app-profil',
  imports: [AsyncPipe],
  changeDetection: ChangeDetectionStrategy.OnPush,
  template: `
    <h2>{{ (profil$ | async)?.nom }}</h2>
    <p>Commandes : {{ (profil$ | async)?.commandes }}</p>
  `,
})
export class ProfilCarte {
  // http.get est froid : il ne fait rien a sa creation, et chaque abonnement
  // envoie sa propre requete. Deux pipes async, deux abonnements : deux GET
  // /api/profil partent au premier rendu, pour afficher la meme reponse.
  protected readonly profil$ = inject(HttpClient).get<Profil>('/api/profil');
}

shareReplay rend le flux chaud : un seul abonnement à la source, partagé, et la dernière valeur rejouée aux abonnés tardifs.

import { AsyncPipe } from '@angular/common';
import { HttpClient } from '@angular/common/http';
import { ChangeDetectionStrategy, Component, inject } from '@angular/core';
import { shareReplay } from 'rxjs';

interface Profil {
  readonly nom: string;
  readonly commandes: number;
}

@Component({
  selector: 'app-profil',
  imports: [AsyncPipe],
  changeDetection: ChangeDetectionStrategy.OnPush,
  template: `
    <h2>{{ (profil$ | async)?.nom }}</h2>
    <p>Commandes : {{ (profil$ | async)?.commandes }}</p>
  `,
})
export class ProfilCarte {
  // shareReplay rend le flux chaud : le premier abonne declenche la requete,
  // les suivants se branchent sur la meme execution, et un abonne tardif recoit
  // la derniere valeur gardee (bufferSize: 1). Une seule requete part.
  // refCount: true desabonne la source quand le dernier abonne s'en va ;
  // sans lui, une source sans fin (un minuteur, un flux d'evenements)
  // continuerait de tourner apres la destruction de tous ses abonnes.
  protected readonly profil$ = inject(HttpClient)
    .get<Profil>('/api/profil')
    .pipe(shareReplay({ bufferSize: 1, refCount: true }));
}

share partage aussi, mais sans mémoire : un abonné arrivé après la réponse relance une requête, parce que share se réinitialise quand la source se termine. Et shareReplay(1) sans refCount: true garde la source abonnée quand plus personne n'écoute : sans conséquence pour une requête, qui se termine, une fuite pour un interval, qui continue de tourner. Dans un composant, le plus simple reste souvent de s'abonner une seule fois, par toSignal, et de lire le signal autant de fois qu'il le faut.

Subject, BehaviorSubject, ReplaySubject

Un Subject est à la fois un observable et un observateur : on l'alimente par next, error et complete, et il diffuse chaque valeur à ses abonnés du moment. C'est le flux chaud élémentaire. Ses variantes ne diffèrent que par ce qu'elles donnent à un abonné qui arrive après coup.

import { BehaviorSubject, ReplaySubject, Subject } from 'rxjs';

// Subject : un evenement. Qui arrive apres coup a manque ce qui est passe.
const clics = new Subject<string>();
clics.next('premier');
clics.subscribe((clic) => console.log('Subject :', clic));
clics.next('second');
// Subject : second

// BehaviorSubject : un etat. Il exige une valeur initiale, la donne a chaque
// nouvel abonne, et l'expose en lecture synchrone par value.
const panier = new BehaviorSubject<number>(0);
panier.next(3);
panier.subscribe((articles) => console.log('BehaviorSubject :', articles));
panier.next(4);
console.log('value :', panier.value);
// BehaviorSubject : 3
// BehaviorSubject : 4
// value : 4

// ReplaySubject(n) : un historique. Il rejoue ses n dernieres valeurs.
const journal = new ReplaySubject<string>(2);
journal.next('connexion');
journal.next('recherche');
journal.next('achat');
journal.subscribe((ligne) => console.log('ReplaySubject :', ligne));
// ReplaySubject : recherche
// ReplaySubject : achat

// Un Subject est aussi un observateur : il diffuse a tous ses abonnes une
// seule execution. complete le ferme pour tout le monde, definitivement.
const a = new Subject<number>();
a.subscribe((n) => console.log('X recoit', n));
a.subscribe((n) => console.log('Y recoit', n));
a.next(1);
a.complete();
a.next(2); // ignore, sans erreur
// X recoit 1
// Y recoit 1
TypeUn abonné tardif reçoitCe qu'il représente
Subjectrien de ce qui est passéun événement
BehaviorSubjectla valeur courante, initiale obligatoireun état
ReplaySubject(n)les n dernières valeursun historique

Le BehaviorSubject privé, exposé par asObservable(), a longtemps été la façon de tenir un état partagé dans un service Angular. Un signal fait ce travail plus simplement : une valeur toujours lisible, sans abonnement à gérer, que le template suit sans pipe async, avec asReadonly() pour n'exposer que la lecture. Le Subject garde sa place pour ce qui n'a pas de valeur courante, une demande d'enregistrement ou un message à afficher, que personne ne doit recevoir en retard. La règle d'exposition vaut pour lui aussi : un complete appelé par n'importe quel composant fermerait le flux pour tous les autres, et asObservable() le rend impossible.

switchMap, mergeMap, concatMap, exhaustMap

Un clic qui déclenche une requête produit un observable d'observables : chaque valeur de la source en lance un autre. Les quatre opérateurs d'aplatissement s'abonnent à ces flux internes et en transmettent les valeurs. Ils ne diffèrent que par ce qu'ils font quand une nouvelle valeur arrive alors qu'un flux interne est encore en cours.

import {
  Observable,
  concatMap,
  exhaustMap,
  map,
  mergeMap,
  switchMap,
  take,
  tap,
  timer,
  type OperatorFunction,
} from 'rxjs';

// Trois clics, a 0, 200 et 400 ms.
const clics = timer(0, 200).pipe(
  take(3),
  map((i) => i + 1),
);

// Chaque clic lance une requete qui repond en 300 ms, et qui dit quand elle
// part et quand on l'abandonne avant sa reponse.
function requete(clic: number): Observable<number> {
  return new Observable<number>((abonne) => {
    console.log(`  depart ${clic}`);
    let repondu = false;
    const minuteur = setTimeout(() => {
      repondu = true;
      abonne.next(clic);
      abonne.complete();
    }, 300);
    return () => {
      if (!repondu) {
        clearTimeout(minuteur);
        console.log(`  annule ${clic}`);
      }
    };
  });
}

function essayer(nom: string, aplatir: OperatorFunction<number, number>): Promise<void> {
  console.log(`--- ${nom}`);
  return new Promise((fini) =>
    clics
      .pipe(
        tap((clic) => console.log(`clic ${clic}`)),
        aplatir,
      )
      .subscribe({ next: (clic) => console.log(`reponse ${clic}`), complete: fini }),
  );
}

await essayer('mergeMap', mergeMap(requete));
await essayer('concatMap', concatMap(requete));
await essayer('switchMap', switchMap(requete));
await essayer('exhaustMap', exhaustMap(requete));

// La sortie, remise en quatre colonnes (une par essai, dans l'ordre). Chaque
// colonne se lit de haut en bas ; une meme rangee ne met pas en regard des
// instants egaux.
//
// --- mergeMap   --- concatMap  --- switchMap  --- exhaustMap
// clic 1         clic 1         clic 1         clic 1
//   depart 1       depart 1       depart 1       depart 1
// clic 2         clic 2         clic 2         clic 2
//   depart 2     reponse 1        annule 1     reponse 1
// reponse 1        depart 2       depart 2     clic 3
// clic 3         clic 3         clic 3           depart 3
//   depart 3     reponse 2        annule 2     reponse 3
// reponse 2        depart 3       depart 3
// reponse 3      reponse 3      reponse 3
OpérateurNouvelle valeur pendant un flux en coursUsage
mergeMaplance un flux de plus, en parallèletraitements indépendants, ordre indifférent
concatMapattend la fin du flux en coursécritures qui doivent arriver dans l'ordre
switchMapannule le flux en courslecture dont seule la dernière demande compte
exhaustMapignore la nouvelle valeuraction à ne pas doubler : connexion, paiement

Le mauvais choix ne lève aucune erreur : il affiche un écran faux. Une recherche écrite en mergeMap laisse gagner la réponse la plus lente, pas la plus récente.

import { Subject, map, mergeMap, timer, type Observable } from 'rxjs';

// Le serveur repond plus lentement aux termes courts, qui ramenent plus de
// resultats : 300 ms pour « an », 100 ms pour « ang ».
function rechercher(terme: string): Observable<string> {
  return timer(terme.length < 3 ? 300 : 100).pipe(map(() => `resultats pour « ${terme} »`));
}

const termes = new Subject<string>();

termes
  .pipe(mergeMap((terme) => rechercher(terme)))
  .subscribe((resultats) => console.log('affiche', resultats));

termes.next('an');
setTimeout(() => termes.next('ang'), 50);

// affiche resultats pour « ang »   a 150 ms
// affiche resultats pour « an »    a 300 ms : la reponse la plus lente arrive
//                                  la derniere et ecrase la bonne. L'ecran
//                                  montre les resultats d'un terme que
//                                  l'utilisateur a deja remplace.

switchMap abandonne la recherche précédente dès qu'un nouveau terme arrive :

import { Subject, map, switchMap, timer, type Observable } from 'rxjs';

function rechercher(terme: string): Observable<string> {
  return timer(terme.length < 3 ? 300 : 100).pipe(map(() => `resultats pour « ${terme} »`));
}

const termes = new Subject<string>();

termes
  // Un nouveau terme desabonne la recherche en cours avant de lancer la
  // suivante. Avec HttpClient, se desabonner annule la requete HTTP elle-meme.
  .pipe(switchMap((terme) => rechercher(terme)))
  .subscribe((resultats) => console.log('affiche', resultats));

termes.next('an');
setTimeout(() => termes.next('ang'), 50);

// affiche resultats pour « ang »   a 150 ms, et plus rien ensuite : la
//                                  recherche de « an » a ete abandonnee a 50 ms.

L'erreur inverse coûte autant. Un enregistrement en switchMap abandonne la requête précédente si l'utilisateur enregistre deux fois de suite : le navigateur l'annule, sans savoir si le serveur l'a déjà traitée, et sa réponse est perdue. concatMap passe les enregistrements l'un après l'autre, dans l'ordre.

Une erreur termine le flux qui la porte, et tous ceux où elle remonte. Posé sur le flux principal, catchError recueille l'erreur d'une recherche, mais trop tard : l'abonnement à la chaîne est déjà clos, et la recherche ne répond plus jamais, bien que le Subject des termes vive toujours.

import { Subject, catchError, of, switchMap, throwError, type Observable } from 'rxjs';

// La recherche echoue pour « zz », reussit pour tout le reste.
function rechercher(terme: string): Observable<string> {
  return terme === 'zz' ? throwError(() => new Error('503')) : of(`resultats pour « ${terme} »`);
}

const termes = new Subject<string>();

termes
  .pipe(
    switchMap((terme) => rechercher(terme)),
    // L'erreur de la recherche remonte dans la chaine principale et la termine.
    // catchError la remplace par of(...), qui emet une valeur puis se
    // termine : l'abonnement a la chaine est clos. Le Subject termes, lui,
    // vit toujours, mais plus rien n'y est abonne.
    catchError(() => of('recherche indisponible')),
  )
  .subscribe({
    next: (texte) => console.log(texte),
    complete: () => console.log('flux termine'),
  });

termes.next('an');
termes.next('zz');
termes.next('ang');

// resultats pour « an »
// recherche indisponible
// flux termine
//                  <- « ang » ne donne rien : plus personne n'ecoute

Posé sur le flux interne, il absorbe l'erreur avant qu'elle ne remonte.

import { Subject, catchError, of, switchMap, throwError, type Observable } from 'rxjs';

function rechercher(terme: string): Observable<string> {
  return terme === 'zz' ? throwError(() => new Error('503')) : of(`resultats pour « ${terme} »`);
}

const termes = new Subject<string>();

termes
  .pipe(
    // catchError est pose sur l'observable interne : c'est lui, et lui seul,
    // qui se termine en erreur, puis est remplace. Le flux des termes ne voit
    // qu'une valeur de plus et continue.
    switchMap((terme) => rechercher(terme).pipe(catchError(() => of('recherche indisponible')))),
  )
  .subscribe({
    next: (texte) => console.log(texte),
    complete: () => console.log('flux termine'),
  });

termes.next('an');
termes.next('zz');
termes.next('ang');

// resultats pour « an »
// recherche indisponible
// resultats pour « ang »

combineLatest et forkJoin

Combiner plusieurs flux pose deux questions : quand émettre, et avec quelles valeurs. combineLatest attend une première valeur de chacun, puis émet à chaque nouvelle valeur de l'un d'eux, avec la dernière de tous. forkJoin attend que chacun se termine et émet une seule fois leurs dernières valeurs : c'est Task.WhenAll, ou Promise.all, pour des observables. Le piège est un flux qui ne se termine jamais. Un état n'a pas de fin : un BehaviorSubject ou les valueChanges d'un formulaire, forkJoin les attend pour toujours, sans la moindre erreur. Un toObservable ne se termine qu'à la destruction de son composant, qui est donc le seul moment où forkJoin émettrait.

import { BehaviorSubject, forkJoin, map, timer } from 'rxjs';

type Devise = 'EUR' | 'USD';

// La devise choisie est un etat : un BehaviorSubject, qui ne se termine jamais.
const devise = new BehaviorSubject<Devise>('EUR');

// Les tarifs arrivent par une requete : une valeur, puis la fin.
const tarifs = timer(100).pipe(map(() => ({ EUR: 12, USD: 14 })));

// forkJoin attend que chaque source se termine, puis emet une seule fois
// leurs dernieres valeurs. devise ne se terminant jamais, rien ne sort :
// ni valeur, ni erreur, ni fin. Le prix ne s'affiche pas, et rien ne dit
// pourquoi.
forkJoin([devise, tarifs]).subscribe({
  next: ([code, grille]) => console.log('prix :', grille[code], code),
  complete: () => console.log('termine'),
});

setTimeout(() => devise.next('USD'), 200);

// (aucune sortie)

Pour suivre un état, combineLatest ; pour rassembler des requêtes, forkJoin.

import { BehaviorSubject, combineLatest, forkJoin, map, timer } from 'rxjs';

type Devise = 'EUR' | 'USD';

const devise = new BehaviorSubject<Devise>('EUR');
const tarifs = timer(100).pipe(map(() => ({ EUR: 12, USD: 14 })));

// combineLatest attend une premiere valeur de chaque source, puis emet a
// chaque nouvelle valeur de l'une d'elles, avec la derniere de chacune.
combineLatest([devise, tarifs]).subscribe(([code, grille]) =>
  console.log('prix :', grille[code], code),
);

setTimeout(() => devise.next('USD'), 200);

// forkJoin reste l'outil des requetes paralleles qui se terminent : une
// seule emission, quand toutes ont repondu. La forme objet nomme les valeurs.
const client = timer(80).pipe(map(() => 'Ada'));
const commandes = timer(120).pipe(map(() => 3));

forkJoin({ client, commandes }).subscribe((vue) => console.log(vue));

// prix : 12 EUR                     a 100 ms, quand les tarifs arrivent
// { client: 'Ada', commandes: 3 }   a 120 ms, une seule fois
// prix : 14 USD                     a 200 ms, quand la devise change

Trois cas limites complètent la règle. Un flux qui se termine sans rien émettre fait terminer forkJoin sans valeur. L'erreur d'un seul flux fait échouer l'ensemble et désabonne les autres. Et combineLatest reste muet tant qu'une de ses sources n'a rien émis.

Se désabonner

Un abonnement vit jusqu'à ce que sa source se termine ou qu'on s'en désabonne. Une requête HTTP se termine d'elle-même après sa réponse ; un Subject de service, router.events, valueChanges ou interval, jamais. Un tel abonnement ouvert par un composant survit à sa destruction : il continue de travailler, et la fonction abonnée, qui référence le composant, l'empêche d'être libéré.

Le pipe async et toSignal se désabonnent seuls. Pour un subscribe écrit à la main, takeUntilDestroyed termine le flux quand le contexte qui l'a créé — composant, directive ou service — est détruit. Sa place dans la chaîne compte. Il se désabonne de ce qui est au-dessus de lui et envoie complete vers le bas ; un opérateur d'aplatissement placé en dessous, dont le flux interne est encore actif, attend la fin de ce flux avant de se terminer, et lui survit. Les deux versions qui suivent s'appuient sur ces services :

// notifications.ts
// HttpClient s'injecte sans rien declarer depuis Angular 21 ; provideHttpClient()
// ne sert qu'a le configurer (voir le cours HttpClient).
import { HttpClient } from '@angular/common/http';
import { Injectable, inject } from '@angular/core';
import { BehaviorSubject, type Observable } from 'rxjs';

@Injectable({ providedIn: 'root' })
export class Session {
  // L'utilisateur connecte : un etat qui dure autant que l'application, et
  // un flux qui ne se termine donc jamais.
  readonly utilisateur$ = new BehaviorSubject('ada');
}

@Injectable({ providedIn: 'root' })
export class Notifications {
  private readonly http = inject(HttpClient);

  nonLues(utilisateur: string): Observable<number> {
    return this.http.get<number>(`/api/notifications/${utilisateur}/non-lues`);
  }
}

Placé en tête, takeUntilDestroyed ne termine que le flux des utilisateurs, et le sondage lancé par switchMap lui survit.

import { ChangeDetectionStrategy, Component, inject, signal } from '@angular/core';
import { takeUntilDestroyed } from '@angular/core/rxjs-interop';
import { switchMap, timer } from 'rxjs';
import { Notifications, Session } from './notifications';

@Component({
  selector: 'app-cloche',
  changeDetection: ChangeDetectionStrategy.OnPush,
  template: `<span>{{ nonLues() }} non lues</span>`,
})
export class Cloche {
  protected readonly nonLues = signal(0);

  constructor() {
    const notifications = inject(Notifications);

    inject(Session)
      .utilisateur$.pipe(
        // takeUntilDestroyed termine ce qui est au-dessus de lui : le flux des
        // utilisateurs. switchMap, en dessous, ne se termine que lorsque sa
        // source ET l'observable interne en cours sont termines ; or timer(0,
        // 5000) ne se termine jamais. A la destruction du composant, le
        // sondage continue toutes les 5 secondes, et retient le composant.
        takeUntilDestroyed(),
        switchMap((utilisateur) =>
          timer(0, 5000).pipe(switchMap(() => notifications.nonLues(utilisateur))),
        ),
      )
      .subscribe((nombre) => this.nonLues.set(nombre));
  }
}

Placé en dernier, il ne laisse rien derrière lui.

import { ChangeDetectionStrategy, Component, inject, signal } from '@angular/core';
import { takeUntilDestroyed } from '@angular/core/rxjs-interop';
import { switchMap, timer } from 'rxjs';
import { Notifications, Session } from './notifications';

@Component({
  selector: 'app-cloche',
  changeDetection: ChangeDetectionStrategy.OnPush,
  template: `<span>{{ nonLues() }} non lues</span>`,
})
export class Cloche {
  protected readonly nonLues = signal(0);

  constructor() {
    const notifications = inject(Notifications);

    inject(Session)
      .utilisateur$.pipe(
        switchMap((utilisateur) =>
          timer(0, 5000).pipe(switchMap(() => notifications.nonLues(utilisateur))),
        ),
        // En dernier, il coupe la chaine entiere : se desabonner du resultat
        // desabonne chaque operateur en remontant, sondage compris. Appele
        // sans argument, il lui faut un contexte d'injection, ici le
        // constructeur ; ailleurs, on lui passe un DestroyRef injecte.
        takeUntilDestroyed(),
      )
      .subscribe((nombre) => this.nonLues.set(nombre));
  }
}

Le sondage appelle le service à 0, 5 et 10 secondes. Dix secondes après la destruction du composant, la première version l'a appelé deux fois de plus, la seconde plus du tout. Sans argument, takeUntilDestroyed injecte lui-même le DestroyRef et exige un contexte d'injection : appelé dans ngOnInit, il lève NG0203. On injecte alors DestroyRef dans un champ, et on le lui passe : takeUntilDestroyed(this.destroyRef).

toSignal et toObservable

toSignal s'abonne à un observable et expose sa dernière valeur en signal ; toObservable fait l'inverse. Le premier est la sortie normale d'un flux vers le template : la vue lit un signal, il n'y a qu'un abonnement, et il prend fin à la destruction. Ce que toSignal impose découle de ce qu'il fait : il s'abonne sur le champ, et il doit savoir quand se désabonner.

  • Un contexte d'injection. Il y cherche le DestroyRef qui mettra fin à l'abonnement. Hors de ce contexte, l'option injector le lui fournit ; manualCleanup: true l'en dispense aussi, pour une source qui se termine d'elle-même.
  • Une valeur initiale. Un signal a toujours une valeur, un observable pas forcément tout de suite. Sans option, le type est Signal<T | undefined> ; initialValue fixe la valeur d'attente et retire undefined du type ; requireSync: true convient à une source qui émet dès l'abonnement, comme un BehaviorSubject, et lève NG0601 si elle ne le fait pas.
  • Des erreurs traitées en amont. Une erreur de l'observable est relancée à chaque lecture du signal, donc à chaque rendu du template qui le lit : le catchError se place avant toSignal.
  • Hors d'un contexte réactif. Appelé depuis un template ou un computed, il créerait un abonnement à chaque évaluation ; en mode développement, Angular le refuse avec NG0602.

Le premier point se viole facilement par habitude : on déplace l'appel dans ngOnInit, là où l'on s'abonnait autrefois.

import {
  ChangeDetectionStrategy,
  Component,
  inject,
  signal,
  type OnInit,
  type Signal,
} from '@angular/core';
import { toSignal } from '@angular/core/rxjs-interop';
import { Session } from './notifications';

@Component({
  selector: 'app-bandeau',
  changeDetection: ChangeDetectionStrategy.OnPush,
  template: `<p>Connecte : {{ utilisateur() }}</p>`,
})
export class Bandeau implements OnInit {
  private readonly session = inject(Session);
  protected utilisateur: Signal<string> = signal('');

  ngOnInit(): void {
    // toSignal s'abonne tout de suite et doit se desabonner a la destruction :
    // il cherche le DestroyRef dans le contexte d'injection courant. ngOnInit
    // n'en est pas un, et le premier rendu echoue (message du mode
    // developpement, Angular 22.1) :
    // NG0203: toSignal() can only be used within an injection context such as
    // a constructor, a factory function, a field initializer, or a function
    // used with `runInInjectionContext`. Find more at
    // https://v22.angular.dev/errors/NG0203
    this.utilisateur = toSignal(this.session.utilisateur$, { requireSync: true });
  }
}

L'initialiseur de champ est le bon endroit, et le code y gagne en brièveté.

import { ChangeDetectionStrategy, Component, inject } from '@angular/core';
import { toSignal } from '@angular/core/rxjs-interop';
import { Session } from './notifications';

@Component({
  selector: 'app-bandeau',
  changeDetection: ChangeDetectionStrategy.OnPush,
  template: `<p>Connecte : {{ utilisateur() }}</p>`,
})
export class Bandeau {
  // Un initialiseur de champ s'execute pendant la construction : c'est un
  // contexte d'injection. Un BehaviorSubject emet des l'abonnement, ce que
  // requireSync exige (sinon NG0601) : le type est Signal<string>, sans
  // undefined et sans valeur initiale a inventer.
  protected readonly utilisateur = toSignal(inject(Session).utilisateur$, { requireSync: true });
}

toObservable lit le signal dans un effect, avec les mêmes exigences de contexte. Il n'émet donc jamais au moment de l'écriture, mais au passage suivant des effects, et trois set successifs n'y produisent qu'une valeur, la dernière. Pour une saisie, c'est sans conséquence, et toObservable devient l'entrée d'une chaîne RxJS branchée sur un signal : les opérateurs de temps au milieu, toSignal à la sortie.

import { HttpClient } from '@angular/common/http';
import { ChangeDetectionStrategy, Component, Injectable, inject, signal } from '@angular/core';
import { toObservable, toSignal } from '@angular/core/rxjs-interop';
import {
  catchError,
  debounceTime,
  distinctUntilChanged,
  map,
  of,
  switchMap,
  type Observable,
} from 'rxjs';

export interface Produit {
  readonly reference: string;
  readonly libelle: string;
}

@Injectable({ providedIn: 'root' })
export class Catalogue {
  private readonly http = inject(HttpClient);

  rechercher(terme: string): Observable<Produit[]> {
    return this.http.get<Produit[]>('/api/produits', { params: { q: terme } });
  }
}

@Component({
  selector: 'app-recherche',
  changeDetection: ChangeDetectionStrategy.OnPush,
  template: `
    <input #champ type="search" [value]="terme()" (input)="terme.set(champ.value)" />
    @for (produit of resultats(); track produit.reference) {
      <p>{{ produit.libelle }}</p>
    } @empty {
      <p>Aucun resultat.</p>
    }
  `,
})
export class Recherche {
  private readonly catalogue = inject(Catalogue);

  // L'etat saisi est un signal, comme tout ce que le template lit.
  protected readonly terme = signal('');

  // Du signal vers l'observable, pour ce que les signals ne savent pas faire :
  // attendre, ecarter, annuler. Puis retour en signal pour le template.
  protected readonly resultats = toSignal(
    toObservable(this.terme).pipe(
      map((terme) => terme.trim()),
      debounceTime(300), // n'emet qu'apres 300 ms sans nouvelle frappe
      // « angu » efface avant 300 ms : debounceTime ne laisse passer que « ang »,
      // egal au precedent, et distinctUntilChanged l'ecarte. Pas de seconde requete.
      distinctUntilChanged(),
      switchMap((terme) =>
        terme.length < 2
          ? of<Produit[]>([])
          : this.catalogue.rechercher(terme).pipe(catchError(() => of<Produit[]>([]))),
      ),
    ),
    // Sans valeur initiale, le type serait Signal<Produit[] | undefined> :
    // rien n'est emis avant 300 ms.
    { initialValue: [] },
  );
}

Le rail de navigation de ce site n'a besoin que de la seconde moitié du trajet : il filtre router.events sur NavigationEnd, en tire l'URL courante et la passe à toSignal. Quand il n'y a rien à attendre ni à espacer, seulement une donnée à charger d'après des paramètres, rxResource, stable depuis Angular 22.0, suffit : c'est le resource du cours Signals, avec un observable pour chargeur, rendu par l'option stream et abandonné, désabonnement compris, quand les paramètres changent.

Ce cours vous a servi ? Offrir un café Signaler une erreur