NestJS Enterprise Backend APIs · Leçon

Sagas pour les processus de longue durée

Coordonnez réactivement des processus en plusieurs étapes avec des sagas CQRS fondées sur RxJS.

Leçon 4 sur 413 étapes

Sagas pour les processus de longue durée est une leçon NestJS Enterprise Backend APIs gratuite sur CoddyKit. Ceci est la leçon 4 sur 4. Tu peux lire la leçon complète ci-dessous gratuitement — puis la pratiquer en direct dans le navigateur avec un éditeur de code intégré et un tuteur IA 24/7. Elle fait partie du parcours d'apprentissage NestJS Enterprise Backend APIs, et ta progression se synchronise sur le web et l'application CoddyKit. Le cours NestJS Enterprise Backend APIs comprend 4 leçons au total.

Certaines parties de cette leçon n'ont pas encore été traduites et s'affichent en anglais.

Why Sagas Exist

In an event-driven system, a single business outcome often requires many steps across multiple aggregates: place order, reserve stock, charge payment, schedule shipping. No single command handler owns that whole flow.

A Saga coordinates such a long-running workflow by listening to events and reacting with new commands. It is the glue that turns one event into the next step of a process.

  • Reactive: a saga is triggered by events, not called directly.
  • Stateless dispatcher: in NestJS, a CQRS saga maps an event stream to a command stream.
  • Decoupled: handlers stay small; the saga owns orchestration.

Sagas in @nestjs/cqrs

In @nestjs/cqrs a saga is a class method decorated with @Saga() that receives an RxJS Observable of all published events and returns an Observable<ICommand>.

Whatever commands the returned stream emits are automatically dispatched through the CommandBus. The framework subscribes to your stream for you.

  • Input type: Observable<any> (the global event stream).
  • Output type: Observable<ICommand>.
  • You shape the flow with RxJS operators like ofType, map, and mergeMap.
import { Injectable } from '@nestjs/common';
import { ICommand, ofType, Saga } from '@nestjs/cqrs';
import { Observable } from 'rxjs';
import { map } from 'rxjs/operators';
import { OrderCreatedEvent } from './events/order-created.event';
import { ReserveStockCommand } from './commands/reserve-stock.command';

@Injectable()
export class OrderSagas {
  @Saga()
  orderCreated = (events$: Observable<any>): Observable<ICommand> => {
    return events$.pipe(
      ofType(OrderCreatedEvent),
      map((event) => new ReserveStockCommand(event.orderId, event.items)),
    );
  };
}

The ofType Operator

The global stream carries every published event. You almost always start a saga by filtering it down to the event types you care about with ofType(...EventClasses).

ofType is a custom RxJS operator shipped by @nestjs/cqrs. It filters by event constructor and, crucially, narrows the TypeScript type of downstream values to that event, so event.orderId type-checks.

  • Pass one or more event classes: ofType(A, B).
  • Always filter early — never map over the raw stream blindly.

From Event to Command

The simplest saga is a one-to-one translation: each matching event produces exactly one command. map is the right operator here because it is synchronous and emits one value per input.

The pattern below reacts to a successful stock reservation by issuing a payment command — moving the workflow forward one step.

import { Injectable } from '@nestjs/common';
import { ICommand, ofType, Saga } from '@nestjs/cqrs';
import { Observable } from 'rxjs';
import { map } from 'rxjs/operators';
import { StockReservedEvent } from './events/stock-reserved.event';
import { ChargePaymentCommand } from './commands/charge-payment.command';

@Injectable()
export class PaymentSagas {
  @Saga()
  stockReserved = (events$: Observable<any>): Observable<ICommand> =>
    events$.pipe(
      ofType(StockReservedEvent),
      map((e) => new ChargePaymentCommand(e.orderId, e.amount)),
    );
}

map vs mergeMap

Use map when one event yields exactly one command synchronously. Use mergeMap (a.k.a. flatMap) when a single event must produce zero, one, or many commands, or when you need an inner Observable.

  • mergeMap + of(...) lets you emit several commands per event.
  • Return EMPTY to emit no command (skip).
  • mergeMap runs inner streams concurrently — good for independent fan-out steps.
import { Injectable } from '@nestjs/common';
import { ICommand, ofType, Saga } from '@nestjs/cqrs';
import { Observable, of, EMPTY } from 'rxjs';
import { mergeMap } from 'rxjs/operators';
import { PaymentConfirmedEvent } from './events/payment-confirmed.event';
import { ScheduleShippingCommand } from './commands/schedule-shipping.command';
import { SendReceiptCommand } from './commands/send-receipt.command';

