Handling failed messages | 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)

  Handling failed messages 
==========================

When a message handler throws an exception, the exception is logged and reported through the Laravel exception handler. What happens to the message next depends on how the consumer is configured.

  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-default-behavior "Permalink")Default behavior
---------------------------------------------------------

Without a dead letter queue, the consumer stops when a message fails, and **the offset of the failed message is not committed**. The consumer is closed, leaving the consumer group, and `consume()` throws a `Junges\Kafka\Exceptions\ConsumerException`. The original exception is available through the `getPrevious` method. The next consumer of the partition, including the same consumer once it is restarted, starts from the failed message, so failed messages are not lost.

         ```
use Junges\Kafka\Exceptions\ConsumerException;

$consumer = \Junges\Kafka\Facades\Kafka::consumer(['orders'])
    ->withGroupId('orders-group')
    ->withHandler(new OrderHandler)
    ->build();

$consumer->consume(); // Throws a ConsumerException when a message fails
```

The consumer process is expected to exit, and to be restarted by a process monitor such as Supervisor. Keep in mind that:

- Messages are delivered at least once. A message can be processed again after a restart, so handlers should be idempotent.
- A message that always fails stops the consumer every time it is consumed, blocking its partition until the cause is fixed. Use a dead letter queue to move such messages out of the way.
- Make sure your process monitor keeps restarting the consumer. Supervisor, for instance, considers a process that exits within `startsecs` seconds of starting as a failed start, and gives up after `startretries` failed starts.

This works both in auto commit mode and in [ manual commit ](../advanced-usage/manual-commit) mode. In auto commit mode, the consumer sets the `enable.auto.offset.store` option to `false` and stores the offset of each message only after it is processed. This keeps librdkafka from committing the offset of a failed message in the background, which it would otherwise do as soon as the message is fetched.

Failed messages can be retried before they are handled as failed. After that, they can also be sent to a dead letter queue, or skipped.

[](#content-retrying-failed-messages "Permalink")Retrying failed messages
-------------------------------------------------------------------------

Failures caused by a temporary problem, such as a dependency that is briefly unavailable, can be retried with the `retryFailedMessages` method. It receives the number of retries and, optionally, the time to wait before each retry in milliseconds:

         ```
$consumer = \Junges\Kafka\Facades\Kafka::consumer(['orders'])
    ->withGroupId('orders-group')
    ->retryFailedMessages(3, backoffInMs: 1000)
    ->withDlq()
    ->withHandler(new OrderHandler)
    ->build();
```

When the handler throws an exception, it is called again with the same message, up to the given number of times. Middlewares run again on every attempt. The `getAttempts` method of the message returns how many times the handler was called with it, including the current call:

         ```
use Junges\Kafka\Contracts\ConsumerMessage;
use Junges\Kafka\Contracts\Consumer;

function (ConsumerMessage $message, Consumer $consumer) {
    if ($message->getAttempts() > 1) {
        logger()->info('Retrying message', ['offset' => $message->getOffset()]);
    }

    // ...
}
```

Once all retries are used, the message is handled as failed: the [ failure callback ](#being-notified-of-failed-messages) is called, and the message is sent to the dead letter queue, stops the consumer, or is skipped, depending on the configuration. Retries also end early when the consumer is asked to stop, for instance by a termination signal.

To wait longer before each retry, pass an array of backoffs instead. Each retry waits for the value at its position, and the last value is used for the remaining retries:

         ```
$consumer = \Junges\Kafka\Facades\Kafka::consumer(['orders'])
    ->retryFailedMessages(5, backoffInMs: [1000, 5000, 10000])
    ->withHandler(new OrderHandler)
    ->build();
```

The consumer waits during the backoff, so no other message is consumed while a message is being retried. Keep the total time spent retrying a message (the number of retries multiplied by the backoff, plus the time the handler takes) well below the `max.poll.interval.ms` consumer option, 5 minutes by default. A consumer that does not poll Kafka within that interval is removed from the consumer group. Longer outages are better handled by a dead letter queue or by letting the consumer stop.

[](#content-sending-failed-messages-to-a-dead-letter-queue "Permalink")Sending failed messages to a dead letter queue
---------------------------------------------------------------------------------------------------------------------

When a dead letter queue is configured with `withDlq`, failed messages are published to the dead letter queue topic before their offsets are committed, and the consumer moves on to the next message instead of stopping. See [ configuring a dead letter queue ](configuring-consumer-options) for details.

[](#content-skipping-failed-messages "Permalink")Skipping failed messages
-------------------------------------------------------------------------

If losing a failed message is acceptable, use the `skipFailedMessages` method. The consumer then commits the offset of the failed message and moves on to the next one, so **the failed message is not consumed again**:

         ```
$consumer = \Junges\Kafka\Facades\Kafka::consumer(['page-views'])
    ->withGroupId('analytics')
    ->skipFailedMessages()
    ->withHandler(new PageViewHandler)
    ->build();
```

A `Junges\Kafka\Events\MessageSkipped` event is dispatched for every skipped message, with the message and the exception that made it fail, so you can monitor them:

         ```
use Illuminate\Support\Facades\Event;
use Junges\Kafka\Events\MessageSkipped;

Event::listen(function (MessageSkipped $event) {
    logger()->warning('Kafka message skipped', [
        'topic' => $event->message->getTopicName(),
        'offset' => $event->message->getOffset(),
        'exception' => $event->throwable->getMessage(),
    ]);
});
```

When a dead letter queue is also configured, failed messages are sent to it instead of being skipped.

[](#content-being-notified-of-failed-messages "Permalink")Being notified of failed messages
-------------------------------------------------------------------------------------------

To run some code when a message is handled as failed, once its retries are used, override the `failed` method of a [ consumer class ](class-structure):

         ```
use Illuminate\Support\Facades\Notification;
use Junges\Kafka\Contracts\ConsumerMessage;
use Throwable;

public function failed(ConsumerMessage $message, Throwable $exception): void
{
    Notification::route('slack', config('services.slack.alerts'))
        ->notify(new KafkaMessageFailed($message, $exception));
}
```

When building a consumer yourself, use the `onMessageFailed` method of the consumer builder:

         ```
$consumer = \Junges\Kafka\Facades\Kafka::consumer(['orders'])
    ->onMessageFailed(function (ConsumerMessage $message, Throwable $exception) {
        // ...
    })
    ->withHandler(new OrderHandler)
    ->build();
```

The callback runs before the message is sent to the dead letter queue, skipped, or stops the consumer, and it can't change what happens to it. An exception thrown by the callback is reported, and the message is handled as usual.

To be notified of the failed messages of every consumer, listen to the `MessageFailed` [ event ](../advanced-usage/events), which is dispatched right before the failure callback is called. The `RetryingMessage` event is dispatched before each retry.

Previous

 [    Configuring consumer options ](https://laravelkafka.com/docs/v3.0/consuming-messages/configuring-consumer-options) 

Next

 [ Custom deserializers    ](https://laravelkafka.com/docs/v3.0/consuming-messages/custom-deserializers) 

 Sponsors

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

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