Node.js Backend Development Bootcamp · Lekcja

Producenci, konsumenci i model AMQP

Łącz się z brokerem i przesyłaj komunikaty przez exchange, kolejki i wiązania za pomocą protokołu AMQP

Lekcja 1 z 413 kroki

Producenci, konsumenci i model AMQP to bezpłatna lekcja Node.js Backend Development Bootcamp na CoddyKit. To lekcja 1 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 Node.js Backend Development Bootcamp, a Twój postęp synchronizuje się między webem a aplikacją CoddyKit. Kurs Node.js Backend Development Bootcamp zawiera 4 lekcji w sumie.

Po co są kolejki komunikatów?

W backendzie Node.js bezpośrednie wywołanie innej usługi ściśle łączy obie usługi: jeśli odbiorca działa wolno lub jest niedostępny, wywołujący czeka albo kończy się błędem. Kolejka komunikatów rozdziela te usługi. Nadawca umieszcza komunikat w brokerze i kontynuuje działanie, a odbiorca pobiera go, gdy jest gotowy.

  • Asynchroniczna — producent nigdy nie czeka na konsumenta.
  • Odporna — komunikaty pozostają dostępne, gdy konsument jest offline.
  • Skalowalna — można dodać więcej konsumentów, aby szybciej opróżniać zaległości.

RabbitMQ to popularny broker obsługujący AMQP 0-9-1, czyli protokół używany w całej tej lekcji.

Model AMQP

AMQP rozdziela czynność publikowania od czynności przechowywania. Producent nigdy nie zapisuje komunikatu bezpośrednio w kolejce. Zamiast tego publikuje go w exchange, a exchange kieruje komunikat do jednej lub większej liczby kolejek na podstawie powiązań.

  • Producent — publikuje komunikaty.
  • Exchange — odbiera komunikaty i decyduje, dokąd je skierować.
  • Powiązanie — reguła łącząca exchange z kolejką, często za pomocą klucza routingu.
  • Kolejka — buforuje komunikaty do czasu ich odczytania przez konsumenta.
  • Konsument — subskrybuje kolejkę i przetwarza komunikaty.

Ta warstwa pośrednia zapewnia elastyczne routowanie: można zmieniać powiązania, a nie kod producenta.

Łączenie z poziomu Node.js

Biblioteka amqplib to standardowy klient AMQP dla Node.js. Najpierw otwierasz połączenie TCP z brokerem, a następnie tworzysz na nim kanał. Kanał jest lekkim wirtualnym połączeniem, za pośrednictwem którego wykonywana jest niemal cała komunikacja AMQP, dlatego nie trzeba otwierać nowego gniazda TCP dla każdego zadania.

Łańcuchy połączeń używają schematu amqp://: amqp://user:pass@host:5672/vhost. Zawsze używaj await podczas nawiązywania połączenia i obsługuj błędy.

const amqp = require('amqplib');

async function connect() {
  const url = process.env.AMQP_URL || 'amqp://guest:guest@localhost:5672';
  const connection = await amqp.connect(url);
  const channel = await connection.createChannel();
  console.log('Connected and channel opened');
  return { connection, channel };
}

connect().catch((err) => {
  console.error('AMQP connection failed:', err.message);
  process.exit(1);
});

Deklarowanie kolejki

Przed wysłaniem komunikatu obie strony zazwyczaj muszą zadeklarować kolejkę. Deklarowanie jest idempotentne: tworzy kolejkę, jeśli jej brakuje, a w przeciwnym razie sprawdza zgodność ustawień. Najważniejszą opcją jest durable.

  • durable: true — definicja kolejki przetrwa ponowne uruchomienie brokera.
  • durable: false — kolejka zostanie utracona po ponownym uruchomieniu.

Trwałość kolejki jest niezależna od trwałości znajdujących się w niej komunikatów — komunikaty trwałe omówimy wkrótce.

async function setupQueue(channel) {
  const queue = 'tasks';
  await channel.assertQueue(queue, { durable: true });
  console.log(`Queue "${queue}" is ready`);
  return queue;
}

Domyślny exchange

RabbitMQ zawiera nienazwany domyślny exchange (pusty ciąg ''). Obowiązuje w nim specjalna reguła: automatycznie wiąże każdą kolejkę, używając jej własnej nazwy jako klucza routingu. Dlatego channel.sendToQueue('tasks', ...) w rzeczywistości publikuje komunikat w domyślnym exchange z kluczem routingu 'tasks'.

To najprostszy sposób na rozpoczęcie pracy, ale ukrywa koncepcję exchange. Rzeczywiste aplikacje deklarują własne exchange na potrzeby jawnego routingu, co zrobimy później.

Producent

