Message Queues aus PHP heraus ansteuern: RabbitMQ und Redis Streams
AI generated
<?php
8.4
PHP · RabbitMQ · Redis · Asynchrone Verarbeitung
Message Queues aus PHP heraus ansteuern
RabbitMQ, Redis Streams und zuverlässige Retry-Strategien

Message Queues entkoppeln PHP-Anwendungen von langsamen oder unzuverlässigen Nachbarsystemen und verlagern Arbeit in asynchrone Worker-Prozesse. Wer versteht, wie Acknowledgments, Prefetch und Dead-Letter-Queues zusammenspielen, kann Message Queues aus PHP heraus so ansteuern, dass Nachrichten weder verloren gehen noch doppelt verarbeitet werden.

18 Min. Lesezeit php-amqplib · RabbitMQ · Redis Streams PHP 8.4 · Worker-Muster

1. Warum PHP Message Queues braucht

Eine klassische PHP-Anfrage läuft synchron: Der Request kommt an, PHP-FPM verarbeitet ihn, und die Antwort geht zurück, sobald alles erledigt ist. Sobald ein Teil dieser Arbeit aber Minuten dauert, etwa der Versand einer Rechnung, die Konvertierung eines Videos oder der Abgleich mit einem langsamen Drittsystem, wird der Request-Response-Zyklus zum Problem. Message Queues lösen das, indem sie die eigentliche Arbeit aus dem Request herauslösen und in eine Warteschlange schreiben, die ein separater Worker-Prozess in seinem eigenen Tempo abarbeitet.

Der zweite große Vorteil von Message Queues ist Entkopplung zwischen Systemen. Statt dass ein PHP-Monolith direkt einen externen Zahlungsdienstleister aufruft und bei dessen Ausfall selbst instabil wird, schreibt er eine Nachricht in eine Queue. Ein Worker verarbeitet diese Nachricht, sobald der Zahlungsdienstleister wieder erreichbar ist. Fällt der Dienstleister für zehn Minuten aus, stauen sich Nachrichten in der Queue, statt dass Kundenanfragen mit Timeout-Fehlern scheitern.

In PHP-Projekten kommen für Message Queues vor allem zwei Technologien zum Einsatz: RabbitMQ als dedizierter Message Broker mit dem AMQP-Protokoll, und Redis Streams als leichtgewichtige Alternative, die ohnehin oft schon für Caching im Einsatz ist. Beide Ansätze werden in diesem Artikel im Detail behandelt, inklusive der Fallstricke, die in produktiven PHP-Anwendungen am häufigsten auftreten.

2. RabbitMQ-Grundlagen: Exchange, Queue, Binding

RabbitMQ organisiert Message Queues über drei zentrale Konzepte: den Exchange, die Queue und das Binding dazwischen. Eine Nachricht wird niemals direkt in eine Queue geschrieben, sondern immer zuerst an einen Exchange gesendet. Der Exchange entscheidet anhand eines Routing-Keys und der konfigurierten Bindings, an welche Queues die Nachricht weitergeleitet wird. Diese Indirektion erlaubt flexible Verteilmuster, etwa eine Nachricht gleichzeitig an mehrere Queues zu senden, ohne dass der Publisher die Empfänger kennen muss.

Der direct-Exchange leitet Nachrichten anhand eines exakten Routing-Key-Abgleichs weiter, der topic-Exchange erlaubt Wildcard-Muster wie order.*.created, und der fanout-Exchange sendet jede Nachricht an alle gebundenen Queues, unabhängig vom Routing-Key. Für die meisten PHP-Anwendungen mit Message Queues reicht ein direct-Exchange mit klar benannten Routing-Keys wie order.created oder invoice.generate völlig aus und bleibt gut nachvollziehbar.

3. Nachrichten mit php-amqplib veröffentlichen

Die Bibliothek php-amqplib/php-amqplib ist der De-facto-Standard, um aus PHP heraus mit RabbitMQ zu sprechen, da sie das AMQP-0-9-1-Protokoll ohne native Erweiterung vollständig in PHP implementiert. Eine Nachricht in Message Queues zu veröffentlichen bedeutet zunächst, eine Verbindung und einen Channel zu öffnen, dann Exchange und Queue zu deklarieren, und erst danach die eigentliche Nachricht zu senden. Wichtig: Deklarationen sind idempotent, ein wiederholter Aufruf von queue_declare() mit identischen Parametern erzeugt keinen Fehler, sondern bestätigt lediglich die bestehende Konfiguration.

