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 Serialisierung —
string,objectoderarrayals Payload, JSON-Encoding/Decoding transparent - Lazy Connection — Verbindung erst beim ersten
publish()oderconsume() - External Connections — bestehende Redis, PDO, AMQP oder Kafka-Clients einbinden
- Graceful Shutdown —
SIGTERM/SIGINT-Handler automatisch installiert
Installation
composer require jardisadapter/messagingGitHub: jardisAdapter/messaging
Optionale PHP-Extensions:
| Extension | Für |
|---|---|
ext-redis | Redis-Transport |
ext-amqp | RabbitMQ-Transport |
ext-rdkafka | Kafka-Transport |
ext-pdo | Database-Transport |
Grundlegende Nutzung
Publish und Consume
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
// 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:
$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:
$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
$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
$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:
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:
// 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:
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:
$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:
$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:
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:
// 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:
use JardisAdapter\Messaging\Validation\MessageValidator;
$publisher = (new MessagePublisher($pubFactory->redis($conn)))
->withValidator(new MessageValidator());
// Wirft PublishException bei Arrays mit:
// - Resources
// - ClosuresHandler-Rückgabewerte
Der Handler-Rückgabewert steuert das ACK-Verhalten:
| Return | Verhalten |
|---|---|
true | ACK — Nachricht als verarbeitet markiert |
false | NACK — Nachricht wird requeued (RabbitMQ) oder Offset nicht committed (Kafka) |
| Exception | NACK + Exception wird weitergeworfen |
Fehlerbehandlung
Alle Exceptions erben von MessageException:
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:
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 ← SynchronVerzeichnisstruktur
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.sqlAPI-Referenz
MessagingService
| Methode | Signatur | Beschreibung |
|---|---|---|
publish | publish(string $topic, string|object|array $message, array $options = []): bool | Nachricht publizieren |
consume | consume(string $topic, MessageHandlerInterface $handler, array $options = []): void | Nachrichten konsumieren |
MessagePublisher
| Methode | Signatur | Beschreibung |
|---|---|---|
publish | publish(string $topic, string|object|array $message, array $options = []): bool | Mit Failover |
withValidator | withValidator(MessageValidator $validator): static | Validator hinzufügen |
MessageConsumer
| Methode | Signatur | Beschreibung |
|---|---|---|
consume | consume(string $topic, MessageHandlerInterface $handler, array $options = []): void | Blocking Loop |
stop | stop(): void | Graceful Shutdown |
Vollständiges Beispiel
Event-Driven Architecture mit Redis Streams und Database-Fallback:
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(),
]);