Configuring consumer options | 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)

  Configuring consumer options 
==============================

The consumer builder, returned by `Kafka::consumer()`, offers the following configuration options. Options shared by every consumer of a Kafka cluster belong in the [ connection configuration ](../advanced-usage/connections) instead.

  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! 

### [](#content-configuring-a-dead-letter-queue "Permalink")Configuring a dead letter queue

In kafka, a Dead Letter Queue (or DLQ), is a simple kafka topic in the kafka cluster which acts as the destination for messages that were not able to make it to the desired destination due to some error.

To create a `dlq` in this package, you can use the `withDlq` method. If you don't specify the DLQ topic name, it will be created based on the topic you are consuming, adding the `-dlq` suffix to the topic name.

Failed messages are published to the dead letter queue with the producer options of the consumer connection, and flushed right away. If the dead letter queue can't be reached, the consumer stops without committing the offset of the message, so it is not lost.

Without a dead letter queue, the consumer stops when a message fails, without committing its offset. See [ handling failed messages ](handling-failed-messages) for the available options.

         ```
$consumer = \Junges\Kafka\Facades\Kafka::consumer()->subscribe('topic')->withDlq();

//Or, specifying the dlq topic name:
$consumer = \Junges\Kafka\Facades\Kafka::consumer()->subscribe('topic')->withDlq('your-dlq-topic-name')
```

When your message is sent to the dead letter queue, we will add three header keys to containing information about what happened to that message:

- `kafka_throwable_message`: The exception message
- `kafka_throwable_code`: The exception code
- `kafka_throwable_class_name`: The exception class name.

#### [](#content-adding-context-metadata-to-dead-letter-queue "Permalink")Adding context metadata to dead letter queue

Sometimes you need additional information (context) in dead letter queue messages. To enrich DLQ message header with custom metadata (e.g. IDs, correlation keys, retry info), throw an exception that implements the `Junges\Kafka\Contracts\ContextAware` interface.

The consumer will merge into the message headers:

- Original message headers (if any)
- Throwable headers as defined above: 
    - `kafka_throwable_message`
    - `kafka_throwable_code`
    - `kafka_throwable_class_name`
- Normalized context from any `ContextAware` exceptions.

Example custom exception:

         ```
use Junges\Kafka\Contracts\ContextAware;
use RuntimeException;
use Throwable;

class OrderProcessingException extends RuntimeException implements ContextAware
{
    public function __construct(
        private array $context,
        string $message = 'Order processing failed',
        int $code = 0,
        ?Throwable $previous = null,
    ) {
        parent::__construct($message, $code, $previous);
    }

    public function getContext(): array
    {
        return $this->context;
    }
}
```

Using it inside a consumer handler:

         ```
$consumer = \Junges\Kafka\Facades\Kafka::consumer()
    ->subscribe('orders')
    ->withDlq()          // DLQ topic will default to "orders-dlq"
    ->withHandler(function($message) {
        $payload = $message->getBody();

        // Simulate failure
        throw new OrderProcessingException([
            'x-order-id' => (string)($payload['order_id'] ?? 'unknown'),
            'x-user-id' => (string)($payload['user_id'] ?? 'unknown'),
            'x-retry-count' => '3',
        ]);
    })
    ->build();

$consumer->consume();
```

Resulting DLQ headers (example):

         ```
[
  'kafka_throwable_message' => 'Order processing failed',
  'kafka_throwable_code' => 0,
  'kafka_throwable_class_name' => OrderProcessingException::class,
  'x-order-id' => '42',
  'x-user-id' => '7',
  'x-retry-count' => '3',
]
```

 Hot tip!

 Header values must be strings. Arrays/objects/numbers as well as empty string keys are ignored. Any headers on the original message are preserved unless overwritten. 

### [](#content-commit-modes-auto-vs-manual "Permalink")Commit modes: Auto vs Manual

The package supports two commit modes for controlling when message offsets are committed to Kafka. Both deliver every message at least once, since failed messages are never committed unless they are sent to a dead letter queue or skipped on purpose.