@Injectable()
export class FulfillmentSagas {
  @Saga()
  paymentConfirmed = (events$: Observable<any>): Observable<ICommand> =>
    events$.pipe(
      ofType(PaymentConfirmedEvent),
      mergeMap((e) =>
        e.amount > 0
          ? of(
              new ScheduleShippingCommand(e.orderId),
              new SendReceiptCommand(e.orderId, e.amount),
            )
          : EMPTY,
      ),
    );
}

Registering Sagas in a Module

A saga class only runs if it is listed in the module's providers. NestJS's CqrsModule discovers every provider that exposes @Saga() methods and subscribes them to the event bus at bootstrap.

  • Import CqrsModule.
  • Add the saga class to providers alongside command and event handlers.
  • No manual subscription, no onModuleInit wiring needed.
import { Module } from '@nestjs/common';
import { CqrsModule } from '@nestjs/cqrs';
import { OrderSagas } from './sagas/order.sagas';
import { PaymentSagas } from './sagas/payment.sagas';
import { ReserveStockHandler } from './commands/reserve-stock.handler';
import { ChargePaymentHandler } from './commands/charge-payment.handler';

@Module({
  imports: [CqrsModule],
  providers: [
    OrderSagas,
    PaymentSagas,
    ReserveStockHandler,
    ChargePaymentHandler,
  ],
})
export class OrderingModule {}

Correlating Multi-Step State

Most sagas need a correlation id — typically the aggregate id (e.g. orderId) — carried on every event so steps can be tied to the same process instance.

The CQRS saga itself is stateless: it just maps events to commands. The workflow state (which step completed, what was reserved) lives in the aggregate or a dedicated read model, updated by command handlers. The saga reacts to the events those handlers emit.

  • Always include the correlation id in event payloads.
  • Keep mutable progress in an aggregate/persistence layer, not in the saga.

Compensation: The Saga's Real Job

Distributed workflows cannot use a single ACID transaction across services. Instead a saga guarantees consistency through compensating actions: if a later step fails, earlier successful steps are semantically undone.

If payment fails after stock was reserved, the saga reacts to PaymentFailedEvent by dispatching a ReleaseStockCommand. Compensation is forward-recovery, not a rollback — you issue a new command that reverses the effect.

import { Injectable } from '@nestjs/common';
import { ICommand, ofType, Saga } from '@nestjs/cqrs';
import { Observable } from 'rxjs';
import { map } from 'rxjs/operators';
import { PaymentFailedEvent } from './events/payment-failed.event';
import { ReleaseStockCommand } from './commands/release-stock.command';

@Injectable()
export class CompensationSagas {
  @Saga()
  paymentFailed = (events$: Observable<any>): Observable<ICommand> =>
    events$.pipe(
      ofType(PaymentFailedEvent),
      map((e) => new ReleaseStockCommand(e.orderId, e.items)),
    );
}

Errors Inside the Saga Stream

A saga returns one long-lived Observable. If an error escapes that stream, RxJS terminates the subscription and your saga stops reacting to all future events — a silent, system-wide failure.

Protect the stream with catchError. Recover by mapping the failure to a compensating command and continuing, rather than letting the observable die.

  • Never let exceptions propagate out of the saga pipe.
  • Prefer per-event isolation: put catchError inside the inner mergeMap so one bad event doesn't kill the outer stream.
import { Injectable, Logger } from '@nestjs/common';
import { ICommand, ofType, Saga } from '@nestjs/cqrs';
import { Observable, of } from 'rxjs';
import { mergeMap, catchError } from 'rxjs/operators';
import { ShippingRequestedEvent } from './events/shipping-requested.event';
import { NotifyOpsCommand } from './commands/notify-ops.command';
import { DispatchCarrierCommand } from './commands/dispatch-carrier.command';

@Injectable()
export class ShippingSagas {
  private readonly logger = new Logger(ShippingSagas.name);

  @Saga()
  shippingRequested = (events$: Observable<any>): Observable<ICommand> =>
    events$.pipe(
      ofType(ShippingRequestedEvent),
      mergeMap((e) =>
        of(new DispatchCarrierCommand(e.orderId)).pipe(
          catchError((err) => {
            this.logger.error(err);
            return of(new NotifyOpsCommand(e.orderId, 'dispatch failed'));
          }),
        ),
      ),
    );
}

Idempotency and At-Least-Once Delivery

When events are delivered over a broker (Kafka, RabbitMQ) the saga may see the same event more than once after a redelivery or restart. Dispatching the same command twice can double-charge or double-ship.

The fix is not in the saga's RxJS plumbing but in the command handler: make it idempotent by recording a processed key (correlation id + step) and ignoring duplicates.

  • Sagas should produce commands that are safe to retry.
  • Track (orderId, step) in a dedup table or use the aggregate's version.
  • Design for at-least-once, not exactly-once.

