NestJS Enterprise Backend APIs · Lekcja

Sagi dla długotrwałych przepływów pracy

Koordynuj wieloetapowe procesy reaktywnie za pomocą sag CQRS opartych na RxJS.

Lekcja 4 z 413 kroki

Sagi dla długotrwałych przepływów pracy to bezpłatna lekcja NestJS Enterprise Backend APIs na CoddyKit. To lekcja 4 z 4. Możesz przeczytać całą lekcję poniżej za darmo — a potem ćwiczyć ją interaktywnie w przeglądarce z wbudowanym edytorem kodu i tutorem AI dostępnym 24/7. To część ścieżki edukacyjnej NestJS Enterprise Backend APIs, a Twój postęp synchronizuje się między webem a aplikacją CoddyKit. Kurs NestJS Enterprise Backend APIs zawiera 4 lekcji w sumie.

Części tej lekcji nie zostały jeszcze przetłumaczone i są wyświetlane po angielsku.

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.
Bezpłatny start

Ucz się TypeScript dzięki korepetycjom AI — za darmo

Pisz i uruchamiaj kod w przeglądarce, otrzymuj natychmiastową pomoc od korepetytora AI dostępnego 24/7 i kontynuuj naukę w sieci lub w aplikacji.

Kursy
20
Lekcje
76

Często zadawane pytania

Czy lekcja „Sagi dla długotrwałych przepływów pracy” jest bezpłatna?

Tak — pełny tekst „Sagi dla długotrwałych przepływów pracy” jest dostępny za darmo tutaj w sieci. Aby ćwiczyć ją interaktywnie (wbudowany edytor kodu i tutor AI dostępny 24/7) i odblokować resztę kursu NestJS Enterprise Backend APIs, przejdź na CoddyKit PRO. Kurs NestJS Enterprise Backend APIs zawiera 4 lekcji w sumie.

Co nauczysz się w „Sagi dla długotrwałych przepływów pracy”?

Koordynuj wieloetapowe procesy reaktywnie za pomocą sag CQRS opartych na RxJS. Ćwiczysz NestJS Enterprise Backend APIs z praktycznym kodem, który uruchamiasz bezpośrednio w przeglądarce, a tutor AI dostępny 24/7 odpowiada na Twoje pytania podczas pracy nad lekcją.

Czy potrzebuję doświadczenia, aby zacząć NestJS Enterprise Backend APIs?

Nie wymagamy żadnego doświadczenia. NestJS Enterprise Backend APIs w CoddyKit jest strukturyzowany dla początkujących i zaawansowanych użytkowników, więc możesz zacząć tutaj lub od początku i uczyć się w swoim tempie. To lekcja 4 z 4.

Ile czasu zajmuje lekcja „Sagi dla długotrwałych przepływów pracy”?

Większość lekcji CoddyKit trwa około 5–10 minut. Każda lekcja to mały, interaktywny krok, dzięki czemu robisz systematyczne postępy i zawsze wracasz dokładnie do tego samego miejsca — na webie i w aplikacji.

Czy mogę pisać i uruchamiać kod w tej lekcji NestJS Enterprise Backend APIs?

Tak. Każda lekcja NestJS Enterprise Backend APIs zawiera wbudowany edytor kodu, więc piszesz i uruchamiasz prawdziwy kod bezpośrednio w przeglądarce i od razu otrzymujesz sprzężenie zwrotne od AI — bez konfiguracji na komputerze.

Wszystkie lekcje w tym kursie

  1. Polecenia, handlery i magistrala poleceń
  2. Zapytania i projekcje modeli odczytu
  3. Zdarzenia domenowe i AggregateRoot
  4. Sagi dla długotrwałych przepływów pracy
← Powrót do NestJS Enterprise Backend APIs