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/v1.13/introduction)
    - [ Requirements ](/docs/v1.13/requirements)
    - [ Installation and Setup ](/docs/v1.13/installation-and-setup)
    - [ Questions and issues ](/docs/v1.13/questions-and-issues)
    - [ Changelog ](/docs/v1.13/changelog)
    - [ Upgrade Guide ](/docs/v1.13/upgrade-guide)
- Producing messages
    ------------------

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

    - [ Creating a kafka consumer ](/docs/v1.13/consuming-messages/creating-consumer)
    - [ Subscribing to kafka topics ](/docs/v1.13/consuming-messages/subscribing-to-kafka-topics)
    - [ Consumer groups ](/docs/v1.13/consuming-messages/consumer-groups)
    - [ Message handlers ](/docs/v1.13/consuming-messages/message-handlers)
    - [ Configuring consumer options ](/docs/v1.13/consuming-messages/configuring-consumer-options)
    - [ Custom deserializers ](/docs/v1.13/consuming-messages/custom-deserializers)
    - [ Consuming messages ](/docs/v1.13/consuming-messages/consuming-messages)
    - [ Handling message batch ](/docs/v1.13/consuming-messages/handling-message-batch)
    - [ Class structure ](/docs/v1.13/consuming-messages/class-structure)
- Advanced usage
    --------------

    - [ Replacing the default serializer/deserializer ](/docs/v1.13/advanced-usage/replacing-default-serializer)
    - [ Graceful shutdown ](/docs/v1.13/advanced-usage/graceful-shutdown)
    - [ SASL Authentication ](/docs/v1.13/advanced-usage/sasl-authentication)
    - [ Custom Committers ](/docs/v1.13/advanced-usage/custom-committers)
    - [ Middlewares ](/docs/v1.13/advanced-usage/middlewares)
    - [ Stop consumer after last messages ](/docs/v1.13/advanced-usage/stop-consumer-after-last-message)
    - [ Writing custom loggers ](/docs/v1.13/advanced-usage/custom-loggers)
    - [ Setting global configurations ](/docs/v1.13/advanced-usage/setting-global-configuration)
- Testing
    -------

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

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

By default, the committers provided by the `DefaultCommitterFactory` are provided.

To set a custom committer on your consumer, add the committer via a factory that implements the `CommitterFactory` interface:

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

class MyCommitter implements Committer
{
    public function commitMessage(Message $message, bool $success) : void {
        // ...
    }

    public function commitDlq(Message $message) : void {
        // ...
    }  
}

class MyCommitterFactory implements CommitterFactory
{
    public function make(KafkaConsumer $kafkaConsumer, Config $config) : Committer {
        // ...
    }
}

$consumer = \Junges\Kafka\Facades\Kafka::createConsumer()
    ->usingCommitterFactory(new MyCommitterFactory())
    ->build();
```

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

If you want to define a new committer for you consumer, you must start by creating a new class that implements the `Committer` interface. The `commitMessage` function has a `$success` param, which is true for all messages that were consumed without throwing exceptions or messages which exceptions were handled successfully by the consumer class. So, the following committer will commit only messages that were consumed withtout throwing an exception:

         ```
class CustomCommitter implements CommitterContract
{
    public function __construct(private KafkaConsumer $consumer) {}

    public function commitMessage(Message $message, bool $success): void
    {
        if (! $success) {
            return;
        }

        $this->consumer->commit($message);
    }

    public function commitDlq(Message $message): void
    {
        $this->consumer->commit($message);
    }
}
```

After creating your custom committer implementation, you must create a committer factory, which is a simples class that implements the `CommitterFactory` interface, which will be used to provide your custom committer implementation to the consumer class:

         ```
class CustomCommitterFactory implements CommitterFactory
{
    public function make(KafkaConsumer $kafkaConsumer, Config $config): CommitterContract
    {
        return new RetryableCommitter(
            new SuccessCommitter(
                $kafkaConsumer
            ),
            new NativeSleeper(),
            $config->getMaxCommitRetries()
        );
    }
}
```

To use this committer implementation, you just need to inform your consumer that you want to use a custom committer class:

         ```
$consumer = Kafka::createConsumer()
    ->usingCommitterFactory(new CustomCommitterFactory())
    ->build();
```

Previous

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

Next

 [ Middlewares    ](https://laravelkafka.com/docs/v1.13/advanced-usage/middlewares) 

 Sponsors

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

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