Skip to content

Messaging

Multi-Transport Messaging mit Redis, Kafka, RabbitMQ, Database und InMemory: eine API für alle.

Einführung

Messaging in PHP-Anwendungen bedeutet oft: sich auf einen Broker festlegen und dessen Client-Library tief in die Anwendung verdrahten. Kafka-Code sieht völlig anders aus als RabbitMQ-Code, und das Wechseln des Transports bedeutet einen Rewrite. Für Tests braucht man dann noch einen laufenden Broker oder Mock-Konstruktionen.

jardisadapter/messaging löst das mit einer einheitlichen Publish/Consume-API, die identisch über fünf Transporte funktioniert:

  • Redis — Pub/Sub für Fire-and-Forget, Streams mit Consumer Groups für persistente Verarbeitung
  • Kafka — Producer/Consumer mit Consumer Groups, SASL-Authentifizierung
  • RabbitMQ — AMQP Exchange/Queue mit ACK/NACK und Routing
  • Database — Transactional Outbox Pattern über PDO, Point-to-Point und Fan-Out
  • InMemory — Synchroner Transport für Tests ohne Broker

Weitere Features:

  • Priority-Failover — mehrere Transporte mit automatischem Fallback
  • Automatische Serialisierungstring, object oder array als Payload, JSON-Encoding/Decoding transparent
  • Lazy Connection — Verbindung erst beim ersten publish() oder consume()
  • External Connections — bestehende Redis, PDO, AMQP oder Kafka-Clients einbinden
  • Graceful ShutdownSIGTERM/SIGINT-Handler automatisch installiert

Installation

bash
composer require jardisadapter/messaging

GitHub: jardisAdapter/messaging

Optionale PHP-Extensions:

ExtensionFür
ext-redisRedis-Transport
ext-amqpRabbitMQ-Transport
ext-rdkafkaKafka-Transport
ext-pdoDatabase-Transport

Grundlegende Nutzung

Publish und Consume

php
use JardisAdapter\Messaging\Factory\ConnectionFactory;
use JardisAdapter\Messaging\Factory\PublisherFactory;
use JardisAdapter\Messaging\Factory\ConsumerFactory;
use JardisAdapter\Messaging\MessagePublisher;
use JardisAdapter\Messaging\MessageConsumer;
use JardisAdapter\Messaging\Handler\CallbackHandler;

$connFactory = new ConnectionFactory();
$pubFactory  = new PublisherFactory();
$conFactory  = new ConsumerFactory();

// Redis-Verbindung
$conn = $connFactory->redis('localhost', 6379);

// Publisher und Consumer erstellen
$publisher = new MessagePublisher($pubFactory->redis($conn));
$consumer  = new MessageConsumer($conFactory->redis($conn));

// Nachricht publizieren (Array → automatisch JSON-encodiert)
$publisher->publish('orders', [
    'order_id' => 42,
    'total' => 299.99,
    'customer_id' => 7,
]);

// Nachrichten konsumieren
$consumer->consume('orders', new CallbackHandler(
    function (string|array $message, array $metadata): bool {
        // $message ist bereits dekodiert: ['order_id' => 42, ...]
        processOrder($message);
        return true;  // ACK
    }
));

Payload-Typen

php
// String — wird direkt gesendet
$publisher->publish('logs', 'Simple log message');

// Array — wird JSON-encodiert
$publisher->publish('events', ['type' => 'OrderCreated', 'data' => [...]]);

// Object — wird JSON-encodiert (public Properties oder JsonSerializable)
$publisher->publish('events', $domainEvent);

Transporte

Redis Pub/Sub

Fire-and-Forget-Fanout. Subscriber müssen zum Zeitpunkt des Publishings verbunden sein:

php
$conn = $connFactory->redis('localhost', 6379, password: 'secret');

$publisher = new MessagePublisher($pubFactory->redis($conn));
$consumer  = new MessageConsumer($conFactory->redis($conn));

Redis Streams

Persistente Nachrichten mit optionalen Consumer Groups:

php
$publisher = new MessagePublisher($pubFactory->redis($conn, useStreams: true));
$consumer  = new MessageConsumer($conFactory->redis($conn, useStreams: true));

