Partition Discovery and Dynamic Assignment | 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)

  Partition Discovery and Dynamic Assignment 
============================================

The Laravel Kafka package provides several methods to discover and work with partition assignments dynamically, which is especially useful when you need to set specific offsets but don't know the partition numbers in advance.

  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-getting-assigned-partitions "Permalink")Getting Assigned Partitions
-------------------------------------------------------------------------------

While consuming, you can retrieve the current partition assignment from the consumer, which is passed to handlers and to the `beforeConsuming` and `afterConsuming` callbacks:

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

$consumer = \Junges\Kafka\Facades\Kafka::consumer(['my-topic'], 'my-group')
    ->withHandler(function (ConsumerMessage $message, Consumer $consumer) {
        // Returns an array of RdKafka\TopicPartition objects
        foreach ($consumer->getAssignedPartitions() as $partition) {
            logger()->debug("Topic: {$partition->getTopic()}, Partition: {$partition->getPartition()}");
        }
    })
    ->build();

$consumer->consume();
```

**Note:** `getAssignedPartitions()` returns an empty array before the consumer starts consuming, until the broker assigns partitions to it, and after it stops consuming. To react to assignments when they happen, use the callbacks described below.

[](#content-partition-assignment-callbacks "Permalink")Partition Assignment Callbacks
-------------------------------------------------------------------------------------

You can set a callback that gets executed whenever partitions are assigned to your consumer:

         ```
$consumer = \Junges\Kafka\Facades\Kafka::consumer(['my-topic'], 'my-group')
    ->onPartitionsAssigned(function ($partitions) {
        echo "Assigned " . count($partitions) . " partitions:\n";

        foreach ($partitions as $partition) {
            echo "- Topic: {$partition->getTopic()}, Partition: {$partition->getPartition()}\n";
        }
    })
    ->withHandler(function ($message) {
        // Handle message
    });
```

This callback is particularly useful for:

- Logging partition assignments for debugging
- Initializing partition-specific resources
- Tracking partition assignment changes during rebalancing

To run code when partitions are taken away from your consumer, use the `onPartitionsRevoked` method. Its callback runs before the partitions are removed from the assignment, so a consumer using [ manual commits ](../advanced-usage/manual-commit) can still commit the offsets of the messages it processed. Both callbacks receive the consumer as their second argument:

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

$consumer = \Junges\Kafka\Facades\Kafka::consumer(['my-topic'], 'my-group')
    ->withManualCommit()
    ->onPartitionsRevoked(function (array $partitions, Consumer $consumer) {
        $consumer->commit();
    })
    ->withHandler(function ($message) {
        // Handle message
    });
```

Every assignment and revocation also dispatches the `PartitionsAssigned` and `PartitionsRevoked` [ events ](../advanced-usage/events).

