Bootcamp backendontwikkeling met Node.js · Les

Producers, consumers en het AMQP-model

Maak verbinding met een broker en verplaats berichten via exchanges, queues en bindings met het AMQP-protocol.

Les 1 van 413 stappen

Producers, consumers en het AMQP-model is een gratis Bootcamp backendontwikkeling met Node.js-les op CoddyKit. Dit is les 1 van 4. Je kunt de volledige les hieronder gratis lezen en daarna in de browser praktisch oefenen met een ingebouwde code-editor en een AI-begeleider die 24/7 beschikbaar is. Deze les maakt deel uit van het leertraject Bootcamp backendontwikkeling met Node.js. Je voortgang wordt gesynchroniseerd op het web en in de CoddyKit-app. De cursus Bootcamp backendontwikkeling met Node.js bevat in totaal 4 lessen.

Waarom berichtenwachtrijen?

In een Node.js-backend koppel je twee services aan elkaar wanneer je de andere rechtstreeks aanroept: als de ontvanger traag of offline is, blokkeert de aanroeper of mislukt de aanroep. Een berichtenwachtrij ontkoppelt ze. De verzender plaatst een bericht in een broker en gaat verder; de ontvanger haalt het op zodra die er klaar voor is.

  • Asynchroon — de producent wacht nooit op de consument.
  • Veerkrachtig — berichten blijven bewaard terwijl een consument offline is.
  • Schaalbaar — voeg meer consumenten toe om een achterstand sneller weg te werken.

RabbitMQ is een populaire broker die AMQP 0-9-1 spreekt, het protocol dat we in deze hele les gebruiken.

Het AMQP-model

AMQP scheidt het publiceren van het opslaan. Een producent schrijft nooit rechtstreeks naar een wachtrij. In plaats daarvan publiceert die naar een exchange, waarna de exchange het bericht op basis van koppelingen naar een of meer wachtrijen routeert.

  • Producent — publiceert berichten.
  • Exchange — ontvangt berichten en bepaalt waar ze naartoe gaan.
  • Koppeling — een regel die een exchange aan een wachtrij koppelt, vaak met een routeringssleutel.
  • Wachtrij — buffert berichten totdat een consument ze leest.
  • Consument — abonneert zich op een wachtrij en verwerkt berichten.

Dankzij deze tussenlaag is routering flexibel: wijzig de koppelingen, niet de code van de producent.

Verbinding maken vanuit Node.js

De bibliotheek amqplib is de standaard AMQP-client voor Node.js. Je opent een TCP-verbinding met de broker en maakt daarover vervolgens een kanaal. Een kanaal is een lichtgewicht virtuele verbinding waarop vrijwel alle AMQP-bewerkingen plaatsvinden, zodat je niet voor elke taak een nieuwe TCP-socket hoeft te openen.

Verbindingsreeksen gebruiken het schema amqp://: amqp://user:pass@host:5672/vhost. Gebruik altijd await voor de verbinding en handel fouten af.

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);
});

Een wachtrij declareren

Voordat je berichten verzendt, declareren beide kanten doorgaans de wachtrij. Declareren is idempotent: de wachtrij wordt aangemaakt als die ontbreekt; anders wordt gecontroleerd of de instellingen overeenkomen. De belangrijkste optie is durable.

  • durable: true — de definitie van de wachtrij blijft bestaan na een herstart van de broker.
  • durable: false — de wachtrij gaat verloren bij een herstart.

De duurzaamheid van de wachtrij staat los van de duurzaamheid van de berichten die erin staan — permanente berichten behandelen we binnenkort.

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

De standaard-exchange

RabbitMQ wordt geleverd met een naamloze standaard-exchange (de lege tekenreeks ''). Deze heeft een speciale regel: elke wachtrij wordt automatisch gekoppeld met de eigen naam van de wachtrij als routeringssleutel. channel.sendToQueue('tasks', ...) publiceert dus in werkelijkheid naar de standaard-exchange met de routeringssleutel 'tasks'.

Dit is de eenvoudigste manier om te beginnen, maar het exchange-concept blijft verborgen. Echte toepassingen declareren hun eigen exchanges voor expliciete routering; dat doen we later.

Een producent

Een producent publiceert een bericht en kan daarna de verbinding sluiten. De berichtinhoud van AMQP bestaat uit onbewerkte bytes, dus serialiseren we deze naar JSON en verpakken we die in een Buffer. De optie persistent: true markeert het bericht, zodat het naar schijf kan worden geschreven in een duurzame wachtrij.

Let op de korte vertraging vóór het sluiten: sendToQueue buffert lokaal, dus wachten we tot het kanaal zijn buffer heeft leeggemaakt voordat we afsluiten.

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);

Een consument

Een consument abonneert zich met channel.consume op een wachtrij. De broker stuurt berichten naar de callback zodra ze binnenkomen; de consument blijft doorgaans actief. Het bericht komt binnen als msg.content, een Buffer die je terug decodeert naar je object.

Standaard bevestigt consume automatisch. In de volgende scène schakelen we dat uit, zodat je precies bepaalt wanneer een bericht als voltooid wordt beschouwd.

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);

Bevestigingen

Een bevestiging (ack) vertelt de broker dat een bericht succesvol is verwerkt en verwijderd kan worden. Stel { noAck: false } in en roep channel.ack(msg) pas aan nadat je werk is voltooid.

  • Als de consument crasht voordat die het bericht bevestigt, plaatst RabbitMQ het bericht terug in de wachtrij voor een andere consument.
  • channel.nack(msg, false, true) weigert het bericht en plaatst het terug in de wachtrij; nack(msg, false, false) verwijdert het of stuurt het naar een dead-letter-wachtrij.

