Search by

ufo-tech / kafka-messenger

Alex MaistrenkoValchik

Kafka transport for Symfony Messenger, built directly on librdkafka: one DSN per environment, confirmed delivery and partition keys.

Package info

github.com/UFO-Tech/kafka-messenger

Type:symfony-bundle

pkg:composer/ufo-tech/kafka-messenger

Statistics

Installs: 1

Dependents: 1

Suggesters: 0

Stars: 0

Open Issues: 0

1.0.0 2026-09-22 10:23 UTC

This package is auto-updated.

Last update: 2026-09-22 10:44:24 UTC


README

Ukraine

Kafka transport for Symfony Messenger, built directly on librdkafka

About this package

This package lets a Symfony application send and consume Kafka messages through the Messenger component.

A message that Kafka never acknowledged should never look delivered.

License Size package_version fork

Environment Requirements

php_version symfony_version

PHP 8.2 or newer with ext-rdkafka, and Symfony Messenger 7.3 or newer. One tag covers both lines — the suite is run against 7.3 and against 8.1, with the same code. Beyond Messenger itself the package pulls in nothing: psr/log and four Symfony components, no HTTP client, no serializer, no schema registry.

Installation

composer require ufo-tech/kafka-messenger

Symfony Flex registers the bundle. Without Flex, add it to config/bundles.php:

Ufo\KafkaMessenger\UfoKafkaMessengerBundle::class => ['all' => true],

Configuration

framework:
    messenger:
        transports:
            events:
                dsn: '%env(KAFKA_DSN)%'
                options:
                    topic: 'orders.created'
                    flush_timeout: 10000      # how long to wait for a delivery report, ms
                    flush_retries: 2          # how many times to repeat the flush
                    receive_timeout: 10000    # how long to wait for a message, ms
                    commit_async: false       # commit offsets without waiting for the broker
                    kafka_conf:               # everything else goes to librdkafka as it is
                        group.id: 'orders-service'
                        auto.offset.reset: 'earliest'

Any other key stops the transport at startup and the error names the ones it accepts. Values in kafka_conf are turned into strings for you, so a YAML false reaches librdkafka as "false" rather than as an empty value.

A consumer needs group.id: without it there is nowhere to keep the offset, and the transport says so instead of letting librdkafka abort. enable.auto.commit is set to false unless you choose otherwise — with auto-commit on, a background thread moves the offset and neither ack nor reject decides anything any more.

DSN

The whole configuration fits in the connection string:

kafka+sasl+ssl://user:pass@b-1:9098,b-2:9098/orders.created?group.id=orders-service&flush_timeout=5000
Part Becomes
scheme security.protocol
user:pass sasl.username and sasl.password, accepted only on a kafka+sasl… scheme
hosts, comma separated the broker list
path the topic
query key with a dot a librdkafka property, same as kafka_conf
query key without a dot a transport option, checked against the same list as above

What the DSN says wins over the options array, the way it does in Symfony's own transports.

Scheme security.protocol
kafka:// plaintext
kafka+ssl:// ssl
kafka+sasl:// sasl_plaintext
kafka+sasl+ssl:// sasl_ssl

A value you set in kafka_conf always beats the scheme. Mixing schemes in one DSN is an error: security.protocol covers the whole client, not one broker.

Partition key

Messages sharing a key land in the same partition and keep their order:

$bus->dispatch(new OrderPlaced($orderId), [new KafkaKeyStamp($orderId)]);

Batches and unreadable messages

The receiver hands back a batch whenever the worker asks for one: the first message waits for receive_timeout, the rest are collected without waiting, as many as are already buffered. Symfony's worker learned to ask (messenger:consume --fetch-size=10) in 8.1, so on 7.3 and 8.0 every call brings a single message.

A message the serializer cannot decode comes back as an envelope carrying MessageDecodingFailedException, the shape ReceiverInterface asks a transport for. The worker routes it through the usual retry and failure path and acks it, so the offset moves on and one bad message cannot stop a partition.

What the failure transport keeps depends on the Symfony line. From 8.1 its serializers return that exception with the original encoded envelope inside it, so the original body travels with it into the failure transport and messenger:failed:retry decodes it again once the reason is fixed. On 7.3 and 8.0 the serializer throws instead, the exception has nowhere to carry the payload, and the failure transport keeps the exception alone.

How it works

Kafka\Connection is the only place that talks to librdkafka: clients, delivery reports, offsets, rebalances, leaving the group. Transport\KafkaSender and Transport\KafkaReceiver speak only Messenger — serialization, stamps, ack and reject — and share one connection. Kafka\Dsn and Kafka\Options parse the configuration and refuse what they do not recognise.

The transport implements CloseableTransportInterface: when a worker stops, the consumer leaves its group at once and the partitions move to its neighbours immediately instead of after session.timeout.ms.

More from UFO-Tech

This and seventeen other packages published under the ufo-tech vendor on Packagist.