// Einfaches Lesen
$consumer->consume('orders', $handler);

// Mit Consumer Group — jede Nachricht wird nur von einem Consumer in der Gruppe verarbeitet
$consumer->consume('orders', $handler, [
    'group'    => 'order-service',
    'consumer' => 'worker-1',
]);

Kafka

php
$producer = $connFactory->kafka('broker1:9092,broker2:9092',
    username: 'app', password: 'secret');

$consumer = $connFactory->kafkaConsumer('broker1:9092,broker2:9092',
    groupId: 'order-service',
    username: 'app', password: 'secret');

$publisher = new MessagePublisher($pubFactory->kafka($producer));
$msgConsumer = new MessageConsumer($conFactory->kafka($consumer));

RabbitMQ

php
$conn = $connFactory->rabbitMq('localhost', 5672, 'guest', 'guest', [
    'exchange_name' => 'app.events',
]);

$publisher = new MessagePublisher($pubFactory->rabbitMq($conn));
$consumer  = new MessageConsumer($conFactory->rabbitMq($conn, 'order-queue'));

Database (Transactional Outbox)

Kein Broker nötig, nutzt die eigene Datenbank als Message-Store:

php
use JardisAdapter\Messaging\Config\DatabaseTransportOptions;

$conn = $connFactory->database('mysql:host=localhost;dbname=app', 'user', 'pass');

$options = new DatabaseTransportOptions(
    table: 'domain_events',
    pollingIntervalMs: 1000,
    batchSize: 10,
    maxAttempts: 3,
    deleteAfterProcessing: false,
);

$publisher = new MessagePublisher($pubFactory->database($conn, $options));
$consumer  = new MessageConsumer($conFactory->database($conn, $options));

Fan-Out (Consumer Groups)

Mehrere unabhängige Consumer-Gruppen verarbeiten dieselben Events:

php
// E-Mail-Service
$consumer->consume('InvoiceCreated', $emailHandler, ['group' => 'email-service']);

// PDF-Service
$consumer->consume('InvoiceCreated', $pdfHandler, ['group' => 'pdf-service']);

Jede Gruppe verarbeitet jedes Event unabhängig. Events ohne group-Option werden im Point-to-Point-Modus konsumiert.

InMemory (Tests)

Synchroner Transport ohne Broker, ideal für Unit-Tests:

php
use JardisAdapter\Messaging\Transport\InMemoryTransport;
use JardisAdapter\Messaging\Publisher\InMemoryPublisher;
use JardisAdapter\Messaging\Consumer\InMemoryConsumer;

$transport = new InMemoryTransport();

$publisher = new MessagePublisher(new InMemoryPublisher($transport));
$consumer  = new MessageConsumer(new InMemoryConsumer($transport));

$publisher->publish('orders', ['order_id' => 123]);
$consumer->consume('orders', $handler);

// Assertions
$transport->getMessageCount('orders');  // 0 (konsumiert)

Oder über die Factories mit shared Transport:

php
$transport = new InMemoryTransport();
$pubFactory->setSharedTransport($transport);
$conFactory->setSharedTransport($transport);

$publisher = new MessagePublisher($pubFactory->inMemory());
$consumer  = new MessageConsumer($conFactory->inMemory());

Priority-Failover

Mehrere Transporte mit automatischem Fallback:

php
$publisher = new MessagePublisher(
    $pubFactory->redis($primaryRedis),      // Erster Versuch
    $pubFactory->redis($secondaryRedis),    // Fallback
    $pubFactory->database($dbConn),         // Zweiter Fallback
);

// Bei Fehler im ersten Transport wird automatisch der nächste versucht
$publisher->publish('orders', $orderData);

Die Reihenfolge im Konstruktor bestimmt die Priorität: der erste Transport ist der höchstpriorisierte. Das Fallback gilt für MessagePublisher und MessageConsumer gleichermaßen und greift nur bei MessageException. Andere Exceptions (z. B. JsonException bei ungültigem Payload) propagieren sofort.

MessagingService (DI-Container)

Lazy-Loading-Orchestrator für DI-Container, Verbindung erst beim ersten Aufruf:

php
use JardisAdapter\Messaging\MessagingService;

