Kurs NestJS · Moduł 8: Cache i wydajność
Message Queues - system posłańców Imperium
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.registerzTransport.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ą zemit(), 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,
nackz ostatnim argumentemtruewraca 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 przezchannel.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// }
96Widzisz błąd w tej lekcji?
Sprawdź się
Odpowiedz na pytania z tej lekcji. Wybierz odpowiedź, a od razu zobaczysz, czy jest poprawna.
1. Middleware compression w NestJS służy do:
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: