Skip to main content

ConsumerRecord

Struct ConsumerRecord 

Source
pub struct ConsumerRecord<K, V> { /* private fields */ }
Expand description

A key/value pair received from Kafka.

Corresponds to Java’s org.apache.kafka.clients.consumer.ConsumerRecord<K,V>.

Per consumer-threading.md §27, the topic is an Arc<str> that is shared cheaply across all records from the same topic-partition (via SubscriptionState), and the headers are owned for this milestone (matching Java’s allocation behavior).

§Thread safety

Mirrors Java’s behavior: this struct is not designed for concurrent mutation; Headers are mutable. Concurrent reads are safe by Rust’s borrow rules.

§Equality

PartialEq / Eq are derived (gated on K: PartialEq, V: PartialEq / K: Eq, V: Eq) so that tests can compare two record batches for structural equality. This mirrors Java’s behavior where ConsumerRecord equality is value-based; the Java class itself does not override equals, but its fields are all value types, so two records with identical fields compare equal via Objects.equals. The cost is paid only by callers that opt in to PartialEq types (e.g. tests using i32 keys); users with non-PartialEq K/V continue to work because the bounds are gated by the derive.

Implementations§

Source§

impl<K, V> ConsumerRecord<K, V>

Source

pub const NO_TIMESTAMP: i64 = -1

Sentinel value indicating no timestamp is associated with a record.

Corresponds to Java’s RecordBatch.NO_TIMESTAMP referenced via ConsumerRecord.NO_TIMESTAMP.

Source

pub const NULL_SIZE: i32 = -1

Sentinel value used for serialized_key_size / serialized_value_size when the key/value is None.

Corresponds to Java’s ConsumerRecord.NULL_SIZE.

Source

pub fn new( topic: impl Into<Arc<str>>, partition: i32, offset: i64, key: Option<K>, value: Option<V>, ) -> Self

Creates a record from a specified topic and partition.

Corresponds to Java’s ConsumerRecord(String, int, long, K, V) (ConsumerRecord.java:83), whose parameters {topic, partition, offset, key, value} are the intersection across the three constructors — so it owns the plain name (CLAUDE.md §2).

The timestamp is set to ConsumerRecord::NO_TIMESTAMP, the timestamp type to TimestampType::NoTimestampType, the serialized sizes to ConsumerRecord::NULL_SIZE, headers to an empty RecordHeaders, and both leader_epoch and delivery_count to None.

Source

pub fn with_options(options: ConsumerRecordOptions<K, V>) -> Self

Creates a record with full metadata.

Corresponds to Java’s widest constructor (ConsumerRecord.java:138), which takes deliveryCount alongside every other field. Its twelve parameters exceed CLAUDE.md §2’s three-parameter cap on derived overload names, so ConsumerRecordOptions is this method’s only parameter and carries all of them.

Java’s intermediate 11-arg constructor (ConsumerRecord.java:107) is not a separate Rust method: its body is literally this one with deliveryCount = Optional.empty(), and under CLAUDE.md §2 both derive the same name with_options once the surplus parameters move into ConsumerRecordOptions. Callers get the 11-arg form by leaving ConsumerRecordOptions::delivery_count at None.

  • options - every parameter of Java’s widest constructor
Source

pub fn topic(&self) -> &str

The topic this record is received from (never null).

Source

pub fn partition(&self) -> i32

The partition from which this record is received.

Source

pub fn offset(&self) -> i64

The position of this record in the corresponding Kafka partition.

Source

pub fn timestamp(&self) -> i64

The timestamp of this record, in milliseconds elapsed since unix epoch.

Source

pub fn timestamp_type(&self) -> TimestampType

The timestamp type of this record.

Source

pub fn serialized_key_size(&self) -> i32

The size of the serialized, uncompressed key in bytes. Returns ConsumerRecord::NULL_SIZE (-1) if the key is None.

Source

pub fn serialized_value_size(&self) -> i32

The size of the serialized, uncompressed value in bytes. Returns ConsumerRecord::NULL_SIZE (-1) if the value is None.

Source

pub fn key(&self) -> Option<&K>

The key (or None if no key was specified).

Source

pub fn value(&self) -> Option<&V>

The value (or None if no value was specified).

Source

pub fn headers(&self) -> &RecordHeaders

The headers (never null).

Source

pub fn leader_epoch(&self) -> Option<i32>

Get the leader epoch for the record if available.

Source

pub fn delivery_count(&self) -> Option<i16>

Get the delivery count for the record if available.

Deliveries are counted for records delivered by share groups.

Trait Implementations§

Source§

impl<K, V> Debug for ConsumerRecord<K, V>
where K: Debug, V: Debug,

Source§

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

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

impl<K, V> Display for ConsumerRecord<K, V>
where K: Debug, V: Debug,

Source§

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

Matches Java’s toString():

ConsumerRecord(topic = T, partition = P, leaderEpoch = E, offset = O, TS_TYPE = TS,
  deliveryCount = D, serialized key size = SKS, serialized value size = SVS,
  headers = H, key = K, value = V)
Source§

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

Source§

fn eq(&self, other: &ConsumerRecord<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 ConsumerRecord<K, V>

Source§

impl<K, V> StructuralPartialEq for ConsumerRecord<K, V>

Auto Trait Implementations§

§

impl<K, V> Freeze for ConsumerRecord<K, V>
where K: Freeze, V: Freeze,

§

impl<K, V> RefUnwindSafe for ConsumerRecord<K, V>

§

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

§

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

§

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

§

impl<K, V> UnsafeUnpin for ConsumerRecord<K, V>
where K: UnsafeUnpin, V: UnsafeUnpin,

§

impl<K, V> UnwindSafe for ConsumerRecord<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> ToString for T
where T: Display + ?Sized,

Source§

fn to_string(&self) -> String

Converts the given value to a String. Read more
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