NestJS Enterprise Backend APIs · Lezione

Saga per workflow di lunga durata

Coordini processi composti da più fasi in modo reattivo con le saga CQRS basate su RxJS.

Lezione 4 di 413 passaggi

Saga per workflow di lunga durata è una lezione NestJS Enterprise Backend APIs gratuita su CoddyKit. Questa è la lezione 4 di 4. Puoi leggere la lezione completa qui gratuitamente — poi esercitati direttamente nel browser con un editor di codice integrato e un tutor IA disponibile 24/7. Fa parte del percorso di apprendimento NestJS Enterprise Backend APIs, e i tuoi progressi si sincronizzano tra il web e l'app CoddyKit. Il corso NestJS Enterprise Backend APIs include 4 lezioni in totale.

Parti di questa lezione non sono ancora state tradotte e vengono mostrate in inglese.

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.
Gratis per iniziare

Impara TypeScript con un tutor IA — gratis

Scrivi ed esegui vero codice nel tuo browser, ricevi aiuto istantaneo da un tutor IA disponibile 24/7, e riprendi da dove hai lasciato sul web o nell'app.

Corsi
20
Lezioni
76

Domande Frequenti

La lezione «Saga per workflow di lunga durata» è gratuita?

Sì — il testo completo di «Saga per workflow di lunga durata» è gratuito qui sul web. Per esercitarvi in modo interattivo (un editor di codice integrato e un tutor IA 24/7) e sbloccare il resto del corso NestJS Enterprise Backend APIs, passa a CoddyKit PRO. Il corso NestJS Enterprise Backend APIs include 4 lezioni in totale.

Cosa imparerò in «Saga per workflow di lunga durata»?

Coordini processi composti da più fasi in modo reattivo con le saga CQRS basate su RxJS. Eserciti NestJS Enterprise Backend APIs con codice pratico che esegui direttamente nel browser, e un tutor IA 24/7 risponde alle tue domande mentre lavori sulla lezione.

Ho bisogno di esperienza per iniziare NestJS Enterprise Backend APIs?

Non è richiesta alcuna esperienza precedente. NestJS Enterprise Backend APIs su CoddyKit è strutturato per principianti e studenti avanzati, quindi puoi iniziare da qui o dall'inizio e procedere al tuo ritmo. Questa è la lezione 4 di 4.

Quanto tempo richiede la lezione «Saga per workflow di lunga durata»?

La maggior parte delle lezioni CoddyKit richiede circa 5–10 minuti. Ogni lezione è breve e interattiva, quindi fai progressi costanti e riprendi esattamente da dove hai lasciato su web e app.

Posso scrivere ed eseguire codice in questa lezione NestJS Enterprise Backend APIs?

Sì. Ogni lezione NestJS Enterprise Backend APIs include un editor di codice integrato, quindi scrivi ed esegui codice reale direttamente nel tuo browser e ricevi feedback istantaneo dall'IA — nessuna configurazione locale necessaria.

Tutte le lezioni di questo corso

  1. Comandi, handler e il bus dei comandi
  2. Query e proiezioni dei read model
  3. Eventi di dominio e AggregateRoot
  4. Saga per workflow di lunga durata
← Torna a NestJS Enterprise Backend APIs