Ein Checkout-Endpunkt, der ein Produktbild verkleinert, eine Bestätigungs-E-Mail verschickt und eine Empfehlungs-Engine aktualisiert, bevor er mit 200 OK antwortet, ist nur so schnell wie sein langsamster Schritt. Läuft auch nur einer davon in ein Timeout, scheitert die ganze Anfrage. Eine Message Queue lässt den Endpunkt diese Arbeit als Nachrichten übergeben und sofort antworten, während separate Worker sie sich nach eigenem Zeitplan holen.

Was eine Message Queue tatsächlich macht

Eine Message Queue sitzt zwischen zwei Teilen eines Systems, die nicht im selben Moment laufen müssen. Ein Producer veröffentlicht eine Nachricht, die eine zu erledigende Aufgabe beschreibt. Die Queue speichert sie. Ein Consumer holt sie ab, verarbeitet sie und bestätigt sie. Die beiden sprechen nie direkt miteinander.

Denken Sie an einen Briefkasten. Der Absender wirft den Brief ein und geht weiter; der Empfänger sieht nach, wenn er Zeit hat, und ist er eine Stunde weg, warten die Briefe einfach.

Daraus ergeben sich drei Dinge, die ein direkter Funktionsaufruf nicht bieten kann:

  • Entkopplung. Der Checkout-Endpunkt muss nicht wissen, wie die Bildverkleinerung funktioniert oder dass sie langsam ist. Er veröffentlicht {"event": "order_placed", "order_id": 4821} und macht weiter.
  • Pufferung. Kommen 500 Bestellungen in derselben Sekunde an, fängt die Queue den Spitzenwert auf. Die Worker verarbeiten sie im Tempo, das sie durchhalten, statt 500 Anfragen gleichzeitig blockiert zu haben.
  • Unabhängige Skalierung und Ausfälle. Betreiben Sie drei Worker fürs Verkleinern und einen für E-Mails, starten Sie den einen neu, ohne den anderen anzufassen, deployen Sie den Checkout-Dienst, ohne seine Worker neu auszuliefern.

RabbitMQ in Docker betreiben

RabbitMQ ist der verbreitetste universelle Message Broker und eine vernünftige Standardwahl, wenn Sie noch keine Queue haben. Laden Sie das Management-Image, damit Sie neben dem Broker auch eine Web-UI bekommen:

1
2
3
4
5
docker run -d --hostname mq --name rabbitmq \
  -e RABBITMQ_DEFAULT_USER=app \
  -e RABBITMQ_DEFAULT_PASS=changeme \
  -p 5672:5672 -p 15672:15672 \
  rabbitmq:4-management

Port 5672 ist das AMQP-Protokoll, mit dem sich Ihre Anwendung verbindet. Port 15672 ist die Management-UI: Öffnen Sie http://localhost:15672, loggen Sie sich mit app / changeme ein, und Sie sehen Queues, Nachrichtenraten und Verbindungen live. Der Standard-Login guest/guest funktioniert nur innerhalb des Containers, setzen Sie also einen eigenen Benutzer für alles, womit Sie sich vom Host aus verbinden.

Wenn Sie Docker noch nie benutzt haben, erklärt Was ist Docker und wofür wird es verwendet Images, Container und Ports. Für einen echten Stack betreiben Sie den Broker zusammen mit Ihrer Anwendung über Compose — siehe Was Ist Docker Compose und Wie Wird Es Verwendet:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
services:
  rabbitmq:
    image: rabbitmq:4-management
    environment:
      - RABBITMQ_DEFAULT_USER=app
      - RABBITMQ_DEFAULT_PASS=changeme
    ports:
      - "5672:5672"
      - "15672:15672"

  worker:
    build: ./worker
    depends_on:
      - rabbitmq
    environment:
      - RABBITMQ_URL=amqp://app:changeme@rabbitmq:5672/

worker erreicht den Broker über den Hostnamen rabbitmq. Compose setzt beide Services in dasselbe user-defined Netzwerk, sodass der Servicename zugleich als auflösbarer Hostname dient — der Mechanismus, den Was ist ein Docker-Netzwerk und wie wird es verwendet beschreibt.

Eine Nachricht veröffentlichen und konsumieren

pika ist der Standard-Python-Client für das Protokoll von RabbitMQ (AMQP 0-9-1). Installieren Sie es mit pip install pika. Ein Producer, der eine Aufgabe veröffentlicht, sieht so aus:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
import pika

connection = pika.BlockingConnection(
    pika.ConnectionParameters(
        host="localhost",
        credentials=pika.PlainCredentials("app", "changeme"),
    )
)
channel = connection.channel()

channel.queue_declare(queue="resize_image", durable=True)

channel.basic_publish(
    exchange="",
    routing_key="resize_image",
    body='{"order_id": 4821, "image_url": "https://example.com/p/4821.jpg"}',
    properties=pika.BasicProperties(delivery_mode=pika.DeliveryMode.Persistent),
)
connection.close()

durable=True bei queue_declare sagt RabbitMQ, die Queue selbst über einen Broker-Neustart hinweg zu behalten. delivery_mode=pika.DeliveryMode.Persistent sagt ihm, jede Nachricht auf die Festplatte zu schreiben statt sie nur im Speicher zu halten. Ohne beides löscht ein docker restart des Brokers stillschweigend alles, was noch nicht konsumiert war.

Der Consumer holt sich Nachrichten aus derselben Queue und bestätigt jede erst, wenn die Arbeit tatsächlich erledigt ist:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
import pika