A Standalone RxJS Saga Pipeline

You can model the exact event-to-command mapping a saga performs using plain RxJS — no NestJS runtime required. This mirrors how ofType + map drive a workflow, and is great for reasoning about the flow in isolation.

import { from } from 'rxjs';
import { filter, map } from 'rxjs/operators';

class OrderCreated { constructor(public orderId: string) {} }
class PaymentFailed { constructor(public orderId: string) {} }
class ReserveStock { constructor(public orderId: string) {} }
class ReleaseStock { constructor(public orderId: string) {} }

const ofType =
  <T>(type: new (...a: any[]) => T) =>
  (source: any) =>
    source.pipe(filter((e: any): e is T => e instanceof type));

const events$ = from([
  new OrderCreated('A1'),
  new PaymentFailed('A1'),
]);

events$
  .pipe(
    map((e) =>
      e instanceof OrderCreated
        ? new ReserveStock(e.orderId)
        : e instanceof PaymentFailed
          ? new ReleaseStock(e.orderId)
          : null,
    ),
    filter((c): c is ReserveStock | ReleaseStock => c !== null),
  )
  .subscribe((cmd) => console.log('dispatch ->', cmd.constructor.name, cmd.orderId));

Quick Check

A teammate writes a CQRS saga whose returned Observable occasionally throws when building a command for a malformed event. After one such event in production, the saga stops reacting to all subsequent events. What is the correct fix?

Recap

You learned how to coordinate long-running workflows with RxJS-based CQRS sagas in NestJS:

  • A @Saga() maps the global Observable of events to an Observable<ICommand> that the CommandBus dispatches automatically.
  • Start every saga with ofType(...) to filter and type-narrow; use map for 1:1 and mergeMap for 0/1/many commands per event.
  • Register sagas as providers in a CqrsModule — no manual subscription.
  • Carry a correlation id on events; keep workflow state in aggregates/read models, not the saga.
  • Achieve distributed consistency via compensating commands, not transactions.
  • Guard the stream with catchError so one failure can't kill the saga, and make command handlers idempotent for at-least-once delivery.
Gratuit pour commencer

Apprends TypeScript avec un tuteur IA — gratuit

Écris et exécute du vrai code dans ton navigateur, obtiens de l'aide instantanée d'un tuteur IA disponible 24h/24, et reprends là où tu t'es arrêté sur le web ou dans l'app.

Cours
20
Leçons
76

Questions Fréquemment Posées

La leçon « Sagas pour les processus de longue durée » est-elle gratuite ?

Oui — le texte complet de « Sagas pour les processus de longue durée » est gratuit à lire ici sur le web. Pour la pratiquer de manière interactive (un éditeur de code intégré et un tuteur IA 24/7) et déverrouiller le reste du cours NestJS Enterprise Backend APIs, passe à CoddyKit PRO. Le cours NestJS Enterprise Backend APIs comprend 4 leçons au total.

Qu'est-ce que j'apprendrai dans « Sagas pour les processus de longue durée » ?

Coordonnez réactivement des processus en plusieurs étapes avec des sagas CQRS fondées sur RxJS. Tu pratiques NestJS Enterprise Backend APIs avec du code pratique que tu exécutes directement dans le navigateur, et un tuteur IA 24/7 répond à tes questions au fur et à mesure que tu avances dans la leçon.

Dois-je avoir de l'expérience pour commencer NestJS Enterprise Backend APIs ?

Aucune expérience préalable n'est requise. NestJS Enterprise Backend APIs sur CoddyKit est structuré pour les débutants jusqu'aux apprenants avancés, donc tu peux commencer ici ou depuis le début et avancer à ton rythme. Ceci est la leçon 4 sur 4.

Combien de temps prend la leçon « Sagas pour les processus de longue durée » ?

La plupart des leçons CoddyKit prennent environ 5–10 minutes. Chacune est courte et interactive, tu progresses régulièrement et tu repiques exactement où tu t'es arrêté sur le web et l'app.

Peux-tu écrire et exécuter du code dans cette leçon NestJS Enterprise Backend APIs ?

Oui. Chaque leçon NestJS Enterprise Backend APIs inclut un éditeur de code intégré, tu écris et exécutes du vrai code directement dans ton navigateur et tu reçois des retours IA instantanés — aucune configuration locale requise.

Toutes les leçons de ce cours

  1. Commandes, gestionnaires et bus de commandes
  2. Requêtes et projections des modèles de lecture
  3. Événements de domaine et AggregateRoot
  4. Sagas pour les processus de longue durée
← Retour à NestJS Enterprise Backend APIs