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

  Events 
========

The package dispatches Laravel events while publishing and consuming messages, which you can listen to for logging, monitoring or metrics. All of them are in the `Junges\Kafka\Events` namespace.

  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! 

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

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

Every event about a single message has a `getMessageIdentifier()` method, returning the [ id of the message ](../producing-messages/configuring-message-payload#message-ids), so you can follow a message across events and applications.

### [](#content-producer-events "Permalink")Producer events

EventDispatched whenProperties`PublishingMessage`A message is about to be queued on the producer.`message`, the `ProducerMessage` being published, and `connection`.`MessagePublished`A message was queued on the producer. It is delivered in the background, so this does not mean Kafka received it.`message`, the published `ProducerMessage`, with its serialized body, and `connection`.`MessageDelivered`Kafka acknowledged a queued message.`topic`, `partition`, `offset`, `key`, `messageIdentifier` and `connection`.`MessageDeliveryFailed`A queued message could not be delivered, for instance because its topic does not exist or it was not acknowledged within `message.timeout.ms`.`topic`, `partition`, `key`, `payload`, `headers`, `errorCode`, `error`, `messageIdentifier` and `connection`.`CouldNotPublishMessage`Flushing the producer failed, after its retries. The exception is also thrown, or reported when the flush happens as the application terminates.`errorCode`, `message`, `throwable` and `connection`.The `connection` of the producer events is the name of the connection the message was published on, as configured in `config/kafka.php`, so listeners can tell the clusters of an application apart, or publish a message again on the same one.

Delivery reports are received while the producer publishes or flushes messages, so `MessageDelivered` and `MessageDeliveryFailed` are usually dispatched when the messages are flushed. See [ producing messages ](../producing-messages/producing-messages) for when queued messages are flushed.

### [](#content-consumer-events "Permalink")Consumer events

Every consumer event has a `consumer` property, the `Junges\Kafka\Contracts\Consumer` that dispatched it. Its `getName()`, `getConnectionName()`, `getGroupId()` and `getTopics()` methods tell which consumer it is, see [ naming consumers ](../consuming-messages/class-structure#naming-consumers).

EventDispatched whenOther properties`ConsumerStarting`The consumer starts consuming, before it connects to Kafka.`PartitionsAssigned`Partitions were assigned to the consumer, on a consumer group rebalance.`partitions`, a list of `RdKafka\TopicPartition`.`PartitionsRevoked`Partitions were revoked from the consumer, on a consumer group rebalance.`partitions`, a list of `RdKafka\TopicPartition`.`StartedConsumingMessage`A message was received, before it is deserialized and handled.`message`, the received `ConsumerMessage`, with its raw body.`MessageConsumed`The handler processed a message.`message`, the `ConsumerMessage` the handler received, with its attempt number.`RetryingMessage`The handler failed and the message will be [ retried ](../consuming-messages/handling-failed-messages#retrying-failed-messages), after the backoff.`message`, with the number of the attempt that failed, and `throwable`.`MessageFailed`A message is handled as failed, once its retries are used, before it is sent to the dead letter queue, skipped, or stops the consumer.`message`, the `ConsumerMessage` the handler last received, and `throwable`. A message that can't be deserialized fails right away, and has its raw body.`MessageSentToDLQ`A failed message was sent to the dead letter queue.`message`, the consumed `ConsumerMessage`, `throwable`, `topic`, the dead letter queue, and `payload`, `key` and `headers`, as published to the dead letter queue.`MessageSkipped`A failed message was skipped, because the consumer [ skips failed messages ](../consuming-messages/handling-failed-messages#skipping-failed-messages) and has no dead letter queue.`message`, the `ConsumerMessage` the handler last received, and `throwable`.`OffsetsCommitted`Offsets were committed, either in the background with auto commit, or through `commit()` and `commitAsync()`.`partitions`, a list of `RdKafka\TopicPartition` holding the committed offsets.`OffsetCommitFailed`Offsets could not be committed.`partitions`, `errorCode` and `error`.`ConsumerStopped`The consumer stopped consuming and left its consumer group.`reason`, a `Junges\Kafka\Consumers\StopReason`, and `exception`, the exception that stopped the consumer, if any.`ConsumerStopped` is dispatched whether the consumer stopped normally or because of an exception, which `consume()` throws right after the event. Its `reason` tells why the consumer stopped:

ReasonThe consumer stopped because`StopReason::Requested`It was asked to stop through [ `stopConsuming()` ](stopping-a-consumer).`StopReason::Signal`The process received a `SIGTERM`, `SIGINT` or `SIGQUIT` signal, see [ graceful shutdown ](graceful-shutdown).`StopReason::Restart`Consumers were asked to restart through the [ `kafka:restart-consumers` ](running-consumers-in-production#restarting-consumers-after-deployments) command.`StopReason::Empty`There were no messages left in its partitions, with [ `stopWhenEmpty()` ](stop-consumer-after-last-message).`StopReason::MessageLimit`It handled the number of messages given to `stopAfterMessages()`.`StopReason::TimeLimit`It ran for the number of seconds given to `stopAfterSeconds()`.`StopReason::Failed`An exception was thrown while consuming, for instance by a failed message.Faked consumers dispatch the same message events, along with `ConsumerStarting` and `ConsumerStopped`, see [ mocking your kafka consumer ](../testing/mocking-your-kafka-consumer).

### [](#content-client-events "Permalink")Client events

These events are dispatched for librdkafka reports that are not tied to a single message. They have a `connection` property, the name of the connection, and a `consumer` property, the consumer that reported them, or `null` when they were reported by a producer.

EventDispatched whenOther properties`StatisticsReported`librdkafka reported its statistics, every `statistics.interval.ms`.`statistics`, the decoded [ statistics ](https://github.com/confluentinc/librdkafka/blob/master/STATISTICS.md), including the lag of each partition of a consumer.`KafkaErrorOccurred`librdkafka reported an error, such as brokers that can't be reached. It recovers from most of them by itself.`errorCode` and `error`.Statistics are disabled by default. Enable them by setting the `statistics.interval.ms` option of a connection, in milliseconds. Consumers report them while consuming, and producers while publishing or flushing messages.

librdkafka logs its errors when nothing receives them, so the `KafkaErrorOccurred` event is only dispatched when it has listeners by the time the consumer or producer is created. Register its listeners in the `boot` method of a service provider.

The `onStatistics()`, `onError()`, `onRebalance()` and `onOffsetCommit()` [ configuration callbacks ](../consuming-messages/configuring-consumer-options) are still called when these events are dispatched, so you can use both.

Previous

 [    Sending multiple messages with the same producer ](https://laravelkafka.com/docs/v3.0/advanced-usage/sending-multiple-messages-with-the-same-producer) 

Next

 [ Consumer lag    ](https://laravelkafka.com/docs/v3.0/advanced-usage/consumer-lag) 

 Sponsors

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

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