Installation:
composer require koco/messenger-kafka
For non-Flex projects, add Koco\Kafka\KocoKafkaBundle::class to config/bundles.php.
Configure DSN in config/packages/messenger.yaml:
messenger:
transports:
kafka:
dsn: '%env(KAFKA_DSN)%' # e.g., kafka://localhost:9092
options:
topic: 'your_topic_name'
group_id: 'your_consumer_group'
Dispatch a message:
use Symfony\Component\Messenger\MessageBusInterface;
$bus->dispatch(new YourMessage());
Consume messages (via Symfony CLI):
php bin/console messenger:consume kafka -vv
Replace a default transport (e.g., async) with Kafka for high-throughput messaging:
framework:
messenger:
transports:
async: ~
kafka:
dsn: '%env(KAFKA_DSN)%'
options:
topic: 'async_messages'
group_id: 'async_workers'
routing:
'App\Message\AsyncTask': kafka
Producer Pattern:
SendEmailMessage to Kafka instead of a database queue.routing:
'App\Message\SendEmail': kafka
Consumer Groups:
workers_*). Each group processes messages independently.group_id: workers_backend for parallel processing.Error Handling:
failed_transport to redirect failed messages to a fallback (e.g., doctrine):transports:
kafka:
dsn: '%env(KAFKA_DSN)%'
failed_transport: doctrine
Serialization:
serializer option:options:
serializer: 'App\Serializer\CustomSerializer'
Laravel-Specific:
Combine with spatie/laravel-messenger for Laravel support:
$this->bus->dispatch(new YourMessage());
Configure in config/messenger.php (Laravel’s adapter layer).
Monitoring:
Use Kafka’s built-in tools (kafka-console-consumer) or integrate with Prometheus via custom metrics.
Retry Logic: Leverage Symfony’s retry middleware:
messenger:
transports:
kafka:
dsn: '%env(KAFKA_DSN)%'
retry_strategy:
max_retries: 3
delay: 1000
Connection Issues:
kafka://host:port) and ensure brokers are reachable.php bin/console debug:container koco_kafka to inspect the transport.Consumer Lag:
Schema Evolution:
SSL/TLS:
kafka+ssl://, ensure certificates are trusted:options:
ssl:
cafile: '/path/to/ca.pem'
certfile: '/path/to/client.pem'
keyfile: '/path/to/client.key'
Enable Verbose Logging:
php bin/console messenger:consume kafka -vv
Look for KafkaTransport logs in var/log/dev.log.
Check Topic Existence:
Use kafka-topics.sh to verify topics:
kafka-topics.sh --list --bootstrap-server localhost:9092
Custom Serializers:
Implement MessageSerializerInterface for non-JSON formats (e.g., Avro):
class AvroSerializer implements MessageSerializerInterface
{
public function serialize($message): string { ... }
public function deserialize(string $data): mixed { ... }
}
Register in config:
options:
serializer: 'App\Serializer\AvroSerializer'
Middleware: Add Kafka-specific middleware (e.g., logging, transformation):
$bus->addMiddleware(new KafkaLoggingMiddleware());
Dynamic Topics: Use a message attribute to route dynamically:
#[Message(['topic' => 'dynamic_topic'])]
class DynamicMessage {}
Requires custom transport extension.
messenger.<transport_name> (e.g., messenger.kafka).options:
group_id: '%env(KAFKA_GROUP_ID)%'
How can I help you explore Laravel packages today?