pub struct ConsumerRecords<K, V> { /* private fields */ }Expand description
A container that holds the list of ConsumerRecord per partition for a
particular topic.
There is one ConsumerRecord list for every topic-partition returned by
a Consumer::poll(...) operation.
Corresponds to Java’s
org.apache.kafka.clients.consumer.ConsumerRecords<K, V>.
Uses an [IndexMap] internally so iteration order matches insertion order
(Java’s ConsumerRecords is documented as iterating per the underlying
map’s iteration order; the test suite assumes ordering matches insertion
when a LinkedHashMap is supplied).
§Equality
PartialEq / Eq are derived (gated on K: PartialEq, V: PartialEq /
K: Eq, V: Eq) so that tests can assert_eq!(actual, expected)
against an entire batch. Mirrors Java’s ConsumerRecordsTest which
uses assertEquals on whole-batch values. The bounds are gated; users
with non-PartialEq keys/values are unaffected.
Implementations§
Source§impl<K, V> ConsumerRecords<K, V>
impl<K, V> ConsumerRecords<K, V>
Sourcepub fn with_next_offsets(
records: IndexMap<TopicPartition, Vec<ConsumerRecord<K, V>>>,
next_offsets: HashMap<TopicPartition, OffsetAndMetadata>,
) -> Self
pub fn with_next_offsets( records: IndexMap<TopicPartition, Vec<ConsumerRecord<K, V>>>, next_offsets: HashMap<TopicPartition, OffsetAndMetadata>, ) -> Self
Create a new ConsumerRecords from per-partition record lists and a
next-offsets map.
Corresponds to Java’s ConsumerRecords(Map, Map)
(ConsumerRecords.java:62). Java’s deprecated ConsumerRecords(Map),
which supplies no next offsets, is not translated (CLAUDE.md §3).
Sourcepub fn empty() -> Self
pub fn empty() -> Self
Returns an empty ConsumerRecords.
Corresponds to Java’s static ConsumerRecords.empty() /
ConsumerRecords.EMPTY.
Sourcepub fn records_partition(
&self,
partition: &TopicPartition,
) -> &[ConsumerRecord<K, V>]
pub fn records_partition( &self, partition: &TopicPartition, ) -> &[ConsumerRecord<K, V>]
Get the records for the given partition.
Returns an empty slice if no records are present for that partition.
Corresponds to Java’s ConsumerRecords.records(TopicPartition).
Sourcepub fn records_topic<'a>(
&'a self,
topic: &'a str,
) -> impl Iterator<Item = &'a ConsumerRecord<K, V>> + 'a
pub fn records_topic<'a>( &'a self, topic: &'a str, ) -> impl Iterator<Item = &'a ConsumerRecord<K, V>> + 'a
Get the records for the given topic, across all partitions in this record set.
Returns an iterator that yields references to records in partition-insertion order.
Corresponds to Java’s ConsumerRecords.records(String).
Sourcepub fn partitions(&self) -> impl Iterator<Item = &TopicPartition>
pub fn partitions(&self) -> impl Iterator<Item = &TopicPartition>
Get the partitions which have records contained in this record set.
Returns an iterator over partition references in insertion order. The set may be empty if no data was returned.
Corresponds to Java’s ConsumerRecords.partitions().
Sourcepub fn count(&self) -> usize
pub fn count(&self) -> usize
The number of records across all topics and partitions in this set.
Sourcepub fn is_empty(&self) -> bool
pub fn is_empty(&self) -> bool
true if this record set has no partitions.
Matches Java’s ConsumerRecords.isEmpty() exactly (Java returns
records.isEmpty(), not count() == 0).
Sourcepub fn next_offsets(&self) -> &HashMap<TopicPartition, OffsetAndMetadata>
pub fn next_offsets(&self) -> &HashMap<TopicPartition, OffsetAndMetadata>
Get the next offsets and metadata for all topic-partitions for which the position has been advanced in this poll call.
Corresponds to Java’s ConsumerRecords.nextOffsets().
Trait Implementations§
Source§impl<K, V> Default for ConsumerRecords<K, V>
impl<K, V> Default for ConsumerRecords<K, V>
Source§impl<'a, K, V> IntoIterator for &'a ConsumerRecords<K, V>
Borrowed iteration over every record across all partitions (in
partition-insertion order).
impl<'a, K, V> IntoIterator for &'a ConsumerRecords<K, V>
Borrowed iteration over every record across all partitions (in partition-insertion order).
Equivalent to Java’s ConsumerRecords.iterator() (which iterates the
concatenated values of the underlying map).
Source§type Item = &'a ConsumerRecord<K, V>
type Item = &'a ConsumerRecord<K, V>
Source§type IntoIter = FlatMap<Values<'a, TopicPartition, Vec<ConsumerRecord<K, V>>>, Iter<'a, ConsumerRecord<K, V>>, fn(&'a Vec<ConsumerRecord<K, V>>) -> Iter<'a, ConsumerRecord<K, V>>>
type IntoIter = FlatMap<Values<'a, TopicPartition, Vec<ConsumerRecord<K, V>>>, Iter<'a, ConsumerRecord<K, V>>, fn(&'a Vec<ConsumerRecord<K, V>>) -> Iter<'a, ConsumerRecord<K, V>>>
Source§impl<K, V> IntoIterator for ConsumerRecords<K, V>
Owned iteration that consumes the ConsumerRecords.
impl<K, V> IntoIterator for ConsumerRecords<K, V>
Owned iteration that consumes the ConsumerRecords.
Source§type Item = ConsumerRecord<K, V>
type Item = ConsumerRecord<K, V>
Source§type IntoIter = FlatMap<IntoValues<TopicPartition, Vec<ConsumerRecord<K, V>>>, IntoIter<ConsumerRecord<K, V>>, fn(Vec<ConsumerRecord<K, V>>) -> IntoIter<ConsumerRecord<K, V>>>
type IntoIter = FlatMap<IntoValues<TopicPartition, Vec<ConsumerRecord<K, V>>>, IntoIter<ConsumerRecord<K, V>>, fn(Vec<ConsumerRecord<K, V>>) -> IntoIter<ConsumerRecord<K, V>>>
impl<K: Eq, V: Eq> Eq for ConsumerRecords<K, V>
impl<K, V> StructuralPartialEq for ConsumerRecords<K, V>
Auto Trait Implementations§
impl<K, V> Freeze for ConsumerRecords<K, V>
impl<K, V> RefUnwindSafe for ConsumerRecords<K, V>where
K: RefUnwindSafe,
V: RefUnwindSafe,
impl<K, V> Send for ConsumerRecords<K, V>
impl<K, V> Sync for ConsumerRecords<K, V>
impl<K, V> Unpin for ConsumerRecords<K, V>
impl<K, V> UnsafeUnpin for ConsumerRecords<K, V>
impl<K, V> UnwindSafe for ConsumerRecords<K, V>where
K: UnwindSafe,
V: UnwindSafe,
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
§impl<Q, K> Equivalent<K> for Q
impl<Q, K> Equivalent<K> for Q
§fn equivalent(&self, key: &K) -> bool
fn equivalent(&self, key: &K) -> bool
§impl<Q, K> Equivalent<K> for Q
impl<Q, K> Equivalent<K> for Q
§fn equivalent(&self, key: &K) -> bool
fn equivalent(&self, key: &K) -> bool
key and return true if they are equal.