$messaging = new MessagingService(
    publisherFactory: fn() => new MessagePublisher(
        $pubFactory->redis($connFactory->redis('localhost'))
    ),
    consumerFactory: fn() => new MessageConsumer(
        $conFactory->redis($connFactory->redis('localhost'))
    ),
);

// Verbindung wird erst hier aufgebaut
$messaging->publish('topic', $payload);
$messaging->consume('topic', $handler);

External Connections

Bestehende Client-Instanzen einbinden:

php
// Bestehende Redis-Instanz
$conn = $connFactory->fromRedis($existingRedis, manageLifecycle: false);

// Bestehende PDO
$conn = $connFactory->fromPdo($existingPdo, manageLifecycle: false);

// Bestehende AMQP-Verbindung
$conn = $connFactory->fromAmqp($existingAmqp, exchangeName: 'app.events');

// Bestehender Kafka Producer
$conn = $connFactory->fromKafkaProducer($existingProducer);

Mit manageLifecycle: false bleibt die Lifecycle-Kontrolle beim externen System: disconnect() ist ein No-Op.

Message Validation

Payload-Validation ist immer aktiv: der MessageValidator wird bereits im Constructor gesetzt. withValidator() tauscht ihn gegen einen eigenen Validator aus (kein Aktivierungsschalter). Geprüft werden Array-Payloads; Objekte gehen direkt durch json_encode:

php
use JardisAdapter\Messaging\Validation\MessageValidator;

$publisher = (new MessagePublisher($pubFactory->redis($conn)))
    ->withValidator(new MessageValidator());

// Wirft PublishException bei Arrays mit:
// - Resources
// - Closures

Handler-Rückgabewerte

Der Handler-Rückgabewert steuert das ACK-Verhalten:

ReturnVerhalten
trueACK — Nachricht als verarbeitet markiert
falseNACK — Nachricht wird requeued (RabbitMQ) oder Offset nicht committed (Kafka)
ExceptionNACK + Exception wird weitergeworfen

Fehlerbehandlung

Alle Exceptions erben von MessageException:

php
use JardisSupport\Contract\Messaging\Exception\MessageException;

try {
    $publisher->publish('topic', $payload);
} catch (MessageException $e) {
    // Publish fehlgeschlagen auf allen Transports
}

Database-Schema

Für den Database-Transport:

sql
CREATE TABLE domain_events (
    id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
    topic VARCHAR(255) NOT NULL,
    payload TEXT NOT NULL,
    created_at DATETIME(6) NOT NULL,
    processed_at DATETIME(6) NULL DEFAULT NULL,
    attempts TINYINT UNSIGNED NOT NULL DEFAULT 0,
    last_error TEXT NULL DEFAULT NULL,
    INDEX idx_unprocessed (processed_at, created_at),
    INDEX idx_topic (topic, processed_at)
);

CREATE TABLE domain_event_subscriptions (
    id BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
    event_id BIGINT UNSIGNED NOT NULL,
    consumer_group VARCHAR(255) NOT NULL,
    processed_at DATETIME(6) NULL DEFAULT NULL,
    attempts TINYINT UNSIGNED NOT NULL DEFAULT 0,
    last_error TEXT NULL DEFAULT NULL,
    UNIQUE INDEX idx_event_group (event_id, consumer_group),
    INDEX idx_pending (consumer_group, processed_at)
);

Architektur

MessagingService                       ← Lazy-Loading Orchestrator
├── MessagePublisher                   ← Publisher-Facade (Priority-Failover)
│   ├── RedisPublisher                 ← publish() / xAdd()
│   ├── KafkaPublisher                 ← produce() + flush()
│   ├── RabbitMqPublisher              ← AMQPExchange::publish()
│   ├── DatabasePublisher              ← INSERT INTO domain_events
│   └── InMemoryPublisher              ← InMemoryTransport
└── MessageConsumer                    ← Consumer-Facade (SIGTERM-Handler)
    ├── RedisConsumer                  ← subscribe() / xRead() / xReadGroup()
    ├── KafkaConsumer                  ← consume() Poll-Loop
    ├── RabbitMqConsumer               ← AMQPQueue::get() + ACK/NACK
    ├── DatabaseConsumer               ← PDO Polling (P2P + Fan-Out)
    └── InMemoryConsumer               ← Synchron

