Zum Inhalt springen

Ein event-driven Microservice mit Node.js bauen

Schritt-für-Schritt zum event-driven Microservice mit Node.js, RabbitMQ und TypeScript: Publishing, Consumption, Dead Letter Queues, Idempotenz.

5 Min. Lesezeit
Event-driven Microservice-Architektur mit Producer-, Message-Broker- und Consumer-Services

Event-driven Microservices kommunizieren über Events statt über direkte API-Aufrufe. Service A publiziert ein Event ("Bestellung wurde aufgegeben"), und jeder Service, der an diesem Event interessiert ist, verarbeitet es unabhängig. Das entkoppelt die Services — der Order-Service muss den Inventory-Service, den Notification-Service oder andere Consumer nicht kennen. Neue Consumer können hinzugefügt werden, ohne den Producer zu ändern.

Wir bauen ein praxisnahes event-driven System mit Node.js, TypeScript und RabbitMQ. Das System verarbeitet Bestell-Events: Wenn eine Bestellung aufgegeben wird, aktualisieren separate Consumer den Lagerbestand, versenden Benachrichtigungs-E-Mails und schreiben Analytics-Logs.

RabbitMQ einrichten

RabbitMQ ist ein Message Broker, der Events über Exchanges und Queues von Producern zu Consumern routet.

ymlyaml
# docker-compose.yml
services:
  rabbitmq:
    image: rabbitmq:3-management
    ports:
      - "5672:5672"    # AMQP protocol
      - "15672:15672"  # Management UI
    environment:
      RABBITMQ_DEFAULT_USER: guest
      RABBITMQ_DEFAULT_PASS: guest
    volumes:
      - rabbitmq_data:/var/lib/rabbitmq
 
volumes:
  rabbitmq_data:
tstypescript
// lib/rabbitmq.ts — connection wrapper
import amqp, { type Connection, type Channel } from 'amqplib';
 
let connection: Connection | null = null;
let channel: Channel | null = null;
 
export async function getChannel(): Promise<Channel> {
  if (channel) return channel;
 
  connection = await amqp.connect(
    process.env.RABBITMQ_URL ?? 'amqp://guest:guest@localhost:5672'
  );
 
  connection.on('error', (err) => {
    console.error('RabbitMQ connection error:', err);
    channel = null;
    connection = null;
  });
 
  channel = await connection.createChannel();
 
  // Prefetch: process one message at a time per consumer
  await channel.prefetch(1);
 
  return channel;
}
 
export async function closeConnection(): Promise<void> {
  if (channel) await channel.close();
  if (connection) await connection.close();
  channel = null;
  connection = null;
}

Events definieren

Events sollten in sich abgeschlossen sein — jeder Consumer sollte das Event verarbeiten können, ohne zusätzliche API-Aufrufe machen zu müssen.

tstypescript
// events/types.ts
interface BaseEvent {
  eventId: string;       // Unique ID for idempotency
  eventType: string;     // Event name
  timestamp: string;     // ISO 8601
  version: number;       // Schema version
  source: string;        // Which service published this
}
 
interface OrderPlacedEvent extends BaseEvent {
  eventType: 'order.placed';
  data: {
    orderId: string;
    customerId: string;
    customerEmail: string;
    items: {
      productId: string;
      quantity: number;
      unitPrice: number;
    }[];
    total: number;
    currency: string;
  };
}
 
interface OrderCancelledEvent extends BaseEvent {
  eventType: 'order.cancelled';
  data: {
    orderId: string;
    customerId: string;
    reason: string;
  };
}
 
type OrderEvent = OrderPlacedEvent | OrderCancelledEvent;
tstypescript
// ❌ Anemic events — consumers need to call back to the producer
interface BadOrderEvent {
  orderId: string;  // Consumer must call GET /orders/123 for details
}
// This creates coupling: the consumer depends on the producer's API
 
// ✅ Rich events — self-contained with all necessary data
interface GoodOrderEvent extends BaseEvent {
  eventType: 'order.placed';
  data: {
    orderId: string;
    customerId: string;
    customerEmail: string;
    items: { productId: string; quantity: number; unitPrice: number }[];
    total: number;
  };
}
// Consumer has everything it needs — no callbacks required

Events publizieren

Der Producer publiziert Events an einen RabbitMQ-Exchange. Ein Topic-Exchange erlaubt es Consumern, sich auf bestimmte Event-Muster zu abonnieren.