Producent publikuje komunikat, a następnie może zamknąć połączenie. Treść komunikatów AMQP składa się z surowych bajtów, dlatego serializujemy dane do JSON i opakowujemy je w Buffer. Opcja persistent: true oznacza komunikat tak, aby mógł zostać zapisany na dysku w trwałej kolejce.

Zwróć uwagę na krótkie opóźnienie przed zamknięciem: sendToQueue buforuje dane lokalnie, więc przed zakończeniem działania czekamy na opróżnienie bufora kanału.

const amqp = require('amqplib');

async function publish() {
  const conn = await amqp.connect('amqp://localhost');
  const channel = await conn.createChannel();
  const queue = 'tasks';
  await channel.assertQueue(queue, { durable: true });

  const job = { id: 42, type: 'resize-image', file: 'cat.png' };
  channel.sendToQueue(queue, Buffer.from(JSON.stringify(job)), {
    persistent: true,
  });
  console.log('Sent job', job.id);

  await channel.close();
  await conn.close();
}

publish().catch(console.error);

Konsument

Konsument subskrybuje kolejkę za pomocą channel.consume. Broker przekazuje komunikaty do funkcji callback w miarę ich napływania, a konsument zazwyczaj pozostaje uruchomiony. Komunikat jest dostępny jako msg.content, czyli jako Buffer, który należy zdekodować z powrotem do obiektu.

Domyślnie consume automatycznie potwierdza odbiór komunikatów. W następnej scenie wyłączymy tę opcję, aby samodzielnie kontrolować moment uznania komunikatu za przetworzony.

const amqp = require('amqplib');

async function consume() {
  const conn = await amqp.connect('amqp://localhost');
  const channel = await conn.createChannel();
  const queue = 'tasks';
  await channel.assertQueue(queue, { durable: true });

  console.log('Waiting for messages...');
  await channel.consume(queue, (msg) => {
    if (!msg) return;
    const job = JSON.parse(msg.content.toString());
    console.log('Processing job', job.id, job.type);
  });
}

consume().catch(console.error);

Potwierdzenia odbioru

Potwierdzenie (ack) informuje brokera, że komunikat został pomyślnie obsłużony i można go usunąć. Ustaw { noAck: false } i wywołaj channel.ack(msg) dopiero po zakończeniu pracy.

  • Jeśli konsument zakończy działanie przed wysłaniem potwierdzenia, RabbitMQ ponownie umieści komunikat w kolejce, aby mógł go odebrać inny konsument.
  • channel.nack(msg, false, true) odrzuca komunikat i ponownie umieszcza go w kolejce, natomiast nack(msg, false, false) go odrzuca lub kieruje do dead-letter queue.

Ręczne potwierdzenia są podstawą dostarczania co najmniej raz — nigdy nie potwierdzaj komunikatu przed zakończeniem pracy.

await channel.consume(
  queue,
  async (msg) => {
    if (!msg) return;
    try {
      const job = JSON.parse(msg.content.toString());
      await processJob(job); // your real work
      channel.ack(msg); // success
    } catch (err) {
      channel.nack(msg, false, false); // discard / dead-letter
    }
  },
  { noAck: false }
);

Sprawiedliwe rozdzielanie za pomocą prefetch

Domyślnie RabbitMQ rozdziela komunikaty cyklicznie między konsumentów, nie uwzględniając obciążenia każdego z nich. Konsument zajęty wolnym zadaniem może gromadzić zaległości, podczas gdy inni pozostają bezczynni.

channel.prefetch(1) rozwiązuje ten problem: broker nie wyśle konsumentowi nowego komunikatu, dopóki nie otrzyma potwierdzenia poprzedniego. Zapewnia to sprawiedliwe rozdzielanie — praca trafia do konsumenta, który rzeczywiście jest wolny.

async function worker() {
  const conn = await amqp.connect('amqp://localhost');
  const channel = await conn.createChannel();
  await channel.assertQueue('tasks', { durable: true });

  // Only one unacked message at a time per consumer
  await channel.prefetch(1);

  await channel.consume('tasks', async (msg) => {
    await handle(JSON.parse(msg.content.toString()));
    channel.ack(msg);
  }, { noAck: false });
}

Niestandardowe exchange i powiązania

W rzeczywistym routingu deklaruje się własny exchange i wiąże z nim kolejki. Typy exchange określają sposób routingu:

  • direct — dokładne dopasowanie klucza routingu.
  • fanout — rozgłaszanie do każdej powiązanej kolejki z pominięciem klucza.
  • topic — dopasowanie do wzorca z symbolami wieloznacznymi, np. order.*.
  • headers — dopasowanie na podstawie nagłówków komunikatu.

Poniżej exchange typu direct o nazwie orders kieruje komunikaty order.created do kolejki. Producent publikuje je za pomocą channel.publish(exchange, routingKey, content).

