Kurs NestJS · Moduł 8: Cache i wydajność

Message Queues - system posłańców Imperium

5 min czytania
W tej lekcji6

Legionista kończy rejestrację, a serwis wysyła mu list powitalny. Serwer pocztowy akurat nie odpowiada - i rejestracja pada. Człowiek stoi przed bramą odrzucony, choć jego wpis do rejestru powiódł się bez zarzutu. Zawiodła rzecz uboczna, a pociągnęła za sobą rzecz główną.

Rzym rozwiązał to inaczej. Dowódca nie czeka, aż posłaniec wróci z Galii - zostawia list w skrzynce kurierskiej i wraca do swoich zajęć. Posłaniec zabierze go, gdy będzie mógł; jeśli koń okuleje, list poczeka. W systemach nazywamy tę skrzynkę kolejką wiadomości.

Co daje kolejka

Zaleta jest jedna i o nią właśnie chodzi: komunikacja asynchroniczna. Producent - ten, kto wysyła - nie czeka na konsumenta, czyli tego, kto odbiera. Zostawia wiadomość i idzie dalej.

Konsekwencje są trzy. Rejestracja kończy się natychmiast, bo nie jest już związana z pocztą. Awaria konsumenta nie przewraca producenta - wiadomości czekają w kolejce, aż wstanie. A gdy listów przybywa, dokładasz kolejnych konsumentów do tej samej kolejki, nie ruszając producenta.

Producent - zostawienie listu

Aplikacja wysyłająca potrzebuje klienta połączonego z brokerem, czyli z serwerem obsługującym kolejki:

1ClientsModule.register([
2  {
3    name: 'LEGION_SERVICE',
4    transport: Transport.RMQ,
5    options: {
6      urls: ['amqp://localhost:5672'],
7      queue: 'legion_queue',
8    },
9  },
10]);

Transport.RMQ wskazuje RabbitMQ - najczęściej używanego brokera. urls to adres serwera, a queue nazwa skrzynki, do której trafią listy. Nazwa LEGION_SERVICE posłuży do wstrzyknięcia klienta tam, gdzie będzie potrzebny.

emit kontra send - dwa rodzaje listów

Mając klienta, wysyłamy wiadomość jedną z dwóch metod. Różnica między nimi jest zasadnicza:

1// emit - wyślij i zapomnij, nie czekamy na odpowiedź
2this.client.emit('legion.created', { id: 1, name: 'Legio X' });
3
4// send - wyślij i poczekaj na odpowiedź
5const result = await this.client.send('legion.count', {}).toPromise();

emit() to fire-and-forget - zostawiasz list i idziesz dalej, nie oczekując odpowiedzi. Tak ogłaszasz fakty: legion powstał, płatność przeszła, użytkownik się zarejestrował. Nadawcy nie interesuje, co odbiorca z tym zrobi.

send() czeka na odpowiedź - to pytanie, nie ogłoszenie. Używasz go, gdy potrzebujesz wyniku: ile jest legionów, czy ten identyfikator istnieje.

Wybór między nimi decyduje o tym, czy w ogóle korzystasz z zalet kolejki. send() przywraca czekanie, a więc i sprzężenie, od którego uciekaliśmy - konsument znów musi działać, żeby producent ruszył dalej. Dlatego do powiadomień, e-maili i logów zawsze emit(); send() zostaw dla przypadków, gdy naprawdę potrzebujesz odpowiedzi.

Konsument - odebranie listu

Po drugiej stronie kolejki stoi serwis nasłuchujący:

1@Controller()
2export class LegionConsumer {
3  @EventPattern('legion.created')
4  async handleLegionCreated(@Payload() data: LegionCreatedDto) {
5    await this.mailService.sendWelcome(data);
6  }
7}

@EventPattern('legion.created') mówi: ta metoda obsługuje wiadomości o takiej nazwie. Nazwa jest umową między producentem a konsumentem - musi zgadzać się co do znaku z tą podaną w emit(). @Payload() wyciąga treść wiadomości, czyli obiekt przekazany jako drugi argument.