connection = pika.BlockingConnection(
    pika.ConnectionParameters(
        host="localhost",
        credentials=pika.PlainCredentials("app", "changeme"),
    )
)
channel = connection.channel()
channel.queue_declare(queue="resize_image", durable=True)
channel.basic_qos(prefetch_count=1)

def handle_message(ch, method, properties, body):
    resize_image(body)  # die eigentliche Arbeit
    ch.basic_ack(delivery_tag=method.delivery_tag)

channel.basic_consume(queue="resize_image", on_message_callback=handle_message)
channel.start_consuming()

Zwei Einstellungen erledigen hier die Zuverlässigkeitsarbeit:

  • basic_qos(prefetch_count=1) verhindert, dass RabbitMQ einem Worker eine zweite Nachricht gibt, bevor er die erste bestätigt hat. Ohne das kann ein ausgelasteter Worker auf zehn Nachrichten sitzen, während ein freier gar nichts bekommt.
  • basic_ack nach der Arbeit, nicht davor. Stürzt der Prozess mitten im Verkleinern ab, wurde die Nachricht nie bestätigt, also liefert RabbitMQ sie erneut an einen anderen Worker, statt sie zu verlieren. Lassen Sie auto_ack deaktiviert, was bei pika ohnehin der Standard ist.

Starten Sie drei Kopien des Consumer-Skripts, und RabbitMQ verteilt die Queue round-robin auf sie. Das ist die gesamte Skalierungsgeschichte für dieses Pattern: kein Koordinationscode, kein geteilter Zustand zwischen Workern.

Manche Nachrichten scheitern bei jeder Zustellung: ein fehlerhaftes Payload, ein Aufruf an einen Dienst, den es endgültig nicht mehr gibt. Die erneute Zustellung macht daraus eine Endlosschleife durch Ihre Consumer. Setzen Sie x-dead-letter-exchange auf der Queue, und Nachrichten, die Sie ablehnen ohne sie erneut einzureihen (basic_nack mit requeue=False), landen in einem separaten Exchange, wo Sie sie von Hand prüfen.

Wie Exchanges in RabbitMQ Nachrichten routen

Die obigen Beispiele veröffentlichen mit exchange="", dem Default Exchange, der eine Nachricht direkt an die in routing_key genannte Queue routet. Das deckt die meisten Task-Queue-Fälle ab. Das eigentliche Modell von RabbitMQ setzt vor jede Queue einen Exchange, und der Exchange-Typ bestimmt das Routing-Verhalten:

Exchange-TypRoutet zuNutzen Sie ihn für
direct (der Default Exchange ist einer)Die Queue, die exakt zum Routing Key passtTask Queues, eine Consumer-Gruppe pro Queue
fanoutJede gebundene Queue, unabhängig vom Routing KeyEin Event an mehrere unabhängige Consumer ausstrahlen
topicQueues, deren Binding-Pattern zum Routing Key passt (order.*.created)Selektive Verteilung, wenn Consumer nur eine Teilmenge der Event-Typen wollen

Ein fanout-Exchange ist, wie Sie mit RabbitMQ Publish/Subscribe bekommen: Sie veröffentlichen order_placed einmal, und der E-Mail-Dienst, der Analytics-Dienst und der Betrugserkennungs-Dienst bekommen jeweils ihre eigene Kopie über ihre eigene Queue, statt sich um dieselben Nachrichten zu streiten.

Wann eine Message Queue das falsche Werkzeug ist

  • Der Aufrufer braucht die Antwort, um selbst zu antworten. Eine Queue ist für “irgendwann erledigen”, nicht für “jetzt berechnen”. Braucht Ihr Checkout-Endpunkt die berechneten Versandkosten, bevor er antworten kann, ist das ein synchroner Aufruf, kein Queue-Job.
  • Sie brauchen strikte Reihenfolge über die gesamte Queue. RabbitMQ garantiert die Reihenfolge pro Queue bei genau einem Consumer, aber fügen Sie für mehr Durchsatz einen zweiten Consumer hinzu, können zwei Nachrichten außer der Reihe ankommen. Zählt die Reihenfolge, wie in einem Event-Log oder einer Zustandsmaschine, brauchen Sie ein Werkzeug für geordnete Streams, etwa Kafka-Partitionen mit einem Key auf der Entitäts-ID.
  • Der Job ist klein, synchron, und Sie betreiben noch keinen Broker. Ein Hintergrund-Thread oder ein In-Process-Job-Runner bedeuten weniger operative Angriffsfläche, als RabbitMQ für eine Cron-große Aufgabe aufzusetzen und zu überwachen.
  • Sie haben bereits Redis und können den gelegentlichen Verlust eines Jobs tolerieren. Redis Streams oder eine Bibliothek wie BullMQ geben Ihnen eine Queue ohne neue Infrastruktur, auf Kosten schwächerer Zustellgarantien bei einem Broker-Absturz. Prüfen Sie zuerst, was Redis in Ihrem Stack bereits tut: Wie Redis-Caching funktioniert und wie Sie es nutzen.

Ihre erste Aufgabe auf eine Queue verlagern

Wählen Sie etwas, das eine Anfrage blockiert, ohne dass es müsste: die Bestätigungs-E-Mail, das Thumbnail, den Analytics-Ping. Betreiben Sie RabbitMQ in Docker, deklarieren Sie eine durable Queue dafür, verlagern Sie es. Setzen Sie prefetch_count=1 und bestätigen Sie von Tag eins an manuell, denn Zuverlässigkeit nachzurüsten, nachdem ein Absturz bereits eine Nachricht gekostet hat, ist teurer als diese zwei Zeilen jetzt zu schreiben. Greifen Sie erst zu einem Fanout- oder Topic-Exchange, wenn ein zweiter Consumer dasselbe Event will.

Verwandte Artikel