Verzeichnisstruktur

src/
├── MessagingService.php            ← Orchestrator
├── MessagePublisher.php            ← Publisher-Facade
├── MessageConsumer.php             ← Consumer-Facade
├── Config/
│   ├── ConnectionConfig.php
│   └── DatabaseTransportOptions.php
├── Connection/
│   ├── RedisConnection.php
│   ├── KafkaConnection.php
│   ├── KafkaConsumerConnection.php
│   ├── RabbitMqConnection.php
│   ├── DatabaseConnection.php
│   ├── External*.php               ← External Wrappers
│   └── NullConnection.php
├── Factory/
│   ├── ConnectionFactory.php
│   ├── PublisherFactory.php
│   └── ConsumerFactory.php
├── Publisher/
│   ├── RedisPublisher.php
│   ├── KafkaPublisher.php
│   ├── RabbitMqPublisher.php
│   ├── DatabasePublisher.php
│   └── InMemoryPublisher.php
├── Consumer/
│   ├── RedisConsumer.php
│   ├── KafkaConsumer.php
│   ├── RabbitMqConsumer.php
│   ├── DatabaseConsumer.php
│   └── InMemoryConsumer.php
├── Transport/
│   └── InMemoryTransport.php
├── Handler/
│   └── CallbackHandler.php
├── Validation/
│   └── MessageValidator.php
└── Schema/
    └── domain_events.sql

API-Referenz

MessagingService

MethodeSignaturBeschreibung
publishpublish(string $topic, string|object|array $message, array $options = []): boolNachricht publizieren
consumeconsume(string $topic, MessageHandlerInterface $handler, array $options = []): voidNachrichten konsumieren

MessagePublisher

MethodeSignaturBeschreibung
publishpublish(string $topic, string|object|array $message, array $options = []): boolMit Failover
withValidatorwithValidator(MessageValidator $validator): staticValidator hinzufügen

MessageConsumer

MethodeSignaturBeschreibung
consumeconsume(string $topic, MessageHandlerInterface $handler, array $options = []): voidBlocking Loop
stopstop(): voidGraceful Shutdown

Vollständiges Beispiel

Event-Driven Architecture mit Redis Streams und Database-Fallback:

php
use JardisAdapter\Messaging\Factory\ConnectionFactory;
use JardisAdapter\Messaging\Factory\PublisherFactory;
use JardisAdapter\Messaging\Factory\ConsumerFactory;
use JardisAdapter\Messaging\MessagePublisher;
use JardisAdapter\Messaging\MessageConsumer;
use JardisAdapter\Messaging\Handler\CallbackHandler;
use JardisAdapter\Messaging\Config\DatabaseTransportOptions;

$connFactory = new ConnectionFactory();
$pubFactory  = new PublisherFactory();
$conFactory  = new ConsumerFactory();

// Redis als primärer Transport, Database als Fallback
$redis = $connFactory->redis('redis.internal', 6379, password: $redisPassword);
$db    = $connFactory->database('mysql:host=localhost;dbname=app', 'user', 'pass');
$dbOpts = new DatabaseTransportOptions(table: 'domain_events');

$publisher = new MessagePublisher(
    $pubFactory->redis($redis, useStreams: true),
    $pubFactory->database($db, $dbOpts),
);

// Domain Event publizieren
$publisher->publish('OrderCreated', [
    'order_id'    => $orderId,
    'customer_id' => $customerId,
    'total'       => 299.99,
    'items'       => $items,
]);

// Worker: Events mit Consumer Group verarbeiten
$consumer = new MessageConsumer(
    $conFactory->redis($redis, useStreams: true),
);

$consumer->consume('OrderCreated', new CallbackHandler(
    function (string|array $message, array $metadata): bool {
        $orderId = $message['order_id'];

        // Bestätigung senden
        $mailer->send(buildConfirmationMail($orderId));

        // Lagerbestand reservieren
        $inventory->reserve($message['items']);

        return true;  // ACK
    }
), [
    'group'    => 'fulfillment-service',
    'consumer' => gethostname(),
]);