tstypescript
// events/publisher.ts
import { getChannel } from '../lib/rabbitmq';
import crypto from 'crypto';
 
const EXCHANGE_NAME = 'order_events';
 
export async function setupPublisher(): Promise<void> {
  const channel = await getChannel();
 
  // Topic exchange: routes messages based on routing key pattern
  await channel.assertExchange(EXCHANGE_NAME, 'topic', {
    durable: true,  // Survives broker restart
  });
}
 
export async function publishEvent(
  event: OrderEvent
): Promise<void> {
  const channel = await getChannel();
 
  const message = Buffer.from(JSON.stringify(event));
  const routingKey = event.eventType; // e.g., "order.placed"
 
  channel.publish(EXCHANGE_NAME, routingKey, message, {
    persistent: true,     // Survives broker restart
    contentType: 'application/json',
    messageId: event.eventId,
    timestamp: Date.now(),
  });
 
  console.log(`Published ${event.eventType}: ${event.eventId}`);
}
 
// Usage in the order service
async function placeOrder(order: Order): Promise<void> {
  // Save order to database first
  await db.query(
    'INSERT INTO orders (id, customer_id, total) VALUES ($1, $2, $3)',
    [order.id, order.customerId, order.total]
  );
 
  // Then publish the event
  await publishEvent({
    eventId: crypto.randomUUID(),
    eventType: 'order.placed',
    timestamp: new Date().toISOString(),
    version: 1,
    source: 'order-service',
    data: {
      orderId: order.id,
      customerId: order.customerId,
      customerEmail: order.customerEmail,
      items: order.items,
      total: order.total,
      currency: 'USD',
    },
  });
}

Events konsumieren

Jeder Consumer bindet seine eigene Queue mit einem Routing-Key-Muster an den Exchange. Mehrere Consumer können dasselbe Event unabhängig voneinander verarbeiten.

tstypescript
// consumers/inventory-consumer.ts
import { getChannel } from '../lib/rabbitmq';
 
const EXCHANGE_NAME = 'order_events';
const QUEUE_NAME = 'inventory_order_events';
 
export async function startInventoryConsumer(): Promise<void> {
  const channel = await getChannel();
 
  // Create a durable queue for this consumer
  await channel.assertQueue(QUEUE_NAME, {
    durable: true,
    deadLetterExchange: 'dlx_order_events',  // Failed messages go here
  });
 
  // Bind queue to exchange with routing key pattern
  await channel.bindQueue(QUEUE_NAME, EXCHANGE_NAME, 'order.*');
 
  console.log('Inventory consumer listening for order events...');
 
  channel.consume(QUEUE_NAME, async (msg) => {
    if (!msg) return;
 
    try {
      const event: OrderEvent = JSON.parse(msg.content.toString());
 
      switch (event.eventType) {
        case 'order.placed':
          await handleOrderPlaced(event);
          break;
        case 'order.cancelled':
          await handleOrderCancelled(event);
          break;
        default:
          console.warn(`Unknown event type: ${event.eventType}`);
      }
 
      // Acknowledge: message processed successfully
      channel.ack(msg);
    } catch (error) {
      console.error('Failed to process message:', error);
      // Reject and send to dead letter queue
      channel.nack(msg, false, false);
    }
  });
}
 
async function handleOrderPlaced(event: OrderPlacedEvent): Promise<void> {
  for (const item of event.data.items) {
    await db.query(
      'UPDATE products SET stock = stock - $1 WHERE id = $2 AND stock >= $1',
      [item.quantity, item.productId]
    );
  }
  console.log(`Inventory updated for order ${event.data.orderId}`);
}
 
async function handleOrderCancelled(event: OrderCancelledEvent): Promise<void> {
  // Restore inventory from the cancelled order
  const order = await db.query(
    'SELECT items FROM orders WHERE id = $1',
    [event.data.orderId]
  );
 
  for (const item of order.rows[0].items) {
    await db.query(
      'UPDATE products SET stock = stock + $1 WHERE id = $2',
      [item.quantity, item.productId]
    );
  }
}

Idempotente Verarbeitung

Nachrichten können mehr als einmal zugestellt werden — Netzwerk-Retries, Broker-Neustarts, Consumer-Abstürze vor dem Acknowledgement. Jeder Consumer muss Duplikate sicher behandeln.

tstypescript
// lib/idempotency.ts
async function isProcessed(
  eventId: string,
  consumerName: string,
  db: Database
): Promise<boolean> {
  const result = await db.query(
    `SELECT 1 FROM processed_events
     WHERE event_id = $1 AND consumer = $2`,
    [eventId, consumerName]
  );
  return result.rows.length > 0;
}
 
