rossel / rossel-kafka
A ready-to-use PHP library for seamless communication with Rossel's Kafka infrastructure, handling both production and consumption of messages.
Requires
- php: >= 8.2
- enqueue/enqueue-bundle: ^0.10.24
- enqueue/fs: ^0.10.19
- enqueue/rdkafka: ^0.10.20
- phpstan/phpstan: ^2.1
- psr/log: ^3.0
- ramsey/uuid: ^4.7
- symfony/config: ^6.4 || ^7.0
- symfony/console: ^6.4 || ^7.0
- symfony/dependency-injection: ^6.4 || ^7.0
- symfony/http-kernel: ^6.4 || ^7.0
- symfony/options-resolver: ^6.4 || ^7.0
- symfony/process: ^7.4
Requires (Dev)
- phpstan/phpstan-phpunit: ^2.0
- phpunit/phpunit: ^11
Suggests
None
Provides
None
Conflicts
None
Replaces
None
- dev-main
- 0.5.0-beta
- 0.4.0-beta
- 0.4.0-alpha-11
- dev-feat/message-key-and-b2b-topics
- dev-fix/raw-kafka-serializer
- dev-feat/update-environment-variables-names
- dev-feat/kafka_broker_authentication
- dev-feat/unit-tests-and-refactoring
- dev-feat/use-new-kafka-broker
- dev-feat/phpunit-tests
- dev-feat/kafka-config-injection-message-types
- dev-develop
This package is auto-updated.
Last update: 2026-09-25 14:21:40 UTC
README
A ready-to-use PHP library for seamless communication with Rossel's Kafka infrastructure, handling both production and consumption of messages.
Installation
composer require rossel/rossel-kafka
Upgrading from 0.4? Read UPGRADE-0.5.md.
Configuration
rossel_kafka: broker: url: '%env(ROSSEL_KAFKA_BROKER_URL)%' topics: account_api_public_log_output_v1_json_delete: '%env(KAFKA_TOPIC_ACCOUNT_API_PUBLIC_LOG_OUTPUT_V1_JSON_DELETE)%' authentication_api_public_log_output_v1_json_delete: '%env(KAFKA_TOPIC_AUTHENTICATION_API_PUBLIC_LOG_OUTPUT_V1_JSON_DELETE)%' public_dead_letter_inout_v1_json_delete_d30: '%env(KAFKA_TOPIC_PUBLIC_DEAD_LETTER_INOUT_V1_JSON_DELETE_D30)%' erp_subscription_api_public_log_output_v1_json_delete: '%env(KAFKA_TOPIC_ERP_SUBSCRIPTION_API_PUBLIC_LOG_OUTPUT_V1_JSON_DELETE)%' public_inheritance_output_v1_json_delete: '%env(KAFKA_TOPIC_PUBLIC_INHERITANCE_OUTPUT_V1_JSON_DELETE)%' inheritance_api_public_log_output_v1_json_delete: '%env(KAFKA_TOPIC_INHERITANCE_API_PUBLIC_LOG_OUTPUT_V1_JSON_DELETE)%' notification_api_public_log_output_v1_json_delete: '%env(KAFKA_TOPIC_NOTIFICATION_API_PUBLIC_LOG_OUTPUT_V1_JSON_DELETE)%' public_notification_input_v1_json_delete: '%env(KAFKA_TOPIC_PUBLIC_NOTIFICATION_INPUT_V1_JSON_DELETE)%' offer_api_public_log_output_v1_json_delete: '%env(KAFKA_TOPIC_OFFER_API_PUBLIC_LOG_OUTPUT_V1_JSON_DELETE)%' public_offer_output_v1_json_delete: '%env(KAFKA_TOPIC_PUBLIC_OFFER_OUTPUT_V1_JSON_DELETE)%' profile_api_public_log_output_v1_json_delete: '%env(KAFKA_TOPIC_PROFILE_API_PUBLIC_LOG_OUTPUT_V1_JSON_DELETE)%' purchase_api_public_log_output_v1_json_delete: '%env(KAFKA_TOPIC_PURCHASE_API_PUBLIC_LOG_OUTPUT_V1_JSON_DELETE)%' public_log_output_v1_json_delete: '%env(KAFKA_TOPIC_PUBLIC_LOG_OUTPUT_V1_JSON_DELETE)%' public_offer_input_v1_json_delete: '%env(KAFKA_TOPIC_PUBLIC_OFFER_INPUT_V1_JSON_DELETE)%' public_subscription_input_v1_json_delete: '%env(KAFKA_TOPIC_PUBLIC_SUBSCRIPTION_INPUT_V1_JSON_DELETE)%' public_subscription_output_v1_json_delete: '%env(KAFKA_TOPIC_PUBLIC_SUBSCRIPTION_OUTPUT_V1_JSON_DELETE)%' public_contact_input_v1_json_delete: '%env(KAFKA_TOPIC_PUBLIC_CONTACT_INPUT_V1_JSON_DELETE)%' input_api_public_log_output_v1_json_delete: '%env(KAFKA_TOPIC_INPUT_API_PUBLIC_LOG_OUTPUT_V1_JSON_DELETE)%' output_api_public_log_output_v1_json_delete: '%env(KAFKA_TOPIC_OUTPUT_API_PUBLIC_LOG_OUTPUT_V1_JSON_DELETE)%' public_profile_output_v1_json_delete: '%env(KAFKA_TOPIC_PUBLIC_PROFILE_OUTPUT_V1_JSON_DELETE)%' public_b2b_customer_output_v1_json_delete: '%env(KAFKA_TOPIC_PUBLIC_B2B_CUSTOMER_OUTPUT_V1_JSON_DELETE)%' public_b2b_sales_rep_output_v1_json_delete: '%env(KAFKA_TOPIC_PUBLIC_B2B_SALES_REP_OUTPUT_V1_JSON_DELETE)%' public_b2b_reference_item_output_v1_json_delete: '%env(KAFKA_TOPIC_PUBLIC_B2B_REFERENCE_ITEM_OUTPUT_V1_JSON_DELETE)%' producer: app_name: '%env(ROSSEL_KAFKA_PRODUCER_APP_NAME)%'
Each topic key maps to the actual Kafka topic name provided via environment variable.
Broker authentication
Authentication is optional — typically required in production, but not in local development.
All authentication options live under broker.authentication in the bundle config (already wired to env vars by the bundle's default config file).
SASL
Set both username and password to enable SASL:
ROSSEL_KAFKA_BROKER_SASL_USERNAME=my-user ROSSEL_KAFKA_BROKER_SASL_PASSWORD=my-password # Mechanism: PLAIN (default), SCRAM-SHA-256 or SCRAM-SHA-512 ROSSEL_KAFKA_BROKER_SASL_MECHANISM=PLAIN
SSL — CA certificate
Two options, mutually exclusive — the local file takes priority:
From a local file:
ROSSEL_KAFKA_BROKER_SSL_CA_CERTIFICATE_PATH=/path/to/ca.pem
From a URL (downloaded automatically, cached in /tmp):
ROSSEL_KAFKA_BROKER_SSL_CA_CERTIFICATE_URL=https://...
mTLS — client certificate
For mutual TLS, provide the client certificate and private key in addition to the CA certificate above. Each variable accepts either a PEM file path or raw PEM content:
ROSSEL_KAFKA_BROKER_SSL_CLIENT_CERTIFICATE=/path/to/client.crt.pem ROSSEL_KAFKA_BROKER_SSL_CLIENT_KEY=/path/to/client.key.pem # Optional: password protecting the private key ROSSEL_KAFKA_BROKER_SSL_CLIENT_KEY_PASSWORD=secret
Security protocol — auto-selected
The bundle selects the security.protocol rdkafka option automatically based on which credentials are provided:
| CA cert / client cert | SASL | Protocol |
|---|---|---|
| no | no | plaintext |
| yes | no | ssl |
| no | yes | sasl_plaintext |
| yes | yes | sasl_ssl |
Unset or empty variables are treated as
null— no configuration needed for unauthenticated brokers.
Usage
Send a message to a topic
use Rossel\RosselKafka\Enum\Config\Broker\TopicConfigKeys; use Rossel\RosselKafka\Enum\MessageHeaders\Area; use Rossel\RosselKafka\Enum\MessageHeaders\MessageType; use Rossel\RosselKafka\Model\Message; use Rossel\RosselKafka\Model\MessageHeaders; use Rossel\RosselKafka\Service\Connector\KafkaConnector; use Rossel\RosselKafka\Service\KafkaTopicsFetcher; // Via Symfony DI (recommended) $topic = $kafkaTopicsFetcher->get(TopicConfigKeys::KAFKA_TOPIC_PUBLIC_SUBSCRIPTION_INPUT_V1_JSON_DELETE); $message = new Message( headers: new MessageHeaders( area: Area::FRANCE, from: 'my-app', messageType: MessageType::CANCEL_B2C_SUBSCRIPTION, ), body: ['foo' => 'bar'], // Optional Kafka record key. Records with the same key go to the same partition, // so their order is kept. Without a key, records are spread across partitions. key: 'subscription-42', ); $kafkaConnector->send($topic, $message);
Consume messages
Implement ConsumerInterface and tag your service with rossel_kafka.consumer:
use Rossel\RosselKafka\Consumer\ConsumerInterface; use Rossel\RosselKafka\Enum\MessageHeaders\MessageType; use Rossel\RosselKafka\Model\Message; use Rossel\RosselKafka\Model\Topic; final class MyConsumer implements ConsumerInterface { public function supportsTopic(Topic $topic): bool { return TopicConfigKeys::KAFKA_TOPIC_PUBLIC_SUBSCRIPTION_INPUT_V1_JSON_DELETE === $topic->getConfigKey(); } public function supportsMessageType(Message $message): bool { return MessageType::CANCEL_B2C_SUBSCRIPTION === $message->getType(); } public function __invoke(Message $message): void { // handle message // $message->getKey() returns the Kafka record key, or null if the record has none } }
Listen to topics
# Listen to all topics (one process per topic) php bin/console rossel:kafka:listen # Listen to specific topics (comma-separated topic values) php bin/console rossel:kafka:listen --topics=public_subscription_input_v1_json_delete,public_offer_input_v1_json_delete
Fetch topic metadata
use Rossel\RosselKafka\Enum\Config\Broker\TopicConfigKeys; use Rossel\RosselKafka\Enum\MessageHeaders\MessageType; // Get a single topic by config key $topic = $kafkaTopicsFetcher->get(TopicConfigKeys::KAFKA_TOPIC_PUBLIC_OFFER_OUTPUT_V1_JSON_DELETE); // Get all topics $topics = $kafkaTopicsFetcher->getAll(); // Get topics supporting a given message type $topics = $kafkaTopicsFetcher->getByMessageType(MessageType::SYNC_B2C_ERP_OFFERS);