Zauważ, że konsument wygląda jak zwykły kontroler - bo działa podobnie. Różnica polega na tym, skąd przychodzi żądanie: nie z sieci HTTP, tylko z kolejki.

Potwierdzenie - kwit odbioru

Broker musi wiedzieć, czy list dotarł i został przetworzony. Bez tej informacji nie może go usunąć - a gdyby usunął od razu po wysłaniu, awaria konsumenta w połowie pracy oznaczałaby bezpowrotną stratę wiadomości.

Stąd potwierdzenie, czyli acknowledgement:

1@EventPattern('legion.created')
2async handleLegionCreated(@Payload() data: LegionCreatedDto, @Ctx() context: RmqContext) {
3  const channel = context.getChannelRef();
4  const originalMsg = context.getMessage();
5
6  try {
7    await this.mailService.sendWelcome(data);
8
9    channel.ack(originalMsg);
10  } catch (error) {
11    channel.nack(originalMsg, false, true);
12  }
13}

@Ctx() daje dostęp do kontekstu RabbitMQ. getChannelRef() zwraca kanał komunikacji z brokerem, a getMessage() - surową wiadomość, której dotyczy potwierdzenie.

Sedno jest w bloku try. channel.ack(originalMsg) potwierdza poprawne przetworzenie - dopiero wtedy broker usuwa wiadomość z kolejki. Zwróć uwagę na miejsce tego wywołania: po wykonaniu pracy, nie przed. Potwierdzenie na początku metody znaczyłoby "odebrałem", a nie "przetworzyłem", i awaria wysyłki e-maila skończyłaby się cichą utratą listu.

W catch stoi nack - odmowa potwierdzenia. Ostatni argument true każe brokerowi wstawić wiadomość z powrotem do kolejki, żeby spróbować później. To jest właśnie ta odporność, dla której kolejki wprowadziliśmy: nieudana próba nie niszczy wiadomości.

Podsumowanie

Skrzynka kurierska działa, listy dochodzą:

  • kolejka wiadomości daje komunikację asynchroniczną - producent nie czeka na konsumenta,
  • dzięki temu awaria konsumenta nie przewraca producenta, a wiadomości czekają, aż odbiorca wstanie,
  • klienta konfigurujemy przez ClientsModule.register z Transport.RMQ, adresem brokera i nazwą kolejki,
  • emit() to fire-and-forget - ogłaszasz fakt i nie czekasz na odpowiedź; send() czeka na odpowiedź, więc przywraca sprzężenie,
  • do powiadomień i e-maili używaj emit(); send() tylko wtedy, gdy naprawdę potrzebujesz wyniku,
  • konsument nasłuchuje przez @EventPattern('nazwa'), a treść wyciąga @Payload(); nazwa musi zgadzać się z tą z emit(),
  • channel.ack(originalMsg) potwierdza przetworzenie i dopiero wtedy broker usuwa wiadomość z kolejki,
  • potwierdzaj po wykonaniu pracy, nie przed - inaczej awaria oznacza cichą utratę wiadomości,
  • nack z ostatnim argumentem true wraca wiadomość do kolejki do ponownej próby,
  • pełny przepływ: producent wysyła przez emit → wiadomość czeka w kolejce → konsument odbiera przez @EventPattern → potwierdza przez channel.ack.

To ostatnia lekcja tego modułu. Umiesz już sprawdzić każdą część osobno, ich współpracę, całą drogę żądania, zmierzyć pokrycie i wytrzymałość, a teraz także rozdzielić serwisy tak, by awaria jednego nie pociągała reszty. A na razie zapamiętaj: kolejka to skrzynka kurierska - dowódca zostawia list i wraca do swoich spraw, a potwierdzenie odbioru przychodzi dopiero wtedy, gdy posłaniec naprawdę dotarł.