async function setupRouting(channel) {
  const exchange = 'orders';
  await channel.assertExchange(exchange, 'direct', { durable: true });

  const queue = 'order_processing';
  await channel.assertQueue(queue, { durable: true });
  await channel.bindQueue(queue, exchange, 'order.created');

  const event = { orderId: 1001, total: 79.9 };
  channel.publish(
    exchange,
    'order.created',
    Buffer.from(JSON.stringify(event)),
    { persistent: true }
  );
}

Połączenie wszystkiego: działające demo

Oto samodzielna symulacja przepływu AMQP, która nie wymaga brokera — modeluje routing z exchange do kolejki oraz pobieranie wiadomości przez konsumenta. Ilustruje ona podstawowy model: producent do exchange, exchange do kolejki za pośrednictwem bindingu, kolejka do konsumenta z potwierdzeniem.

Uruchom ją, aby zobaczyć, jak routing key wybiera docelową kolejkę oraz jak każda wiadomość jest potwierdzana dokładnie raz.

// In-memory model of the AMQP routing flow (no broker needed)
class Broker {
  constructor() { this.queues = {}; this.bindings = {}; }
  assertQueue(q) { this.queues[q] = this.queues[q] || []; }
  bind(exchange, queue, key) {
    (this.bindings[exchange] ||= []).push({ queue, key });
  }
  publish(exchange, routingKey, body) {
    for (const b of this.bindings[exchange] || []) {
      if (b.key === routingKey) this.queues[b.queue].push(body);
    }
  }
  consume(queue, handler) {
    let msg;
    while ((msg = this.queues[queue].shift())) handler(msg, () => {});
  }
}

const broker = new Broker();
broker.assertQueue('order_processing');
broker.bind('orders', 'order_processing', 'order.created');

broker.publish('orders', 'order.created', { orderId: 1001 });
broker.publish('orders', 'order.deleted', { orderId: 1002 }); // no binding

broker.consume('order_processing', (msg, ack) => {
  console.log('Consumed:', msg);
  ack();
});

Szybki test

Uruchamia Pan/Pani kilku konsumentów-workerów na jednej trwałej kolejce. Niektóre zadania trwają znacznie dłużej niż inne i zauważa Pan/Pani, że szybcy workerzy pozostają bezczynni, podczas gdy kilku innych jest przeciążonych. Która pojedyncza zmiana najlepiej równoważy obciążenie?

Podsumowanie

Połączył(a) się Pan/Pani z brokerem i przesyłał(a) wiadomości zgodnie z modelem AMQP.

  • Połączenie a kanał — jedno połączenie TCP i wiele lekkich kanałów.
  • Producent do exchange, a następnie do kolejki — producenci nigdy nie zapisują bezpośrednio do kolejek; exchange kierują wiadomości za pomocą bindingów i routing keys.
  • Typy exchange — direct, fanout, topic i headers określają logikę routingu.
  • Trwałość i persystencja — trwałe kolejki oraz persystentne wiadomości przetrwają ponowne uruchomienie.
  • Potwierdzenia — ręczne ack/nack po wykonaniu pracy zapewniają dostarczenie co najmniej raz; niepotwierdzone wiadomości są ponownie umieszczane w kolejce.
  • Prefetch — prefetch(1) umożliwia sprawiedliwe rozdzielanie pracy między konkurujących konsumentów.

Następnie może Pan/Pani poznać exchange typu dead-letter, TTL wiadomości oraz potwierdzenia publikacji zapewniające niezawodne dostarczanie.

Bezpłatny start

Ucz się JavaScript 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
22
Lekcje
92

Często zadawane pytania

Czy lekcja „Producenci, konsumenci i model AMQP” jest bezpłatna?

Tak — pełny tekst „Producenci, konsumenci i model AMQP” 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 Node.js Backend Development Bootcamp, przejdź na CoddyKit PRO. Kurs Node.js Backend Development Bootcamp zawiera 4 lekcji w sumie.

Co nauczysz się w „Producenci, konsumenci i model AMQP”?

Łącz się z brokerem i przesyłaj komunikaty przez exchange, kolejki i wiązania za pomocą protokołu AMQP Ćwiczysz Node.js Backend Development Bootcamp 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ąć Node.js Backend Development Bootcamp?

Nie wymagamy żadnego doświadczenia. Node.js Backend Development Bootcamp 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 1 z 4.

Ile czasu zajmuje lekcja „Producenci, konsumenci i model AMQP”?

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 Node.js Backend Development Bootcamp?

Tak. Każda lekcja Node.js Backend Development Bootcamp 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. Producenci, konsumenci i model AMQP
  2. Typy exchange: Direct, Topic, Fanout i Headers
  3. Potwierdzenia, kolejki wiadomości odrzuconych i ponawianie prób
  4. Kolejki zadań, prefetch i konkurujący konsumenci
← Powrót do Node.js Backend Development Bootcamp