Für produktiven Einsatz sollte jede veröffentlichte Nachricht als persistent markiert werden, damit RabbitMQ sie auf die Festplatte schreibt und ein Neustart des Brokers keine Nachrichten verliert. Zusätzlich empfiehlt sich publisher confirms, ein Mechanismus, bei dem RabbitMQ dem Publisher explizit bestätigt, dass die Nachricht sicher im Broker angekommen ist, bevor die PHP-Anwendung davon ausgeht, dass die Zustellung erfolgreich war.


<?php

declare(strict_types=1);

use PhpAmqpLib\Connection\AMQPStreamConnection;
use PhpAmqpLib\Message\AMQPMessage;

final class OrderEventPublisher
{
    private AMQPStreamConnection $connection;

    /**
     * Opens a persistent AMQP connection and enables publisher confirms.
     */
    public function __construct(string $host, int $port, string $user, string $password)
    {
        $this->connection = new AMQPStreamConnection($host, $port, $user, $password);
    }

    public function publishOrderCreated(int $orderId, array $payload): void
    {
        $channel = $this->connection->channel();
        $channel->exchange_declare('orders', 'direct', false, true, false);
        $channel->queue_declare('order.created', false, true, false, false);
        $channel->queue_bind('order.created', 'orders', 'order.created');

        // Confirm mode: broker acknowledges receipt before we consider it safe.
        $channel->confirm_select();

        $message = new AMQPMessage(
            json_encode(['order_id' => $orderId] + $payload, JSON_THROW_ON_ERROR),
            [
                'delivery_mode' => AMQPMessage::DELIVERY_MODE_PERSISTENT,
                'content_type' => 'application/json',
            ]
        );

        $channel->basic_publish($message, 'orders', 'order.created');
        $channel->wait_for_pending_acks(5.0);

        $channel->close();
    }
}

4. Einen zuverlässigen Consumer schreiben

Ein Consumer für Message Queues registriert einen Callback bei RabbitMQ und läuft dann in einer Endlosschleife, die auf eingehende Nachrichten wartet. Der entscheidende Unterschied zu einem klassischen PHP-Request ist, dass dieser Prozess dauerhaft läuft, typischerweise als eigener Systemd-Service oder Supervisor-Worker, nicht als Teil von PHP-FPM. Das bringt eigene betriebliche Herausforderungen mit sich: Speicherlecks über Stunden hinweg, Verbindungsabbrüche zur Datenbank, und die Notwendigkeit, den Prozess bei einem Deployment kontrolliert neu zu starten, statt ihn mitten in der Verarbeitung einer Nachricht zu killen.

Innerhalb des Callbacks muss jeder Fehlerfall explizit behandelt werden. Eine unbehandelte Exception im Consumer-Callback darf niemals dazu führen, dass die Nachricht stillschweigend verloren geht. Stattdessen muss der Consumer entscheiden, ob die Nachricht erneut zugestellt (nack mit Requeue), verworfen, oder in eine separate Fehler-Queue verschoben wird.


<?php

declare(strict_types=1);

use PhpAmqpLib\Connection\AMQPStreamConnection;
use PhpAmqpLib\Message\AMQPMessage;

$connection = new AMQPStreamConnection('rabbitmq.internal', 5672, 'app', 'secret');
$channel = $connection->channel();

$channel->queue_declare('order.created', false, true, false, false);

// Prefetch limits how many unacknowledged messages this worker holds at once.
$channel->basic_qos(null, 10, null);

$callback = function (AMQPMessage $message) use ($channel): void {
    try {
        $payload = json_decode($message->getBody(), true, 512, JSON_THROW_ON_ERROR);
        processOrder($payload);

        // Positive acknowledgment: message is safely removed from the queue.
        $message->ack();
    } catch (\JsonException $e) {
        // Malformed payload will never succeed on retry: discard, no requeue.
        $message->nack(false);
    } catch (\Throwable $e) {
        error_log('Order processing failed: ' . $e->getMessage());
        // Transient failure: requeue for another attempt.
        $message->nack(true);
    }
};

$channel->basic_consume('order.created', '', false, false, false, false, $callback);

while ($channel->is_consuming()) {
    $channel->wait();
}

5. Acknowledgments und Prefetch richtig konfigurieren

