Skip to main content

ConsumerHandle

Struct ConsumerHandle 

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

A Clone + Send + Sync handle to a consumer that exposes Consumer::wakeup and the reentrant-safe consumer operations, callable from a task or thread other than the one owning the consumer.

No Java class counterpart — it recovers a Java capability. Java’s Consumer reference is itself a freely-shareable, thread-safe reference. Application code relies on this in two ways that a bare Rust &mut self consumer cannot express:

  1. Cross-task wakeup() — consumer.wakeup() is called from another thread while the owning thread blocks in poll() / position() (e.g. CompletableFuture.runAsync(() -> consumer.wakeup()), PlaintextConsumerTest.java:1501).
  2. In-callback reentrancy — a ConsumerRebalanceListener calls consumer.assign/seek/pause/resume/position/committed/ beginningOffsets/commit from inside onPartitionsAssigned / onPartitionsRevoked by capturing the consumer variable in the (anonymous-inner-class) listener (PlaintextConsumerCallbackTest.java).

In Rust the consumer is owned via &mut self for the duration of a blocking call, and Box<dyn Consumer> is not Clone, so neither pattern is expressible with a bare reference. This handle captures only the already-Arc-shared, internally-synchronized consumer state, so both patterns are expressible without unsafe. The user captures the handle into their listener struct — the Rust equivalent of Java capturing the consumer variable.

Obtain one via Consumer::handle. Cheap to clone — clones share the same underlying state.

§Operations

Sync: wakeup, assignment, subscription, paused.

Async (reentrant-safe consumer ops): assign, seek_with_offset, seek_to_beginning, seek_to_end, pause, resume, position, committed, beginning_offsets, end_offsets, offsets_for_times, commit_sync, commit_async.

Lifecycle / ownership operations (poll, subscribe, unsubscribe, close) are intentionally NOT exposed — Java does not invoke these reentrantly from callbacks.

§Concrete async fn, no #[async_trait]

ConsumerHandle is a concrete struct, so its async methods are concrete async fn returning an anonymous future (no Pin<Box<dyn Future>>), per CLAUDE.md §13. None of its methods are on a per-record hot path.

Implementations§

Source§

impl ConsumerHandle

Source

pub fn wakeup(&self)

Fires the consumer’s wakeup() from this handle. Equivalent to calling Consumer::wakeup on the owning consumer, but callable from any task / thread without holding a reference to the consumer.

Source

pub fn assignment(&self) -> HashSet<TopicPartition>

Consumer::assignment via the shared SubscriptionState.

Source

pub fn subscription(&self) -> HashSet<String>

Consumer::subscription via the shared SubscriptionState.

Source

pub fn paused(&self) -> HashSet<TopicPartition>

Consumer::paused via the shared SubscriptionState.

Source

pub async fn assign(&self, partitions: Vec<TopicPartition>) -> Result<(), Error>

Source

pub async fn seek_with_offset( &self, partition: TopicPartition, offset: i64, ) -> Result<(), Error>

AsyncKafkaConsumer::seek.

Source

pub async fn seek_with_offset_and_metadata( &self, partition: TopicPartition, offset_and_metadata: OffsetAndMetadata, ) -> Result<(), Error>

Source

pub async fn seek_to_beginning( &self, partitions: &[TopicPartition], ) -> Result<(), Error>

Source

pub async fn seek_to_end( &self, partitions: &[TopicPartition], ) -> Result<(), Error>

Source

pub async fn pause(&self, partitions: &[TopicPartition]) -> Result<(), Error>

Source

pub async fn resume(&self, partitions: &[TopicPartition]) -> Result<(), Error>

Source

pub async fn position(&self, partition: &TopicPartition) -> Result<i64, Error>

Source

pub async fn position_with_timeout( &self, partition: &TopicPartition, timeout: Duration, ) -> Result<i64, Error>

Source

pub async fn committed( &self, partitions: &[TopicPartition], ) -> Result<HashMap<TopicPartition, OffsetAndMetadata>, Error>

Source

pub async fn beginning_offsets( &self, partitions: &[TopicPartition], ) -> Result<HashMap<TopicPartition, i64>, Error>

Source

pub async fn end_offsets( &self, partitions: &[TopicPartition], ) -> Result<HashMap<TopicPartition, i64>, Error>

Source

pub async fn offsets_for_times( &self, timestamps_to_search: HashMap<TopicPartition, i64>, ) -> Result<HashMap<TopicPartition, OffsetAndTimestamp>, Error>

Source

pub async fn commit_sync(&self) -> Result<(), Error>

Consumer::commit_sync. Commits the offsets the bg task has consumed (Java commitSync() with no offsets — commit allConsumed).

Source

pub async fn commit_sync_with_offsets( &self, offsets: HashMap<TopicPartition, OffsetAndMetadata>, ) -> Result<(), Error>

Source

pub async fn commit_async(&self) -> Result<(), Error>

Consumer::commit_async. Fire-and-forget commit of the offsets the bg task has consumed.

Source

pub async fn commit_async_offsets( &self, offsets: HashMap<TopicPartition, OffsetAndMetadata>, ) -> Result<(), Error>

AsyncKafkaConsumer::commit_async_offsets.

Trait Implementations§

Source§

impl Clone for ConsumerHandle

Source§

fn clone(&self) -> ConsumerHandle

Returns a duplicate of the value. Read more
1.0.0 · Source§

fn clone_from(&mut self, source: &Self)

Performs copy-assignment from source. 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> CloneToUninit for T
where T: Clone,

Source§

unsafe fn clone_to_uninit(&self, dest: *mut u8)

🔬This is a nightly-only experimental API. (clone_to_uninit)
Performs copy-assignment from self to dest. 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> ToOwned for T
where T: Clone,

Source§

type Owned = T

The resulting type after obtaining ownership.
Source§

fn to_owned(&self) -> T

Creates owned data from borrowed data, usually by cloning. Read more
Source§

fn clone_into(&self, target: &mut T)

Uses borrowed data to replace owned data, usually by cloning. 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