Custom Committers | Laravel Kafka              

 [ K Laravel Kafka ](/) [Docs](/docs) [Blog](https://laravelkafka.com/blog) [GitHub](https://github.com/mateusjunges/laravel-kafka) 

     Search ⌘K       [    Login with GitHub Login ](https://laravelkafka.com/oauth/github/redirect) 

    Docs for version         

   selected

 v3.0     

 v2.13     

 v2.12     

 v2.11     

 v2.10     

 v2.9     

 v2.8     

 v1.13   

- - [ Introduction ](/docs/v3.0/introduction)
    - [ Requirements ](/docs/v3.0/requirements)
    - [ Installation and Setup ](/docs/v3.0/installation-and-setup)
    - [ Questions and issues ](/docs/v3.0/questions-and-issues)
    - [ Changelog ](/docs/v3.0/changelog)
    - [ Upgrade guide ](/docs/v3.0/upgrade-guide)
    - [ Example docker-compose file ](/docs/v3.0/example-docker-compose)
- Producing messages
    ------------------

    - [ Producing messages ](/docs/v3.0/producing-messages/producing-messages)
    - [ Configuring your kafka producer ](/docs/v3.0/producing-messages/configuring-producers)
    - [ Configuring message payload ](/docs/v3.0/producing-messages/configuring-message-payload)
    - [ Custom serializers ](/docs/v3.0/producing-messages/custom-serializers)
    - [ Publishing to kafka ](/docs/v3.0/producing-messages/publishing-to-kafka)
- Consuming messages
    ------------------

    - [ Creating a kafka consumer ](/docs/v3.0/consuming-messages/creating-consumer)
    - [ Subscribing to kafka topics ](/docs/v3.0/consuming-messages/subscribing-to-kafka-topics)
    - [ Using regex to subscribe to kafka topics ](/docs/v3.0/consuming-messages/using-regex-to-subscribe-to-kafka-topics)
    - [ Assigning consumers to a topic partition ](/docs/v3.0/consuming-messages/assigning-partitions)
    - [ Consuming messages from specific offsets ](/docs/v3.0/consuming-messages/consuming-from-specific-offsets)
    - [ Consumer groups ](/docs/v3.0/consuming-messages/consumer-groups)
    - [ Partition Discovery and Dynamic Assignment ](/docs/v3.0/consuming-messages/partition-discovery)
    - [ Message handlers ](/docs/v3.0/consuming-messages/message-handlers)
    - [ Configuring consumer options ](/docs/v3.0/consuming-messages/configuring-consumer-options)
    - [ Handling failed messages ](/docs/v3.0/consuming-messages/handling-failed-messages)
    - [ Custom deserializers ](/docs/v3.0/consuming-messages/custom-deserializers)
    - [ Consuming messages ](/docs/v3.0/consuming-messages/consuming-messages)
    - [ Pausing partitions ](/docs/v3.0/consuming-messages/pausing-partitions)
    - [ Consumer classes ](/docs/v3.0/consuming-messages/class-structure)
    - [ Using consumers with Laravel Telescope ](/docs/v3.0/consuming-messages/laravel-telescope)
- Advanced usage
    --------------

    - [ Connections ](/docs/v3.0/advanced-usage/connections)
    - [ Replacing the default serializer/deserializer ](/docs/v3.0/advanced-usage/replacing-default-serializer)
    - [ Graceful shutdown ](/docs/v3.0/advanced-usage/graceful-shutdown)
    - [ Running consumers in production ](/docs/v3.0/advanced-usage/running-consumers-in-production)
    - [ SASL Authentication ](/docs/v3.0/advanced-usage/sasl-authentication)
    - [ Custom Committers ](/docs/v3.0/advanced-usage/custom-committers)
    - [ Manual Commit ](/docs/v3.0/advanced-usage/manual-commit)
    - [ Middlewares ](/docs/v3.0/advanced-usage/middlewares)
    - [ Stop consumer when there are no messages left ](/docs/v3.0/advanced-usage/stop-consumer-after-last-message)
    - [ Stop consumer on demand ](/docs/v3.0/advanced-usage/stopping-a-consumer)
    - [ Writing custom loggers ](/docs/v3.0/advanced-usage/custom-loggers)
    - [ Before and after callbacks ](/docs/v3.0/advanced-usage/before-callbacks)
    - [ Setting global configurations ](/docs/v3.0/advanced-usage/setting-global-configuration)
    - [ Sending multiple messages with the same producer ](/docs/v3.0/advanced-usage/sending-multiple-messages-with-the-same-producer)
    - [ Events ](/docs/v3.0/advanced-usage/events)
    - [ Consumer lag ](/docs/v3.0/advanced-usage/consumer-lag)
- Testing
    -------

    - [ Kafka fake ](/docs/v3.0/testing/fake)
    - [ Assert Published ](/docs/v3.0/testing/assert-published)
    - [ Assert published On ](/docs/v3.0/testing/assert-published-on)
    - [ Assert not published ](/docs/v3.0/testing/assert-not-published)
    - [ Assert nothing published ](/docs/v3.0/testing/assert-nothing-published)
    - [ Assert published times ](/docs/v3.0/testing/assert-published-times)
    - [ Assert published on times ](/docs/v3.0/testing/assert-published-on-times)
    - [ Mocking your kafka consumer ](/docs/v3.0/testing/mocking-your-kafka-consumer)

  Custom Committers 
===================

In auto commit mode, the consumer stores the offset of each message after it is processed, and librdkafka commits the stored offsets in the background. Committers are used when handlers commit offsets themselves, by calling the `commit` or `commitAsync` methods of the consumer, usually in [ manual commit ](manual-commit) mode.

  Sponsorship 

Support Laravel Kafka by sponsoring me!

Laravel Kafka is free and Open Source software, built to empower developers like you. Your support helps maintain and enhance the project. If you find it valuable, please consider sponsoring me on GitHub. Every contribution makes a difference and keeps the development going strong! Thank you!

 [   Become a Sponsor ](https://github.com/sponsors/mateusjunges) Want to hide this message? Sponsor at any tier of $10/month or more! 

The `Junges\Kafka\Contracts\Committer` interface has two methods:

- `commit(ConsumerMessage|Message|array|null $messageOrOffsets = null): void`, used for synchronous commits.
- `commitAsync(ConsumerMessage|Message|array|null $messageOrOffsets = null): void`, used for asynchronous commits.

Both receive what the handler passed to the consumer: nothing, to commit the offsets of the current assignment, a `Junges\Kafka\Contracts\ConsumerMessage` or `RdKafka\Message`, or an array of `RdKafka\TopicPartition`.

### [](#content-usage-example "Permalink")Usage example

The following committer retries synchronous commits while the consumer group is rebalancing:

         ```
use Junges\Kafka\Contracts\Committer;
use Junges\Kafka\Contracts\ConsumerMessage;
use RdKafka\Exception;
use RdKafka\KafkaConsumer;
use RdKafka\Message;
use RdKafka\TopicPartition;

class RetryingCommitter implements Committer
{
    public function __construct(private KafkaConsumer $consumer) {}

    public function commit(ConsumerMessage|Message|array|null $messageOrOffsets = null): void
    {
        retry(3, fn () => $this->consumer->commit($this->offsets($messageOrOffsets)), 100, function (Exception $exception) {
            return $exception->getCode() === RD_KAFKA_RESP_ERR_REBALANCE_IN_PROGRESS;
        });
    }

    public function commitAsync(ConsumerMessage|Message|array|null $messageOrOffsets = null): void
    {
        $this->consumer->commitAsync($this->offsets($messageOrOffsets));
    }

    private function offsets(mixed $messageOrOffsets): mixed
    {
        if (! $messageOrOffsets instanceof ConsumerMessage) {
            return $messageOrOffsets;
        }

        return [new TopicPartition(
            $messageOrOffsets->getTopicName(),
            $messageOrOffsets->getPartition(),
            $messageOrOffsets->getOffset() + 1
        )];
    }
}
```

To use it, create a committer factory, which is a class that implements the `Junges\Kafka\Contracts\CommitterFactory` interface, and pass it to the consumer:

         ```
use Junges\Kafka\Config\Config;
use Junges\Kafka\Contracts\Committer;
use Junges\Kafka\Contracts\CommitterFactory;
use Junges\Kafka\Facades\Kafka;
use RdKafka\KafkaConsumer;

class RetryingCommitterFactory implements CommitterFactory
{
    public function make(KafkaConsumer $kafkaConsumer, Config $config): Committer
    {
        return new RetryingCommitter($kafkaConsumer);
    }
}

$consumer = Kafka::consumer(['orders'])
    ->withManualCommit()
    ->usingCommitterFactory(new RetryingCommitterFactory)
    ->withHandler(function ($message, $consumer) {
        // ...
        $consumer->commit($message);
    })
    ->build();
```

Previous

 [    SASL Authentication ](https://laravelkafka.com/docs/v3.0/advanced-usage/sasl-authentication) 

Next

 [ Manual Commit    ](https://laravelkafka.com/docs/v3.0/advanced-usage/manual-commit) 

 Sponsors

 [ version="1.0" encoding="UTF-8"?       EasyCal ](https://easycal.app/) 

 [       Search  ⌘ K   ](https://typesense.org/)