Das Acknowledgment-Modell ist das Herzstück zuverlässiger Message Queues. Bestätigt ein Consumer eine Nachricht nicht explizit mit ack(), betrachtet RabbitMQ sie weiterhin als unbearbeitet und liefert sie bei einem Verbindungsabbruch des Consumers erneut aus. Das schützt vor Datenverlust, wenn ein Worker-Prozess mitten in der Verarbeitung abstürzt, erzeugt aber auch die Notwendigkeit, Verarbeitungslogik idempotent zu gestalten, da dieselbe Nachricht im Fehlerfall mehrfach ankommen kann.

Der Prefetch-Wert, gesetzt über basic_qos(), begrenzt, wie viele unbestätigte Nachrichten ein Consumer gleichzeitig erhalten darf. Ein Prefetch von 1 verarbeitet Nachrichten streng sequenziell und ist sicher, aber langsam. Ein zu hoher Prefetch-Wert kann dazu führen, dass ein abstürzender Worker hunderte Nachrichten gleichzeitig als unbestätigt zurücklässt, die dann alle erneut zugestellt werden. Für die meisten Message Queues-Anwendungsfälle in PHP ist ein Prefetch zwischen 5 und 20 ein guter Ausgangspunkt, abhängig von der durchschnittlichen Verarbeitungszeit pro Nachricht.

6. Retry-Strategien und Dead-Letter-Queues

Nicht jeder Fehler in Message Queues ist gleich zu behandeln. Ein Netzwerk-Timeout beim Aufruf eines externen Dienstes rechtfertigt einen erneuten Versuch, ein fehlerhaftes JSON-Format hingegen wird auch beim zehnten Versuch nicht plötzlich gültig. Eine robuste Retry-Strategie unterscheidet deshalb zwischen transienten und permanenten Fehlern und begrenzt die Anzahl der Wiederholungsversuche, um zu verhindern, dass eine defekte Nachricht endlos zwischen Queue und Consumer pendelt.

RabbitMQ bietet dafür x-dead-letter-exchange als Queue-Argument: Überschreitet eine Nachricht die maximale Anzahl an Wiederholungen, oder läuft ihre TTL ab, verschiebt RabbitMQ sie automatisch in eine konfigurierte Dead-Letter-Queue. Dort können defekte Nachrichten manuell inspiziert, korrigiert und erneut eingespielt werden, statt dass sie unbemerkt verloren gehen oder die Hauptqueue blockieren. Ein zusätzlicher Header wie x-retry-count, den der Consumer selbst mitführt und bei jedem erneuten Publish erhöht, macht die Anzahl der bisherigen Versuche für die Anwendungslogik sichtbar.


<?php

declare(strict_types=1);

use PhpAmqpLib\Message\AMQPMessage;

// Queue arguments: route exhausted messages to a dead letter exchange.
$arguments = new \PhpAmqpLib\Wire\AMQPTable([
    'x-dead-letter-exchange' => 'orders.dlx',
    'x-dead-letter-routing-key' => 'order.failed',
    'x-message-ttl' => 60000, // 60s before expiry if unprocessed
]);

$channel->queue_declare('order.created', false, true, false, false, false, $arguments);

function republishWithRetryCount(AMQPMessage $message, \PhpAmqpLib\Channel\AMQPChannel $channel, int $maxRetries = 3): void
{
    $headers = $message->get('application_headers');
    $retryCount = $headers !== null && $headers->hasKey('x-retry-count')
        ? (int) $headers->getNativeData()['x-retry-count'] + 1
        : 1;

    if ($retryCount > $maxRetries) {
        $message->nack(false); // exceeded retries, let dead-lettering take over
        return;
    }

    $retryMessage = new AMQPMessage($message->getBody(), [
        'delivery_mode' => AMQPMessage::DELIVERY_MODE_PERSISTENT,
        'application_headers' => new \PhpAmqpLib\Wire\AMQPTable(['x-retry-count' => $retryCount]),
    ]);

    $channel->basic_publish($retryMessage, 'orders', 'order.created');
    $message->ack(); // remove original, replaced by the retry copy
}

7. Redis Streams als leichtgewichtige Alternative

Wer bereits Redis für Caching oder Sessions im Einsatz hat, kann Redis Streams als deutlich leichtgewichtigere Alternative zu RabbitMQ für Message Queues nutzen, ohne einen zusätzlichen Broker betreiben zu müssen. Ein Redis Stream ist eine append-only Log-Struktur, in die mit XADD Einträge geschrieben werden, und aus der Consumer-Gruppen mit XREADGROUP Nachrichten konsumieren, ähnlich dem Kafka-Modell, aber ohne dessen Betriebskomplexität.