Handmatige bevestigingen vormen de basis van levering van minstens één keer — bevestig nooit voordat het werk is voltooid.

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 }
);

Eerlijke verdeling met prefetch

Standaard verdeelt RabbitMQ berichten beurtelings over consumenten, zonder rekening te houden met hoe druk elke consument is. Een consument die vastzit aan een trage taak kan een achterstand opbouwen terwijl andere consumenten niets doen.

channel.prefetch(1) lost dit op: de broker stuurt pas een nieuw bericht naar een consument nadat die het vorige heeft bevestigd. Dit zorgt voor een eerlijke verdeling — werk gaat naar de consument die daadwerkelijk beschikbaar is.

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 });
}

Aangepaste exchanges en koppelingen

Voor echte routering declareer je je eigen exchange en koppel je daar wachtrijen aan. Het type exchange bepaalt de routeringslogica:

  • direct — exacte overeenkomst met de routeringssleutel.
  • fanout — verzendt naar elke gekoppelde wachtrij en negeert de sleutel.
  • topic — overeenkomst met een jokertekenpatroon, bijvoorbeeld order.*.
  • headers — overeenkomst op basis van berichtheaders.

Hieronder routeert een directe exchange met de naam orders berichten met order.created naar een wachtrij. De producent publiceert met 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 }
  );
}

Alles samenbrengen: een uitvoerbare demo

Hier zie je een zelfstandige simulatie van de AMQP-stroom waarvoor geen broker nodig is — de simulatie modelleert een exchange die berichten naar een queue routeert en een consumer die deze queue leegmaakt. Dit illustreert het denkmodel: producer naar exchange, exchange naar queue via een binding, queue naar consumer met een bevestiging.

Voer de simulatie uit om te zien hoe de routing key de bestemmingsqueue selecteert en hoe elk bericht precies één keer wordt bevestigd.

// 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();
});

Korte controle

Je voert meerdere worker-consumers uit op één duurzame queue. Sommige taken duren veel langer dan andere, en je merkt dat snelle workers niets doen terwijl enkele workers overbelast zijn. Welke enkele wijziging zorgt het best voor een evenwichtige verdeling van de werklast?

Samenvatting

Je hebt verbinding gemaakt met een broker en berichten door het AMQP-model gestuurd.

  • Verbinding versus kanaal — één TCP-verbinding, veel lichtgewicht kanalen.
  • Producer naar exchange naar queue — producers schrijven nooit rechtstreeks naar queues; exchanges routeren berichten via bindings en routing keys.
  • Typen exchanges — direct, fanout, topic en headers bepalen de routeringslogica.
  • Duurzaamheid en persistentie — duurzame queues en persistente berichten blijven na herstarts behouden.
  • Bevestigingen — handmatig ack/nack nadat het werk is uitgevoerd, zorgt voor aflevering van ten minste één keer; niet-bevestigde berichten worden opnieuw in de queue geplaatst.
  • Prefetch — prefetch(1) maakt een eerlijke verdeling mogelijk tussen concurrerende consumers.

Daarna kun je dead-letter exchanges, TTL's voor berichten en publisher confirms verkennen voor betrouwbare aflevering.

Gratis beginnen

Leer JavaScript met een AI-tutor — gratis

Schrijf echte code en voer die uit in je browser, krijg direct hulp van een AI-tutor die 24/7 beschikbaar is en ga verder waar je gebleven bent op het web of in de app.

Cursussen
22
Lessen
92

Veelgestelde vragen

Is de les “Producers, consumers en het AMQP-model” gratis?

Ja — de volledige tekst van “Producers, consumers en het AMQP-model” kun je hier gratis op het web lezen. Als je interactief wilt oefenen met een ingebouwde code-editor en een AI-begeleider die 24/7 beschikbaar is, en de rest van de cursus Bootcamp backendontwikkeling met Node.js wilt ontgrendelen, kun je upgraden naar CoddyKit PRO. De cursus Bootcamp backendontwikkeling met Node.js bevat in totaal 4 lessen.

Wat leer ik in “Producers, consumers en het AMQP-model”?

Maak verbinding met een broker en verplaats berichten via exchanges, queues en bindings met het AMQP-protocol. Je oefent met Bootcamp backendontwikkeling met Node.js door code rechtstreeks in de browser uit te voeren. Een AI-begeleider die 24/7 beschikbaar is beantwoordt je vragen terwijl je de les doorwerkt.

Heb ik ervaring nodig om met Bootcamp backendontwikkeling met Node.js te beginnen?

Ervaring vooraf is niet nodig. Bootcamp backendontwikkeling met Node.js op CoddyKit is opgebouwd voor beginners tot gevorderden, zodat je hier of bij het begin kunt starten en in je eigen tempo kunt leren. Dit is les 1 van 4.

Hoe lang duurt de les “Producers, consumers en het AMQP-model”?

De meeste lessen van CoddyKit duren ongeveer 5–10 minuten. Elke les is kort en interactief, zodat je gestaag vooruitgaat en op het web en in de app precies verdergaat waar je was gebleven.

Kan ik code schrijven en uitvoeren in deze les over Bootcamp backendontwikkeling met Node.js?

Ja. Elke les over Bootcamp backendontwikkeling met Node.js bevat een ingebouwde code-editor, zodat je rechtstreeks in je browser echte code kunt schrijven en uitvoeren en direct feedback van AI krijgt — lokale installatie is niet nodig.

Alle lessen in deze cursus

  1. Producers, consumers en het AMQP-model
  2. Exchangetypen: direct, topic, fanout en headers
  3. Acknowledgements, dead-letter-queues en retries
  4. Work queues, prefetch en concurrerende consumers
← Terug naar Bootcamp backendontwikkeling met Node.js