Middlewares | 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)

  Middlewares 
=============

Middlewares provide a convenient way to inspect, filter or transform your Kafka messages before they reach the handler. A middleware receives the message and the next step of the pipeline, and whatever it passes to `$next` is what the handler receives. Not calling `$next` skips the handler for that message, which is then considered handled.

  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-writing-middlewares "Permalink")Writing middlewares

A middleware can be a closure:

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

function (ConsumerMessage $message, callable $next) {
    // Perform some work here
    return $next($message);
}
```

Or a class implementing the `Junges\Kafka\Contracts\Middleware` interface. Middleware classes given by name are resolved from the service container, so their constructor can receive any dependency they need. They are resolved once, when the consumer handles its first message, and the same instance handles every message, so don't keep state about a single message in their properties:

         ```
use Illuminate\Log\LogManager;
use Junges\Kafka\Contracts\ConsumerMessage;
use Junges\Kafka\Contracts\Middleware;

class LogMessages implements Middleware
{
    public function __construct(private LogManager $log) {}

    public function __invoke(ConsumerMessage $message, callable $next): mixed
    {
        $result = $next($message);

        $this->log->info('Kafka message handled', [
            'topic' => $message->getTopicName(),
            'offset' => $message->getOffset(),
        ]);

        return $result;
    }
}
```

### [](#content-registering-middlewares "Permalink")Registering middlewares

In a [ consumer class ](../consuming-messages/class-structure), return the middlewares from the `middleware` method:

         ```
public function middleware(): array
{
    return [
        LogMessages::class,
        new RequireHeader('tenant'),
        function (ConsumerMessage $message, callable $next) {
            return $next($message);
        },
    ];
}
```

When building a consumer yourself, use the `withMiddleware` method, once for each middleware:

         ```
$consumer = \Junges\Kafka\Facades\Kafka::consumer(['orders'])
    ->withMiddleware(LogMessages::class)
    ->withMiddleware(new RequireHeader('tenant'))
    ->withHandler($handler);
```

Middlewares run in the order they are registered, so the first one wraps all the others. When failed messages are [ retried ](../consuming-messages/handling-failed-messages), the middlewares run again on every attempt, and an exception thrown by a middleware makes the message fail like one thrown by the handler.

### [](#content-global-middlewares "Permalink")Global middlewares

To run middlewares for every consumer, register them with the `consumerMiddleware` method of the `Kafka` facade, usually in the `boot` method of a service provider:

         ```
use Junges\Kafka\Facades\Kafka;

public function boot(): void
{
    Kafka::consumerMiddleware([
        LogMessages::class,
        SetTenantFromHeaders::class,
    ]);
}
```

Global middlewares apply to every consumer created through the `Kafka` facade, including [ consumer classes ](../consuming-messages/class-structure), on any connection. They run before the middlewares of each consumer. They are kept when using `Kafka::fake()`, so your tests go through them as well.

Previous

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

Next

 [ Stop consumer when there are no messages left    ](https://laravelkafka.com/docs/v3.0/advanced-usage/stop-consumer-after-last-message) 

 Sponsors

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

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