Consumer-Gruppen bei Redis Streams bieten ein eigenes Acknowledgment-Modell über XACK, funktional vergleichbar mit RabbitMQs Ack-Mechanismus. Nicht bestätigte Einträge bleiben in der Pending Entries List sichtbar und können mit XCLAIM von einem anderen Consumer übernommen werden, falls der ursprüngliche Worker ausgefallen ist. Der große Vorteil gegenüber RabbitMQ liegt in der operativen Einfachheit, der Nachteil in weniger ausgefeilten Routing-Möglichkeiten, es gibt keine Exchange-Konzepte, nur benannte Streams.


<?php

declare(strict_types=1);

$redis = new \Redis();
$redis->connect('redis.internal', 6379);

// Publisher: append an entry to the stream.
$redis->xAdd('orders:created', '*', [
    'order_id' => '10245',
    'total' => '129.90',
]);

// Consumer group setup, run once during deployment.
try {
    $redis->xGroup('CREATE', 'orders:created', 'invoice-workers', '0');
} catch (\RedisException $e) {
    // Group already exists — safe to ignore on redeploy.
}

// Worker loop: read new entries assigned to this consumer group.
while (true) {
    $entries = $redis->xReadGroup(
        'invoice-workers',
        'worker-1',
        ['orders:created' => '>'],
        10,
        5000
    );

    foreach ($entries['orders:created'] ?? [] as $id => $fields) {
        try {
            generateInvoice($fields);
            $redis->xAck('orders:created', 'invoice-workers', $id);
        } catch (\Throwable $e) {
            error_log('Invoice generation failed for ' . $id . ': ' . $e->getMessage());
            // Left unacknowledged — visible in pending entries for retry via XCLAIM.
        }
    }
}

8. Idempotenz: doppelte Verarbeitung sicher verhindern

Sowohl RabbitMQ als auch Redis Streams garantieren im Standardbetrieb "at least once"-Zustellung, niemals "exactly once". Das bedeutet, jede Nachricht kann theoretisch mehr als einmal beim Consumer ankommen, etwa wenn die Bestätigung nach erfolgreicher Verarbeitung, aber vor dem Senden des Acks verloren geht. Wer Message Queues produktiv einsetzt, muss deshalb jede Verarbeitungslogik idempotent gestalten, damit eine doppelte Zustellung keine doppelten Effekte erzeugt.

Die gängigste Lösung ist ein eindeutiger Idempotenz-Schlüssel pro Nachricht, meist eine UUID, die vor der eigentlichen Verarbeitung gegen eine Tabelle oder einen Redis-Set geprüft wird. Existiert der Schlüssel bereits, wird die Nachricht als bereits verarbeitet markiert und übersprungen, ohne die Geschäftslogik erneut auszuführen. Für Vorgänge wie das Versenden einer Rechnungs-E-Mail oder das Abbuchen eines Zahlungsbetrags ist diese Absicherung nicht optional, sondern notwendige Voraussetzung, um Message Queues überhaupt sicher in Zahlungsprozessen einzusetzen.

9. RabbitMQ vs. Redis Streams im Vergleich

Die Wahl zwischen RabbitMQ und Redis Streams für Message Queues in PHP-Projekten hängt stark von den bestehenden betrieblichen Gegebenheiten und den Anforderungen an Routing-Flexibilität ab. Die folgende Tabelle fasst die wichtigsten Unterschiede zusammen.

Kriterium RabbitMQ Redis Streams
Betrieblicher Aufwand Zusätzlicher dedizierter Broker Nutzt vorhandene Redis-Instanz
Routing-Flexibilität Exchanges, Topics, Fanout Nur benannte Streams
Dead-Letter-Queues Nativ über x-dead-letter-exchange Manuell über Pending Entries List
Persistenz-Garantien Ausgereift, disk-basiert Abhängig von Redis-Persistenzkonfiguration
Einstiegshürde in PHP Mittel, php-amqplib nötig Niedrig, phpredis meist schon vorhanden

Für komplexe Verteilmuster mit mehreren unabhängigen Empfängern und strengen Zustellungsgarantien bleibt RabbitMQ die robustere Wahl unter den Message Queues-Technologien. Für kleinere PHP-Projekte, die bereits auf Redis setzen und einfache Worker-Warteschlangen benötigen, sparen Redis Streams operative Komplexität, ohne auf grundlegende Zuverlässigkeitsmerkmale wie Consumer-Gruppen und Acknowledgments zu verzichten.

