Skip to main content

ConsumerRecords

Struct ConsumerRecords 

Source
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>

Source

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

Source

pub fn empty() -> Self

Returns an empty ConsumerRecords.

Corresponds to Java’s static ConsumerRecords.empty() / ConsumerRecords.EMPTY.

Source

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

Source

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

Source

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().

Source

pub fn count(&self) -> usize

The number of records across all topics and partitions in this set.

Source

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

Source

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: Debug, V: Debug> Debug for ConsumerRecords<K, V>

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more
Source§

impl<K, V> Default for ConsumerRecords<K, V>

Source§

fn default() -> Self

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

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>

The type of the elements being iterated over.
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>>>

Which kind of iterator are we turning this into?
Source§

fn into_iter(self) -> Self::IntoIter

Creates an iterator from a value. Read more
Source§

impl<K, V> IntoIterator for ConsumerRecords<K, V>

Owned iteration that consumes the ConsumerRecords.

Source§

type Item = ConsumerRecord<K, V>

The type of the elements being iterated over.
Source§

type IntoIter = FlatMap<IntoValues<TopicPartition, Vec<ConsumerRecord<K, V>>>, IntoIter<ConsumerRecord<K, V>>, fn(Vec<ConsumerRecord<K, V>>) -> IntoIter<ConsumerRecord<K, V>>>

Which kind of iterator are we turning this into?
Source§

fn into_iter(self) -> Self::IntoIter

Creates an iterator from a value. Read more
Source§

impl<K: PartialEq, V: PartialEq> PartialEq for ConsumerRecords<K, V>

Source§

fn eq(&self, other: &ConsumerRecords<K, V>) -> bool

Tests for self and other values to be equal, and is used by ==.
1.0.0 · Source§

fn ne(&self, other: &Rhs) -> bool

Tests for !=. The default implementation is almost always sufficient, and should not be overridden without very good reason.
Source§

impl<K: Eq, V: Eq> Eq for ConsumerRecords<K, V>

Source§

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>

§

impl<K, V> Send for ConsumerRecords<K, V>
where K: Send, V: Send,

§

impl<K, V> Sync for ConsumerRecords<K, V>
where K: Sync, V: Sync,

§

impl<K, V> Unpin for ConsumerRecords<K, V>
where K: Unpin, V: Unpin,

§

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> 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
§

impl<Q, K> Equivalent<K> for Q
where Q: Eq + ?Sized, K: Borrow<Q> + ?Sized,

§

fn equivalent(&self, key: &K) -> bool

Checks if this value is equivalent to the given key. Read more
§

impl<Q, K> Equivalent<K> for Q
where Q: Eq + ?Sized, K: Borrow<Q> + ?Sized,

§

fn equivalent(&self, key: &K) -> bool

Compare self to key and return true if they are equal.
§

impl<Q, K> Equivalent<K> for Q
where Q: Eq + ?Sized, K: Borrow<Q> + ?Sized,

§

fn equivalent(&self, key: &K) -> bool

Checks if this value is equivalent to the given key. 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