Пошук уроків, статей та іншого контенту
Налаштуєте EventBus, обробники подій, підписки та обробку помилок для внутрішньої комунікації модулів.
EventBus у NestJS — це внутрішній механізм публікації подій між модулями застосунку. Він входить до пакета @nestjs/cqrs.
Подія описує факт, який уже відбувся:
користувача зареєстровано;
замовлення оплачено;
файл оброблено;
підписку скасовано.
Модуль, який створює подію, не повинен знати, які інші модулі на неї реагують. Це зменшує зв’язаність між компонентами.
Наприклад, модуль користувачів може опублікувати UserRegisteredEvent, а такі обробники можуть реагувати на неї незалежно:
модуль повідомлень надсилає вітальний лист;
модуль аудиту записує подію;
модуль аналітики оновлює статистику.
EventBus працює лише в межах одного процесу Node.js. Він не є брокером повідомлень і не забезпечує довготривале зберігання подій.
Встановіть пакет:
npm install @nestjs/cqrsПідключіть CqrsModule у кореневому модулі:
import { Module } from '@nestjs/common';
import { CqrsModule } from '@nestjs/cqrs';
@Module({
imports: [CqrsModule.forRoot()],
})
export class AppModule {}CqrsModule.forRoot() реєструє EventBus, а також інші компоненти CQRS. Після цього EventBus можна інжектити у провайдери, сервіси й контролери.
Подія — це звичайний клас, який реалізує IEvent.
import { IEvent } from '@nestjs/cqrs';
export class UserRegisteredEvent implements IEvent {
constructor(
public readonly userId: string,
public readonly email: string,
public readonly registeredAt: Date,
) {}
}Подія повинна містити дані, необхідні обробникам для роботи. Її властивості зазвичай роблять readonly, щоб після публікації подію не можна було змінити.
Клас події важливий не лише як тип TypeScript. EventBus використовує конструктор класу під час пошуку відповідних обробників:
eventBus.publish(new UserRegisteredEvent(
'user-1',
'user@example.com',
new Date(),
));Інтерфейс без класу для цього не підходить, оскільки після компіляції TypeScript інтерфейси не існують у JavaScript.
Обробник позначається декоратором @EventsHandler() і реалізує IEventHandler<T>.
import { EventsHandler, IEventHandler } from '@nestjs/cqrs';
import { Injectable, Logger } from '@nestjs/common';
import { UserRegisteredEvent } from './user-registered.event';
@Injectable()
@EventsHandler(UserRegisteredEvent)
export class SendWelcomeEmailHandler
implements IEventHandler<UserRegisteredEvent>
{
private readonly logger = new Logger(SendWelcomeEmailHandler.name);
async handle(event: UserRegisteredEvent): Promise<void> {
this.logger.log(
`Надсилання вітального листа для ${event.email}`,
);
// Тут може бути виклик поштового сервісу.
await Promise.resolve();
}
}Обробник потрібно додати до providers відповідного модуля:
import { Module } from '@nestjs/common';
import { SendWelcomeEmailHandler } from './send-welcome-email.handler';
@Module({
providers: [SendWelcomeEmailHandler],
})
export class NotificationsModule {}Якщо обробник не зареєстрований як провайдер, NestJS не зможе створити його екземпляр і підписати на подію.
Для однієї події можна мати кілька обробників:
@EventsHandler(UserRegisteredEvent)
export class WriteAuditLogHandler
implements IEventHandler<UserRegisteredEvent>
{
async handle(event: UserRegisteredEvent): Promise<void> {
// Запис події в журнал аудиту.
console.log(`Аудит: зареєстровано користувача ${event.userId}`);
}
}Обидва обробники отримають ту саму подію.
Для публікації використовується EventBus:
import { Controller, Post } from '@nestjs/common';
import { EventBus } from '@nestjs/cqrs';
import { UserRegisteredEvent } from './user-registered.event';
@Controller('users')
export class UsersController {
constructor(private readonly eventBus: EventBus) {}
@Post()
registerUser(): { status: string } {
const userId = crypto.randomUUID();
const email = 'user@example.com';
// Тут зазвичай спочатку зберігають користувача в базі даних.
this.eventBus.publish(
new UserRegisteredEvent(userId, email, new Date()),
);
return { status: 'registered' };
}
}publish() передає подію всім зареєстрованим обробникам.
Важлива особливість: публікація події не очікує завершення всіх асинхронних обробників. Навіть якщо handle() оголошений як async, виклик eventBus.publish() не повертає Promise, який можна очікувати через await.
Тому відповідь HTTP може бути відправлена раніше, ніж завершиться надсилання листа або запис аудиту.
Це підходить для другорядних побічних дій. Якщо результат операції є обов’язковим для поточного запиту, його не слід приховувати за асинхронною подією.
Нижче наведено мінімальний приклад застосунку з:
подією реєстрації користувача;
двома обробниками;
прямою підпискою через ofType;
налаштуванням обробки неперехоплених помилок.
user-registered.event.tsimport { IEvent } from '@nestjs/cqrs';
export class UserRegisteredEvent implements IEvent {
constructor(
public readonly userId: string,
public readonly email: string,
) {}
}send-welcome-email.handler.tsimport { Injectable, Logger } from '@nestjs/common';
import { EventsHandler, IEventHandler } from '@nestjs/cqrs';
import { UserRegisteredEvent } from './user-registered.event';
@Injectable()
@EventsHandler(UserRegisteredEvent)
export class SendWelcomeEmailHandler
implements IEventHandler<UserRegisteredEvent>
{
private readonly logger = new Logger(SendWelcomeEmailHandler.name);
async handle(event: UserRegisteredEvent): Promise<void> {
this.logger.log(`Вітальний лист: ${event.email}`);
// Імітація асинхронної операції.
await new Promise((resolve) => setTimeout(resolve, 100));
}
}audit-subscription.tsimport {
Injectable,
Logger,
OnModuleDestroy,
OnModuleInit,
} from '@nestjs/common';
import { EventBus } from '@nestjs/cqrs';
import { Subscription } from 'rxjs';
import { UserRegisteredEvent } from './user-registered.event';
@Injectable()
export class AuditSubscription
implements OnModuleInit, OnModuleDestroy
{
private readonly logger = new Logger(AuditSubscription.name);
private subscription?: Subscription;
constructor(private readonly eventBus: EventBus) {}
onModuleInit(): void {
this.subscription = this.eventBus
.ofType(UserRegisteredEvent)
.subscribe((event) => {
this.logger.log(
`Аудит: користувач ${event.userId} зареєстрований`,
);
});
}
onModuleDestroy(): void {
this.subscription?.unsubscribe();
}
}ofType() створює RxJS-потік, який отримує лише події вказаного класу.
Такий підхід корисний, коли підписка:
належить до довгоживучого сервісу;
має складну RxJS-логіку;
повинна бути створена або знищена разом із модулем.
Для звичайної бізнес-реакції переважно використовувати @EventsHandler, оскільки NestJS автоматично керує її реєстрацією.
users.controller.tsimport { Controller, Post } from '@nestjs/common';
import { EventBus } from '@nestjs/cqrs';
import { randomUUID } from 'node:crypto';
import { UserRegisteredEvent } from './user-registered.event';
@Controller('users')
export class UsersController {
constructor(private readonly eventBus: EventBus) {}
@Post()
register(): { userId: string } {
const userId = randomUUID();
const email = `user-${userId}@example.com`;
// У реальному застосунку тут перед публікацією зберігають користувача.
this.eventBus.publish(
new UserRegisteredEvent(userId, email),
);
return { userId };
}
}app.module.tsimport { Logger, Module } from '@nestjs/common';
import { CqrsModule } from '@nestjs/cqrs';
import { UsersController } from './users.controller';
import { SendWelcomeEmailHandler } from './send-welcome-email.handler';
import { AuditSubscription } from './audit-subscription';
@Module({
imports: [
CqrsModule.forRoot({
unhandledExceptionHandler: (error, event) => {
const logger = new Logger('EventBus');
logger.error(
`Помилка під час обробки події ${event.constructor.name}`,
error instanceof Error ? error.stack : String(error),
);
},
}),
],
controllers: [UsersController],
providers: [
SendWelcomeEmailHandler,
AuditSubscription,
],
})
export class AppModule {}main.tsimport { NestFactory } from '@nestjs/core';
import { AppModule } from './app.module';
async function bootstrap(): Promise<void> {
const app = await NestFactory.create(AppModule);
await app.listen(3000);
}
void bootstrap();Після запуску застосунку HTTP-запит:
curl -X POST http://localhost:3000/usersпризводить до такої послідовності:
Контролер створює UserRegisteredEvent.
EventBus публікує подію.
SendWelcomeEmailHandler виконує свою реакцію.
AuditSubscription отримує подію через RxJS-підписку.
Помилки в обробниках подій потрібно розглядати окремо від помилок основного HTTP-запиту.
Очікувані помилки краще обробляти безпосередньо в обробнику:
import { Injectable, Logger } from '@nestjs/common';
import { EventsHandler, IEventHandler } from '@nestjs/cqrs';
import { UserRegisteredEvent } from './user-registered.event';
@Injectable()
@EventsHandler(UserRegisteredEvent)
export class SendWelcomeEmailHandler
implements IEventHandler<UserRegisteredEvent>
{
private readonly logger = new Logger(SendWelcomeEmailHandler.name);
async handle(event: UserRegisteredEvent): Promise<void> {
try {
await this.sendEmail(event.email);
} catch (error) {
this.logger.error(
`Не вдалося надіслати лист для ${event.email}`,
error instanceof Error ? error.stack : String(error),
);
// Тут можна створити окрему подію для повторної спроби
// або передати помилку до системи моніторингу.
}
}
private async sendEmail(email: string): Promise<void> {
this.logger.log(`Надсилання листа на ${email}`);
await Promise.resolve();
}
}Такий варіант підходить, коли помилка не повинна переривати інші реакції на подію.
Якщо обробник не перехопив помилку, EventBus передає її налаштованому unhandledExceptionHandler.
CqrsModule.forRoot({
unhandledExceptionHandler: (error, event) => {
const eventName = event.constructor.name;
const message =
error instanceof Error ? error.message : String(error);
console.error(
`Помилка в обробнику події ${eventName}: ${message}`,
);
},
});Централізований обробник може:
записати помилку в журнал;
відправити її до системи моніторингу;
зберегти інформацію для повторної обробки;
додати контекст події.
Не слід автоматично повторно публікувати ту саму подію в цьому обробнику без обмеження кількості спроб. Інакше помилка може створити нескінченний цикл.
rethrowUnhandledЗа потреби можна налаштувати повторне викидання неперехоплених помилок:
CqrsModule.forRoot({
rethrowUnhandled: true,
unhandledExceptionHandler: (error, event) => {
console.error(
`Неперехоплена помилка в ${event.constructor.name}`,
error,
);
},
});Цей режим слід вмикати усвідомлено. Події часто обробляються після завершення основної операції, тому помилка обробника не завжди повинна завершувати процес застосунку.
Кожен обробник має виконувати одну логічну реакцію:
@EventsHandler(UserRegisteredEvent)
export class UpdateStatisticsHandler
implements IEventHandler<UserRegisteredEvent>
{
async handle(event: UserRegisteredEvent): Promise<void> {
// Оновлення статистики реєстрацій.
console.log(`Статистика оновлена для ${event.userId}`);
}
}Не варто створювати один великий обробник, який одночасно:
надсилає лист;
записує аудит;
оновлює статистику;
викликає зовнішній API.
Окремі обробники легше тестувати, змінювати та повторно використовувати.
Події часто розміщують у спільному каталозі контрактів, який можуть імпортувати різні модулі:
src/
events/
user-registered.event.ts
users/
users.controller.ts
notifications/
send-welcome-email.handler.ts
audit/
audit-subscription.tsМодуль, який створює подію, залежить лише від класу події та EventBus. Він не імпортує обробники інших модулів.
Обробники, навпаки, імпортують подію, але не повинні змушувати модуль-власник події знати про їхнє існування.
Це дозволяє додати нову реакцію без зміни коду реєстрації користувача.
Публікувати подію потрібно після успішного виконання основної операції.
Невдалий порядок:
this.eventBus.publish(new UserRegisteredEvent(userId, email));
// Якщо збереження завершиться помилкою,
// інші модулі вже вважатимуть користувача зареєстрованим.
await this.usersRepository.save(user);Кращий порядок:
const user = await this.usersRepository.save({
id: userId,
email,
});
this.eventBus.publish(
new UserRegisteredEvent(user.id, user.email),
);Однак навіть у цьому випадку між завершенням транзакції та публікацією події може виникнути помилка процесу. EventBus не надає гарантії доставки після перезапуску застосунку.
Якщо події мають бути надійно збережені та доставлені, потрібен окремий підхід із журналом подій або зовнішнім брокером повідомлень. Для внутрішньої комунікації модулів у межах одного процесу EventBus зазвичай достатній.
providersДекоратор @EventsHandler() сам по собі не створює провайдер:
@EventsHandler(UserRegisteredEvent)
export class SomeHandler {}Потрібно також додати клас до модуля:
@Module({
providers: [SomeHandler],
})
export class SomeModule {}Цей код не працюватиме як очікується:
interface UserRegistered {
userId: string;
}Для @EventsHandler() потрібен клас, доступний під час виконання:
export class UserRegisteredEvent implements IEvent {
constructor(public readonly userId: string) {}
}await eventBus.publish()publish() не призначений для очікування завершення всіх обробників:
await this.eventBus.publish(event);Це не зробить обробку події синхронною. Якщо результат обробки необхідний негайно, використовуйте звичайний виклик сервісу або інший механізм, призначений для запит-відповідь взаємодії.
Підписка через ofType() створює RxJS-підписку. Її потрібно знищувати разом із провайдером:
onModuleDestroy(): void {
this.subscription?.unsubscribe();
}Інакше при повторній ініціалізації компонента можна отримати дублювання обробки подій.
Не змінюйте дані події після її публікації:
event.email = 'another@example.com';Краще створювати події з незмінними властивостями та передавати в них лише необхідні значення.
Якщо обробник мовчки ігнорує помилку, система може втратити важливу побічну дію:
try {
await sendEmail();
} catch {
// Помилка повністю загублена.
}Помилку потрібно хоча б записати в журнал або передати до системи моніторингу.
EventBus забезпечує внутрішню комунікацію модулів у межах одного процесу.
Події описуються класами, які реалізують IEvent.
Обробники позначаються @EventsHandler() і додаються до providers.
Подія публікується через eventBus.publish().
Для однієї події можна зареєструвати кілька незалежних обробників.
RxJS-підписки можна створювати через eventBus.ofType(), але їх потрібно своєчасно знищувати.
publish() не очікує завершення асинхронних обробників.
Очікувані помилки обробляють у самому обробнику, а неперехоплені — через unhandledExceptionHandler.
EventBus не гарантує довготривале зберігання або доставку подій після перезапуску застосунку.