[](#content-dynamic-partition-assignment-with-offsets "Permalink")Dynamic Partition Assignment with Offsets
-----------------------------------------------------------------------------------------------------------

The most powerful feature is the ability to dynamically assign offsets based on discovered partitions:

         ```
$consumer = \Junges\Kafka\Facades\Kafka::consumer(['my-topic'], 'my-group')
    ->resolveOffsetsUsing(function ($partitions) {
        $partitionsWithOffsets = [];

        foreach ($partitions as $partition) {
            // Set different offset strategies based on partition
            if ($partition->getPartition() === 0) {
                // Start from the beginning for partition 0
                $partition->setOffset(RD_KAFKA_OFFSET_BEGINNING);
            } elseif ($partition->getPartition() === 1) {
                // Start from the end for partition 1
                $partition->setOffset(RD_KAFKA_OFFSET_END);
            } else {
                // Start from a specific offset for other partitions
                $partition->setOffset(1000);
            }

            $partitionsWithOffsets[] = $partition;
        }

        return $partitionsWithOffsets;
    })
    ->withHandler(function ($message) {
        // Handle message
    });
```

[](#content-practical-examples "Permalink")Practical Examples
-------------------------------------------------------------

### [](#content-resume-from-stored-offsets "Permalink")Resume from Stored Offsets

         ```
$consumer = \Junges\Kafka\Facades\Kafka::consumer(['user-events'], 'analytics-group')
    ->resolveOffsetsUsing(function ($partitions) {
        $partitionsWithOffsets = [];

        foreach ($partitions as $partition) {
            // Get stored offset from database or cache
            $storedOffset = Cache::get("kafka_offset_{$partition->getTopic()}_{$partition->getPartition()}", 0);

            $partition->setOffset($storedOffset);
            $partitionsWithOffsets[] = $partition;

            Log::info("Resuming from offset {$storedOffset} for partition {$partition->getPartition()}");
        }

        return $partitionsWithOffsets;
    })
    ->withHandler(function ($message) {
        // Handle message

        // Store current offset
        Cache::put("kafka_offset_{$message->getTopicName()}_{$message->getPartition()}", $message->getOffset() + 1);
    });
```

### [](#content-partition-specific-processing "Permalink")Partition-Specific Processing

         ```
$consumer = \Junges\Kafka\Facades\Kafka::consumer(['orders'], 'order-processor')
    ->onPartitionsAssigned(function ($partitions) {
        // Initialize partition-specific resources
        foreach ($partitions as $partition) {
            $partitionId = $partition->getPartition();

            // Each partition might handle different regions
            initializeRegionProcessor($partitionId);

            Log::info("Initialized processor for region partition {$partitionId}");
        }
    })
    ->withHandler(function ($message) {
        $partitionId = $message->getPartition();
        processOrderForRegion($message, $partitionId);
    });
```

[](#content-manual-assignment "Permalink")Manual Assignment
-----------------------------------------------------------

When you know the partitions to consume in advance, assign them with the `assignPartitions` method instead. Each partition can have its own starting offset:

         ```
$consumer = \Junges\Kafka\Facades\Kafka::consumer(['my-topic'], 'my-group')
    ->assignPartitions([
        new \RdKafka\TopicPartition('my-topic', 0, 100),  // Start from offset 100
        new \RdKafka\TopicPartition('my-topic', 1, RD_KAFKA_OFFSET_END),  // Start from end
    ])
    ->withHandler(function ($message) {
        // Handle message
    });

```

[](#content-important-notes "Permalink")Important Notes
-------------------------------------------------------

1. **Timing**: Partition assignments happen during consumer group rebalancing, which occurs when consumers join or leave the group.
2. **Consumer Groups**: If you're using consumer groups, partition assignments are managed by Kafka's partition assignment strategy. Manual assignments override consumer group behavior.
3. **Rebalancing**: When using `onPartitionsAssigned()`, `onPartitionsRevoked()` or `resolveOffsetsUsing()`, your callbacks will be called every time a rebalance occurs. They can be used together, in any order: the partitions are assigned with the offsets returned by `resolveOffsetsUsing()`, and then passed to the `onPartitionsAssigned()` callback. They rely on the default partition assignment, which `onRebalance()` replaces, so they can't be combined with it. With the cooperative sticky [ rebalance strategy ](consumer-groups#partition-assignment-strategies), partitions are added to and removed from the assignment one rebalance at a time, so the callbacks receive only the partitions that were just assigned.
4. **Error Handling**: Always handle potential errors in your callbacks, as exceptions can disrupt the rebalancing process.
5. **Performance**: Partition assignment callbacks should be fast, as they block the rebalancing process.

This partition discovery functionality solves the common problem of needing to know partition numbers before consumption starts, making it much easier to implement features like offset management, partition-aware processing, and resumable consumers.

Previous

 [    Consumer groups ](https://laravelkafka.com/docs/v3.0/consuming-messages/consumer-groups) 

Next

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

 Sponsors

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

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