#### [](#content-auto-commit-default "Permalink")Auto Commit (Default)

With auto commit, the offset of each message is stored after your handler processes it, and librdkafka commits the stored offsets in the background. This is the default and simplest mode:

         ```
$consumer = \Junges\Kafka\Facades\Kafka::consumer()
    ->withHandler(function ($message, $consumer) {
        // Process your message. Its offset is stored once the handler returns.
    });
```

The `auto_commit` key of the connection defines whether consumers use auto commit, and `withAutoCommit()` overrides it for a single consumer.

#### [](#content-manual-commit "Permalink")Manual Commit

With manual commit, the handler decides when offsets are committed, by calling the `commit` or `commitAsync` methods of the consumer:

         ```
$consumer = \Junges\Kafka\Facades\Kafka::consumer()
    ->withManualCommit()
    ->withHandler(function ($message, $consumer) {
        processMessage($message);

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

Let exceptions propagate from the handler. Catching the exception of a failed message without rethrowing it makes the consumer move on as if it was processed.

See the [ manual commit guide ](../advanced-usage/manual-commit) for the available commit methods and patterns.

### [](#content-stopping-after-a-number-of-messages-or-seconds "Permalink")Stopping after a number of messages or seconds

If you want to consume a limited amount of messages, use the `stopAfterMessages` method, and to consume for a limited amount of time, use the `stopAfterSeconds` method:

         ```
$consumer = \Junges\Kafka\Facades\Kafka::consumer()->stopAfterMessages(100);

$consumer = \Junges\Kafka\Facades\Kafka::consumer()->stopAfterSeconds(3600);
```

To stop once there are no messages left, see [ stopping the consumer when there are no messages left ](../advanced-usage/stop-consumer-after-last-message).

### [](#content-setting-kafka-configuration-options "Permalink")Setting Kafka configuration options

To set configuration options, you can use two methods: `withOptions`, passing an array of option and option value or, using the `withOption method and passing two arguments, the option name and the option value.

         ```
$consumer = \Junges\Kafka\Facades\Kafka::consumer()
    ->withOptions([
        'option-name' => 'option-value'
    ]);
// Or:
$consumer = \Junges\Kafka\Facades\Kafka::consumer()
    ->withOption('option-name', 'option-value');
```

### [](#content-configuration-callbacks "Permalink")Configuration callbacks

librdkafka reports some information through callbacks, which you can register on the consumer builder:

         ```
$consumer = \Junges\Kafka\Facades\Kafka::consumer(['orders'])
    ->onError(function ($kafka, int $error, string $reason) {
        logger()->error("Kafka error: {$reason}");
    })
    ->onLog(function ($kafka, int $level, string $facility, string $message) {
        logger()->debug($message);
    })
    ->onStatistics(function ($kafka, string $json, int $length) {
        // Emitted every "statistics.interval.ms", when that option is set
    })
    ->onOffsetCommit(function ($kafka, int $error, array $partitions) {
        // Called with the result of every commit, including background and asynchronous ones
    });
```

The same methods are available on [ connections ](../advanced-usage/connections), where they apply to the producer and to every consumer of the connection. Each of them holds a single callback, so a consumer callback replaces the one of its connection. To receive these reports in several places, listen to the `StatisticsReported`, `KafkaErrorOccurred`, `PartitionsAssigned`, `PartitionsRevoked`, `OffsetsCommitted` and `OffsetCommitFailed` [ events ](../advanced-usage/events) instead. The `onRebalance` method sets the rebalance callback, see [ partition discovery ](partition-discovery) for simpler ways to react to partition assignments, and `onOAuthBearerTokenRefresh` is described in [ SASL authentication ](../advanced-usage/sasl-authentication).

Previous

 [    Message handlers ](https://laravelkafka.com/docs/v3.0/consuming-messages/message-handlers) 

Next

 [ Handling failed messages    ](https://laravelkafka.com/docs/v3.0/consuming-messages/handling-failed-messages) 

 Sponsors

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

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