async function markProcessed(
  eventId: string,
  consumerName: string,
  db: Database
): Promise<void> {
  await db.query(
    `INSERT INTO processed_events (event_id, consumer, processed_at)
     VALUES ($1, $2, NOW())
     ON CONFLICT (event_id, consumer) DO NOTHING`,
    [eventId, consumerName]
  );
}
 
// Idempotent consumer wrapper
async function processIdempotently(
  event: BaseEvent,
  consumerName: string,
  handler: (event: BaseEvent) => Promise<void>,
  db: Database
): Promise<void> {
  if (await isProcessed(event.eventId, consumerName, db)) {
    console.log(`Skipping duplicate event: ${event.eventId}`);
    return;
  }
 
  await handler(event);
  await markProcessed(event.eventId, consumerName, db);
}
tstypescript
// ❌ Non-idempotent: processes duplicates, corrupts data
async function handleOrderPlaced(event: OrderPlacedEvent) {
  await db.query(
    'UPDATE accounts SET balance = balance + $1 WHERE id = $2',
    [event.data.total, event.data.customerId]
  );
  // If this event is delivered twice, balance is credited twice!
}
 
// ✅ Idempotent: safely handles duplicate delivery
async function handleOrderPlaced(event: OrderPlacedEvent) {
  await processIdempotently(
    event,
    'payment-consumer',
    async (e) => {
      const orderEvent = e as OrderPlacedEvent;
      await db.query(
        'UPDATE accounts SET balance = balance + $1 WHERE id = $2',
        [orderEvent.data.total, orderEvent.data.customerId]
      );
    },
    db
  );
}

Dead Letter Queues

Nachrichten, deren Verarbeitung fehlschlägt, sollten nicht endlos wiederholt werden. Dead Letter Queues fangen fehlgeschlagene Nachrichten zur Inspektion und manuellen Neuverarbeitung ab.

tstypescript
// Setup dead letter exchange and queue
async function setupDeadLetterQueue(): Promise<void> {
  const channel = await getChannel();
 
  // Dead letter exchange
  await channel.assertExchange('dlx_order_events', 'fanout', {
    durable: true,
  });
 
  // Dead letter queue — stores failed messages
  await channel.assertQueue('dead_letter_order_events', {
    durable: true,
  });
 
  await channel.bindQueue(
    'dead_letter_order_events',
    'dlx_order_events',
    ''
  );
}
 
// Monitoring: check dead letter queue depth
async function getDeadLetterCount(): Promise<number> {
  const channel = await getChannel();
  const info = await channel.checkQueue('dead_letter_order_events');
  return info.messageCount;
}
 
// Reprocess dead letters (manual intervention)
async function reprocessDeadLetters(limit: number = 10): Promise<void> {
  const channel = await getChannel();
 
  for (let i = 0; i < limit; i++) {
    const msg = await channel.get('dead_letter_order_events');
    if (!msg) break;
 
    // Republish to the original exchange for retry
    channel.publish(
      'order_events',
      msg.fields.routingKey,
      msg.content,
      { persistent: true }
    );
 
    channel.ack(msg);
    console.log(`Requeued dead letter: ${msg.properties.messageId}`);
  }
}

Wichtigste Erkenntnisse

  1. Publiziere reichhaltige Events — füge alle Daten bei, die Consumer brauchen; anämische Events, die Rückfragen erfordern, erzeugen Kopplung zwischen Services
  2. Nutze Topic-Exchanges für flexibles Routing — Consumer abonnieren Muster wie order.* und erhalten neue Event-Typen automatisch
  3. Jeder Consumer muss idempotent sein — verfolge verarbeitete Event-IDs, um doppelte Zustellungen sicher zu behandeln
  4. Acknowledge nach der Verarbeitung, nicht davor — stürzt der Consumer nach dem Ack, aber vor der Verarbeitung ab, geht die Nachricht verloren
  5. Dead Letter Queues verhindern endlose Retries — fehlgeschlagene Nachrichten werden zur Inspektion abgefangen, statt die Queue zu blockieren
  6. Prefetch nur eine Nachricht auf einmal — so puffern langsame Consumer keine Nachrichten, die sie nicht verarbeiten können, was Speicherprobleme vermeidet
Wilfredo Rujel

Wilfredo Rujel

Full-Stack-Softwareentwickler

Diesen Beitrag teilenX