pub trait Partitioner<K, V>: Send + Sync {
// Required method
fn partition(
&self,
topic: &str,
key: Option<&K>,
key_bytes: Option<&[u8]>,
value: Option<&V>,
value_bytes: Option<&[u8]>,
cluster: &Cluster,
) -> i32;
// Provided methods
fn configure(&mut self, _configs: &HashMap<String, String>) { ... }
fn close(&self) { ... }
}Expand description
Partitioner interface: compute the partition for a record.
A Partitioner maps each produced record to a partition of its topic. The
producer holds a single shared instance and consults it on the send path
(see KafkaProducer); the built-in
key-hash / sticky logic is used only when no partitioner is configured.
§Translation notes (Definition-of-Done §7 — justified deviations)
Java’s interface extends Configurable, Closeable and takes Object key /
Object value. Rust has neither Configurable nor Closeable, and the
producer is generic rather than erased, so the shape differs deliberately:
K, Vgenerics translate Java’sObject key/Object value. In the generic Rust producer the key and value have concrete types, so the partitioner receives typed borrows instead ofObject— no downcasts.Send + Syncsupertraits: the producer is shared across tasks, so the boxed partitioner it owns must be shareable. Because these are supertrait bounds,Box<dyn Partitioner<K, V>>is itselfSend + Sync.partition(&self), not&mut self: Java’sKafkaProduceris thread-safe and calls the one shared partitioner instance concurrently, so Java implementations must already be internally synchronized (RoundRobinPartitioneruses aConcurrentMap+AtomicInteger).&selfis the faithful translation and keeps the send path lock-free (DoD §10): any per-record state is the implementer’s job via interior mutability (atomics / a concurrent map), exactly as in Java.close(&self), not&mut self:Producer::closetakes&self, so the partitioner must be closable through a shared reference. Wrapping it in aMutexto obtain&mutwould add a lock to every send for no behavioral gain.configure(&mut self)with a default no-op mirrors the consumer interceptor precedent (ConsumerInterceptor::configure) and is called exactly once, before the instance is boxed and shared with the producer, so&mut selfis available then.- Java’s
Plugin<Partitioner>/Monitorablemetrics wrapper (the javadoc’s “implementMonitorableto register metrics”) has no Rust counterpart and is not translated.
Required Methods§
Sourcefn partition(
&self,
topic: &str,
key: Option<&K>,
key_bytes: Option<&[u8]>,
value: Option<&V>,
value_bytes: Option<&[u8]>,
cluster: &Cluster,
) -> i32
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.
Corresponds to Java’s int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster).
topic— the topic name.key— the typed key to partition on, orNoneif there is no key.key_bytes— the serialized key to partition on, orNoneif there is no key.value— the typed value to partition on, orNone.value_bytes— the serialized value to partition on, orNone.cluster— the current cluster metadata.
The returned partition must be non-negative; the producer rejects a
negative result with an IllegalArgumentError.
Note: on the borrowed zero-copy send path
(KafkaProducer<Vec<u8>, Vec<u8>>::send(ProducerRecord<&[u8], &[u8]>))
the typed key / value are None even when key_bytes / value_bytes
are present, because on that path Java’s record.key() is the same
byte[] as keyBytes; materializing a typed &Vec<u8> from the borrowed
bytes would allocate, violating the zero-copy contract. Partitioners that
need the key/value should read the *_bytes parameters, which are always
supplied when a key/value exists.
Provided Methods§
Sourcefn configure(&mut self, _configs: &HashMap<String, String>)
fn configure(&mut self, _configs: &HashMap<String, String>)
Configure this partitioner.
Corresponds to Java’s Configurable.configure(Map<String, ?> configs).
Called once, before the instance is shared with the producer, with the
producer’s original configuration plus the (possibly generated)
client.id. The default implementation is a no-op.