장시간 실행되는 워크플로를 위한 사가
RxJS 기반 CQRS 사가로 여러 단계의 프로세스를 반응형으로 조정합니다.
장시간 실행되는 워크플로를 위한 사가은(는) CoddyKit의 무료 NestJS Enterprise Backend APIs 강의입니다. 이것은 4개 중 4번째 강의입니다. 아래에서 전체 강의를 무료로 읽을 수 있으며, 내장 코드 에디터와 24/7 AI 튜터와 함께 브라우저에서 직접 실습할 수 있습니다. 이 강의는 NestJS Enterprise Backend APIs 학습 경로의 일부이며, 진행 상황이 웹과 CoddyKit 앱에 동기화됩니다. NestJS Enterprise Backend APIs 강의에는 총 4개의 강의가 포함되어 있습니다.
이 강의의 일부는 아직 번역되지 않았으며 영어로 표시됩니다.
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, andmergeMap.
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
mapover 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
EMPTYto emit no command (skip). mergeMapruns 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
providersalongside command and event handlers. - No manual subscription, no
onModuleInitwiring 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
catchErrorinside the innermergeMapso 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 globalObservableof events to anObservable<ICommand>that theCommandBusdispatches automatically. - Start every saga with
ofType(...)to filter and type-narrow; usemapfor 1:1 andmergeMapfor 0/1/many commands per event. - Register sagas as
providersin aCqrsModule— 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
catchErrorso one failure can't kill the saga, and make command handlers idempotent for at-least-once delivery.
자주 묻는 질문
“장시간 실행되는 워크플로를 위한 사가” 강의는 무료인가요?
네 — “장시간 실행되는 워크플로를 위한 사가” 전체 내용을 이 웹사이트에서 무료로 읽을 수 있습니다. 인터랙티브하게 실습하려면(내장 코드 에디터와 24/7 AI 튜터), CoddyKit PRO로 업그레이드하면 NestJS Enterprise Backend APIs 강의 전체를 잠금 해제할 수 있습니다. NestJS Enterprise Backend APIs 강의에는 총 4개의 강의가 포함되어 있습니다.
“장시간 실행되는 워크플로를 위한 사가”에서 뭘 배우나요?
RxJS 기반 CQRS 사가로 여러 단계의 프로세스를 반응형으로 조정합니다. 브라우저에서 직접 실행하는 실습 코드로 NestJS Enterprise Backend APIs을(를) 배우며, 24/7 AI 튜터가 강의를 진행하면서 질문에 답변해줍니다.
NestJS Enterprise Backend APIs을(를) 시작하는 데 경험이 필요한가요?
사전 경험은 필요하지 않습니다. CoddyKit의 NestJS Enterprise Backend APIs은(는) 초급자부터 고급 학습자까지를 위해 구성되어 있으므로, 여기서 시작하거나 처음부터 시작할 수 있으며 자신의 속도대로 진행할 수 있습니다. 이것은 4개 중 4번째 강의입니다.
“장시간 실행되는 워크플로를 위한 사가” 강의는 얼마나 걸리나요?
대부분의 CoddyKit 강의는 약 5~10분이 소요됩니다. 각 강의는 간결하고 인터랙티브하여 꾸준한 진행이 가능하며, 웹과 앱에서 중단한 부분부터 바로 시작할 수 있습니다.
이 NestJS Enterprise Backend APIs 강의에서 코드를 작성하고 실행할 수 있나요?
네. 모든 NestJS Enterprise Backend APIs 강의에는 내장 코드 에디터가 포함되어 있으므로, 브라우저에서 바로 실제 코드를 작성하고 실행한 후 즉시 AI 피드백을 받을 수 있습니다 — 로컬 설정이 필요 없습니다.
이 강의의 모든 강의
- 명령, 처리기 및 CommandBus
- 쿼리와 읽기 모델 프로젝션
- 도메인 이벤트와 AggregateRoot
- 장시간 실행되는 워크플로를 위한 사가