Download the PHP package micromus/kafka-bus-commiter without Composer
On this page you can find all versions of the php package micromus/kafka-bus-commiter. It is possible to download/install these versions without Composer. Possible dependencies are resolved automatically.
Download micromus/kafka-bus-commiter
More information about micromus/kafka-bus-commiter
Files in micromus/kafka-bus-commiter
Package kafka-bus-commiter
Short Description This is my package kafka-bus-repeater
License MIT
Homepage https://github.com/micromus/kafka-bus-commiter
Informations about the package kafka-bus-commiter
Kafka Bus Commiter for PHP
A middleware package for micromus/kafka-bus that provides idempotent Kafka message processing. It tracks which messages have already been handled, prevents duplicate processing, and allows limiting the maximum number of read attempts.
How It Works
ConsumerCommiterMiddleware is inserted into the consumer pipeline and handles four outcomes:
- Already committed — if the message was already successfully processed (
commitedAtis notnull), the middleware logs a warning and stops the pipeline. - Max attempts exceeded — if
maxAttemptis set and the attempt count has exceeded it, the middleware logs an error and stops the pipeline. - Successful processing — if both checks pass, the message continues down the pipeline; once handled,
commit()is called to record it as processed. - Handler error — if downstream processing throws an exception, the middleware calls
failed()and rethrows the exception.
Installation
Usage
Basic Example
Implement RepositorySourceInterface to persist message state (for example in a database or Redis), then pass the
middleware into your worker options:
Middleware Options
| Parameter | Type | Default | Description |
|---|---|---|---|
repository |
RepositorySourceInterface |
— | Storage for message state |
logger |
LoggerInterface |
NullLogger |
PSR-3 compatible logger |
maxAttempt |
int |
-1 |
Maximum number of processing attempts. -1 means unlimited |
Implementing the Repository
You need to provide your own implementation of RepositorySourceInterface:
The middleware reads/writes processing state through RepositorySourceInterface.
If you need key derivation from message data, use IdempotencyMessageRepository as an adapter that maps ConsumerMessageInterface to repository keys.
Idempotency Keys
Out of the box the consumer uses Kafka's msgId() (a combination of topic, partition and offset) as the storage key.
That works as long as the same physical message is never replayed under a different offset. As soon as you have
retries, producer-side resends, or cross-cluster mirroring, the same logical event can arrive with a different
msgId() and slip past the duplicate check.
An idempotency key is a stable identifier that the producer attaches to a message so the consumer can recognize
duplicates regardless of where or how they arrive. This package transports the key through the
x-idempotency-key Kafka header (IdempotencyMessageRepository::HEADER_NAME).
Producing: HasIdempotency + ProducerIdempotencyMiddleware
On the producer side you mark your message class with the HasIdempotency interface and return the stable key.
The ProducerIdempotencyMiddleware reads that key and writes it into the x-idempotency-key header before the
message hits the broker.
Wire the middleware into the publisher route for that message class:
Middleware is registered per publisher route, so you opt individual message classes into the header. Messages
that do not implement HasIdempotency pass through the middleware untouched — no header is added.
Consuming: IdempotencyMessageRepository
On the consumer side, plug IdempotencyMessageRepository into ConsumerCommiterMiddleware. It reads the
x-idempotency-key header and builds the storage key as "{header}-{topicName}", so the same idempotency key
in two different topics is still treated as two distinct events. If the header is missing, it falls back to
msgId() so legacy producers keep working.
Picking a key
Pick something that uniquely identifies the business event, not the transport. Good choices are an aggregate
id plus a version (order-42-v3), an outbox row id, or any value the upstream system already treats as unique.
Avoid values that change on retry (timestamps, random UUIDs generated per send attempt) — they defeat the whole
mechanism.
Testing
Changelog
Please see CHANGELOG for more information on what has changed recently.
Contributing
Please see CONTRIBUTING for details.
Security Vulnerabilities
Please review our security policy on how to report security vulnerabilities.
Credits
- Kirill Popkov
- All Contributors
License
The MIT License (MIT). Please see License File for more information.