kovenant / messenger-kafka
Symfony Messenger Kafka Transport
Package info
github.com/kovenant/messenger-kafka
Type:symfony-bundle
pkg:composer/kovenant/messenger-kafka
Requires
- php: ^8.1
- ext-json: *
- psr/http-client: ^1.0
- psr/http-factory: ^1.0
- psr/http-message: ^2.0
- psr/log: ^1.0.1|^2.0 |^3.0
- symfony/config: ^6.4|^7.3|^8.0
- symfony/dependency-injection: ^6.4|^7.3|^8.0
- symfony/http-kernel: ^6.4|^7.3|^8.0
- symfony/messenger: ^6.4|^7.3|^8.0
Requires (Dev)
- friendsofphp/php-cs-fixer: ^3.6
- kwn/php-rdkafka-stubs: ^2.1
- nyholm/psr7: ^1.4
- phpstan/phpstan: ^1.4.6
- symfony/error-handler: ^6.4|^7.3|^8.0
- symfony/framework-bundle: ^6.4|^7.3|^8.0
- symfony/phpunit-bridge: ^6.4|^7.3|^8.0
- symfony/property-access: ^6.4|^7.0|^8.0
- symfony/serializer: ^6.4|^7.3|^8.0
Suggests
- ext-rdkafka: ^4.0; Needed to support Kafka connectivity
- koco/avro-regy: Confluent Schema Registry integration
Provides
None
Conflicts
- jobcloud/messenger-kafka: *
- koco/messenger-kafka: *
Replaces
None
This package is auto-updated.
Last update: 2026-09-22 15:42:33 UTC
README
This bundle aims to provide a simple Kafka transport for Symfony Messenger. Kafka REST Proxy support coming soon.
This is a fork of KonstantinCodes/messenger-kafka
with compatibility updates from Jobcloud.
It requires PHP 8.1 or newer and supports Symfony 6.4, 7.3+ and 8.x.
The Composer package is kovenant/messenger-kafka; the Koco\Kafka namespace and
bundle registration are preserved.
Installation
Install the latest stable release from Packagist:
composer require kovenant/messenger-kafka:^1.0
Bundle registration
This fork has no Symfony Flex recipe. Enable the bundle in your project's
config/bundles.php:
return [ // ... Koco\Kafka\KocoKafkaBundle::class => ['all' => true], ];
Tests
With PHP and ext-rdkafka installed:
composer install vendor/bin/simple-phpunit tests/Unit
The full suite (vendor/bin/simple-phpunit) also requires a disposable Kafka
broker at 127.0.0.1:9092. Tests fail on Symfony deprecations.
Configuration
DSN
Specify a DSN starting with either kafka:// or kafka+ssl://. Multiple brokers are separated by ,.
kafka://my-local-kafka:9092kafka+ssl://my-staging-kafka:9093kafka+ssl://prod-kafka-01:9093,kafka+ssl://prod-kafka-02:9093,kafka+ssl://prod-kafka-03:9093
Example
The configuration options for kafka_conf and topic_conf can be found here.
It is highly recommended to set enable.auto.offset.store to false for consumers. Otherwise, every message will be acknowledged, regardless of any error thrown by the message handlers.
framework: messenger: transports: producer: dsn: '%env(KAFKA_URL)%' # serializer: App\Infrastructure\Messenger\MySerializer options: flushTimeout: 10000 flushRetries: 5 topic: name: 'events' kafka_conf: security.protocol: 'sasl_ssl' ssl.ca.location: '%kernel.project_dir%/config/kafka/ca.pem' sasl.username: '%env(KAFKA_SASL_USERNAME)%' sasl.password: '%env(KAFKA_SASL_PASSWORD)%' sasl.mechanisms: 'SCRAM-SHA-256' consumer: dsn: '%env(KAFKA_URL)%' # serializer: App\Infrastructure\Messenger\MySerializer options: commitAsync: true receiveTimeout: 10000 topic: name: "events" kafka_conf: enable.auto.offset.store: 'false' group.id: 'my-group-id' # should be unique per consumer security.protocol: 'sasl_ssl' ssl.ca.location: '%kernel.project_dir%/config/kafka/ca.pem' sasl.username: '%env(KAFKA_SASL_USERNAME)%' sasl.password: '%env(KAFKA_SASL_PASSWORD)%' sasl.mechanisms: 'SCRAM-SHA-256' max.poll.interval.ms: '45000' topic_conf: auto.offset.reset: 'earliest'
Serializer
You will most likely want to implement your own Serializer. Please see: https://symfony.com/doc/current/messenger.html#serializing-messages
The fields key, headers, and body are available in the decode() and encode() methods.
<?php namespace App\Infrastructure\Messenger; use App\Catalogue\Domain\Model\Event\ProductCreated; use Symfony\Component\Messenger\Envelope; use Symfony\Component\Messenger\Transport\Serialization\SerializerInterface; final class MySerializer implements SerializerInterface { public function decode(array $encodedEnvelope): Envelope { $record = json_decode($encodedEnvelope['body'], true); return new Envelope(new ProductCreated( $record['id'], $record['name'], $record['description'], )); } public function encode(Envelope $envelope): array { /** @var ProductCreated $event */ $event = $envelope->getMessage(); return [ 'key' => $event->getId(), 'headers' => [], 'body' => json_encode([ 'id' => $event->getId(), 'name' => $event->getName(), 'description' => $event->getDescription(), ]), ]; } }
How do I work with Avro?
Same as with the basic example above, you need to build your own serializer.
Within the decode() and encode() you can make use of flix-tech/avro-serde-php.
What about the Confluent Schema Registry?
To connect with Schema Registry and control various settings, you can use this bundle:
$ composer require koco/avro-regy
And configure it to match your setup:
avro_regy: base_uri: '%env(SCHEMA_REGISTRY_URL)%' file_naming_strategy: subject options: register_missing_schemas: true register_missing_subjects: true serializers: catalogue: schema_dir: '%kernel.project_dir%/src/Catalogue/Domain/Model/Event/Avro/' orders: schema_dir: '%kernel.project_dir%/src/Orders/Domain/Model/Event/Avro/' file_naming_strategy: qualified_name options: register_missing_schemas: false register_missing_subjects: false
Please see https://github.com/KonstantinCodes/avro-regy for the full documentation.