Mironsoft

Asynchrone Architekturen und Message-Queue-Integration

Blockieren langsame Drittsysteme Ihre PHP-Requests?

Wir entwerfen und implementieren Message-Queue-Architekturen mit RabbitMQ oder Redis Streams, inklusive Retry-Strategien, Dead-Letter-Queues und idempotenter Verarbeitung für zuverlässige PHP-Worker.

Architektur-Design

RabbitMQ oder Redis Streams je nach Anforderung auswählen

Worker-Implementierung

Zuverlässige Consumer mit Retry, DLQ und Idempotenz

Monitoring

Queue-Tiefe, Fehlerraten und Consumer-Health im Blick behalten

10. Zusammenfassung

Message Queues lösen ein zentrales Problem synchroner PHP-Anwendungen: langsame oder unzuverlässige Arbeit muss nicht mehr den Request-Response-Zyklus blockieren. RabbitMQ mit php-amqplib bietet ausgereiftes Routing über Exchanges, verlässliche Acknowledgments und native Dead-Letter-Queues. Redis Streams punktet mit geringerem Betriebsaufwand, wenn Redis ohnehin schon im Stack vorhanden ist, bei etwas einfacherer Routing-Logik.

Unabhängig von der gewählten Technologie gelten dieselben Grundregeln: Nachrichten persistent speichern, Prefetch begrenzen, transiente von permanenten Fehlern unterscheiden, und jede Verarbeitung idempotent gestalten. Wer diese Prinzipien beim Einsatz von Message Queues konsequent umsetzt, baut PHP-Systeme, die auch bei Ausfällen von Drittsystemen stabil und nachvollziehbar bleiben.

Message Queues aus PHP heraus ansteuern — Das Wichtigste auf einen Blick

RabbitMQ mit php-amqplib

Exchange, Queue und Binding trennen Publisher von Consumern. Publisher Confirms sichern Zustellung ab.

Acknowledgments & Prefetch

Explizites ack() verhindert Datenverlust. Prefetch begrenzt gleichzeitig unbestätigte Nachrichten pro Worker.

Redis Streams

Leichtgewichtige Alternative mit Consumer-Gruppen, XACK und XCLAIM, ohne zusätzlichen Broker.

Idempotenz

At-least-once-Zustellung verlangt idempotente Verarbeitungslogik über eindeutige Schlüssel pro Nachricht.

11. FAQ: Message Queues aus PHP heraus ansteuern

1Wann Message Queues einsetzen?
Wenn Arbeit deutlich länger dauert als eine Request-Antwortzeit, oder ein instabiles Drittsystem den Hauptprozess nicht blockieren soll.
2Exchange vs. Queue?
Exchange empfängt und routet Nachrichten, Consumer lesen ausschließlich aus Queues, nie direkt aus dem Exchange.
3Genau einmal verarbeitet?
Nein, Standard ist at-least-once. Doppelte Zustellung ist möglich, Verarbeitung muss idempotent sein.
4Wofür Dead-Letter-Queue?
Macht dauerhaft fehlschlagende Nachrichten sichtbar, statt sie zu verlieren oder die Hauptqueue zu blockieren.
5Transient vs. permanent?
Transiente Fehler wie Timeouts rechtfertigen Requeue, permanente Fehler wie ungültiges JSON sollten direkt verworfen werden.
6Warum Redis Streams statt RabbitMQ?
Wenn Redis bereits vorhanden ist und keine komplexen Routing-Muster nötig sind, spart das einen zusätzlichen Broker.
7Was macht Prefetch?
Begrenzt gleichzeitig unbestätigte Nachrichten pro Consumer. Zu hoher Wert lässt bei Absturz viele Nachrichten unbestätigt zurück.
8Wie läuft ein Consumer dauerhaft?
Als eigener Prozess über Systemd oder Supervisor, mit einer Schleife auf channel->wait() oder XREADGROUP.
9Was bei Absturz während Verarbeitung?
Unbestätigte Nachricht wird erneut zugestellt, bei Redis Streams bleibt sie in der Pending Entries List für XCLAIM sichtbar.
10Müssen alle Nachrichten persistent sein?
Für geschäftskritische Vorgänge wie Bestellungen ja, für unkritische Benachrichtigungen kann auf Persistenz verzichtet werden.