Skip to main content

RoundRobinPartitioner

Struct RoundRobinPartitioner 

Source
pub struct RoundRobinPartitioner { /* private fields */ }
Expand description

The “Round-Robin” partitioner.

This partitioning strategy can be used when user wants to distribute the writes to all partitions equally. This is the behaviour regardless of record key hash.

The record key (if present) is ignored: the partitioner cycles through the topic’s partitions in order, choosing only from partitions that currently have a leader (available partitions), so a leaderless partition is skipped.

A single instance is shared by the producer across tasks, so the per-topic counter is an AtomicI32 inside a concurrent map — the direct translation of Java’s ConcurrentMap<String, AtomicInteger>.

Implementations§

Source§

impl RoundRobinPartitioner

Source

pub fn new() -> Self

Creates a new RoundRobinPartitioner.

Trait Implementations§

Source§

impl Default for RoundRobinPartitioner

Source§

fn default() -> RoundRobinPartitioner

Returns the “default value” for a type. Read more
Source§

impl<K, V> Partitioner<K, V> for RoundRobinPartitioner

Source§

fn partition( &self, topic: &str, _key: Option<&K>, _key_bytes: Option<&[u8]>, _value: Option<&V>, _value_bytes: Option<&[u8]>, cluster: &Cluster, ) -> i32

Compute the partition for the given record.

If there are available (leader-having) partitions for the topic, the next counter value picks one of them round-robin; otherwise it falls back to the full partition count. The key/value are ignored.

§Precondition / panics

When the topic has no partitions at all in cluster, the fallback Utils::to_positive(next) % 0 divides by zero — Java’s identical code throws ArithmeticException there. This is unreachable from KafkaProducer: wait_on_metadata guarantees the topic has partitions before partition is called. Per CLAUDE.md §12.1 this class of panic (division by zero) is acceptable, so the i32 return type is preserved rather than made fallible.

Source§

fn configure(&mut self, _configs: &HashMap<String, String>)

Configure this partitioner. Read more
Source§

fn close(&self)

Close this partitioner. Read more

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

§

fn vzip(self) -> V