Kod do tej lekcji: src/message-queues.ts
1// Message Queues - System Poslancow Imperium Rzymskiego
2import { Injectable } from '@nestjs/common';
3
4// 1. Prosta implementacja kolejki wiadomosci
5class MessageQueue<T> {
6  private queue: Array<{ id: string; data: T; timestamp: Date }> = [];
7  private handlers: Map<string, (data: T) => Promise<void>> = new Map();
8
9  // Wyslij wiadomosc do kolejki
10  async publish(channel: string, data: T): Promise<string> {
11    const id = `msg-${Date.now()}`;
12    this.queue.push({ id, data, timestamp: new Date() });
13
14    const handler = this.handlers.get(channel);
15    if (handler) {
16      await handler(data);
17    }
18    return id;
19  }
20
21  // Zasubskrybuj kanal
22  subscribe(channel: string, handler: (data: T) => Promise<void>) {
23    this.handlers.set(channel, handler);
24  }
25
26  getQueueStatus() {
27    return {
28      size: this.queue.length,
29      channels: Array.from(this.handlers.keys()),
30      oldestMessage: this.queue[0]?.timestamp,
31    };
32  }
33}
34
35// 2. NestJS Bull Queue - konfiguracja
36// @Module({
37//   imports: [
38//     BullModule.forRoot({
39//       redis: { host: 'localhost', port: 6379 },
40//     }),
41//     BullModule.registerQueue({ name: 'tribute-collection' }),
42//   ],
43// })
44// export class QueueModule {}
45
46// 3. Producent - dodawanie zadan do kolejki
47@Injectable()
48class TributeQueueService {
49  private queue = new MessageQueue<any>();
50
51  constructor() {
52    this.queue.subscribe('tribute', async (data) => {
53      console.log(
54        `Collecting tribute: ${data.province} - ${data.amount} gold`
55      );
56    });
57
58    this.queue.subscribe('alert', async (data) => {
59      console.log(`ALERT: ${data.message}`);
60    });
61  }
62
63  async scheduleTributeCollection(province: string, amount: number) {
64    return this.queue.publish('tribute', { province, amount });
65  }
66
67  async sendAlert(message: string, priority: 'low' | 'high') {
68    return this.queue.publish('alert', { message, priority });
69  }
70
71  getStatus() {
72    return this.queue.getQueueStatus();
73  }
74}
75
76// 4. Wzorce komunikacji:
77// - Point-to-Point: producent -> konsument
78// - Pub/Sub: producent -> wielu konsumentow
79// - Request/Reply: producent czeka na odpowiedz
80// - Fan-out: wiadomosc do wielu kolejek
81
82// 5. Processor - przetwarzanie zadan
83// @Processor('tribute-collection')
84// class TributeProcessor {
85//   @Process('collect')
86//   async handleCollect(job: Job<TributeData>) {
87//     const { province, amount } = job.data;
88//     return { collected: true, province };
89//   }
90//
91//   @OnQueueCompleted()
92//   onCompleted(job: Job, result: any) {
93//     console.log(`Job ${job.id} completed`);
94//   }
95// }
96

Widzisz błąd w tej lekcji?

Sprawdź się

Odpowiedz na pytania z tej lekcji. Wybierz odpowiedź, a od razu zobaczysz, czy jest poprawna.

  1. 1. Middleware compression w NestJS służy do:

  2. 2. Dekorator @Exclude() z class-transformer powoduje:

Zadania praktyczne w grze

  • Edytor kodu

    Zaimplementuj LegionResponseDto z @Expose() na polach id, name, province i @Exclude() na polach secretCode i passwordHash. Dodaj @Transform() na polu soldiers zaokrąglając wartość.

  • Układanie w pionie

    Ułóż elementy odpowiedzi paginowanej od danych do metadanych:

  • Układanie w pionie

    Ułóż nagłówki cache HTTP od najczęściej używanego do najbardziej zaawansowanego:

  • Edytor kodu

    Stwórz @Processor() i @Process() dla kolejki wysyłania emaili z obsługą błędów

  • Układanie w poziomie

    Ułóż poprawną rejestrację Bull queue w module NestJS:

Przydatne artykuły