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:
- Cross-task
wakeup()—consumer.wakeup()is called from another thread while the owning thread blocks inpoll()/position()(e.g.CompletableFuture.runAsync(() -> consumer.wakeup()),PlaintextConsumerTest.java:1501). - In-callback reentrancy — a
ConsumerRebalanceListenercallsconsumer.assign/seek/pause/resume/position/committed/ beginningOffsets/commitfrom insideonPartitionsAssigned/onPartitionsRevokedby capturing theconsumervariable 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
impl ConsumerHandle
Sourcepub fn wakeup(&self)
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.
Sourcepub fn assignment(&self) -> HashSet<TopicPartition>
pub fn assignment(&self) -> HashSet<TopicPartition>
Consumer::assignment via the shared SubscriptionState.
Sourcepub fn subscription(&self) -> HashSet<String>
pub fn subscription(&self) -> HashSet<String>
Consumer::subscription via the shared SubscriptionState.
Sourcepub fn paused(&self) -> HashSet<TopicPartition>
pub fn paused(&self) -> HashSet<TopicPartition>
Consumer::paused via the shared SubscriptionState.
Sourcepub async fn seek_with_offset(
&self,
partition: TopicPartition,
offset: i64,
) -> Result<(), Error>
pub async fn seek_with_offset( &self, partition: TopicPartition, offset: i64, ) -> Result<(), Error>
AsyncKafkaConsumer::seek.
Sourcepub async fn seek_with_offset_and_metadata(
&self,
partition: TopicPartition,
offset_and_metadata: OffsetAndMetadata,
) -> Result<(), Error>
pub async fn seek_with_offset_and_metadata( &self, partition: TopicPartition, offset_and_metadata: OffsetAndMetadata, ) -> Result<(), Error>
Sourcepub async fn seek_to_beginning(
&self,
partitions: &[TopicPartition],
) -> Result<(), Error>
pub async fn seek_to_beginning( &self, partitions: &[TopicPartition], ) -> Result<(), Error>
Sourcepub async fn seek_to_end(
&self,
partitions: &[TopicPartition],
) -> Result<(), Error>
pub async fn seek_to_end( &self, partitions: &[TopicPartition], ) -> Result<(), Error>
Sourcepub async fn position_with_timeout(
&self,
partition: &TopicPartition,
timeout: Duration,
) -> Result<i64, Error>
pub async fn position_with_timeout( &self, partition: &TopicPartition, timeout: Duration, ) -> Result<i64, Error>
Sourcepub async fn committed(
&self,
partitions: &[TopicPartition],
) -> Result<HashMap<TopicPartition, OffsetAndMetadata>, Error>
pub async fn committed( &self, partitions: &[TopicPartition], ) -> Result<HashMap<TopicPartition, OffsetAndMetadata>, Error>
Sourcepub async fn beginning_offsets(
&self,
partitions: &[TopicPartition],
) -> Result<HashMap<TopicPartition, i64>, Error>
pub async fn beginning_offsets( &self, partitions: &[TopicPartition], ) -> Result<HashMap<TopicPartition, i64>, Error>
Sourcepub async fn end_offsets(
&self,
partitions: &[TopicPartition],
) -> Result<HashMap<TopicPartition, i64>, Error>
pub async fn end_offsets( &self, partitions: &[TopicPartition], ) -> Result<HashMap<TopicPartition, i64>, Error>
Sourcepub async fn offsets_for_times(
&self,
timestamps_to_search: HashMap<TopicPartition, i64>,
) -> Result<HashMap<TopicPartition, OffsetAndTimestamp>, Error>
pub async fn offsets_for_times( &self, timestamps_to_search: HashMap<TopicPartition, i64>, ) -> Result<HashMap<TopicPartition, OffsetAndTimestamp>, Error>
Sourcepub async fn commit_sync(&self) -> Result<(), Error>
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).
Sourcepub async fn commit_sync_with_offsets(
&self,
offsets: HashMap<TopicPartition, OffsetAndMetadata>,
) -> Result<(), Error>
pub async fn commit_sync_with_offsets( &self, offsets: HashMap<TopicPartition, OffsetAndMetadata>, ) -> Result<(), Error>
Sourcepub async fn commit_async(&self) -> Result<(), Error>
pub async fn commit_async(&self) -> Result<(), Error>
Consumer::commit_async. Fire-and-forget commit of the
offsets the bg task has consumed.
Sourcepub async fn commit_async_offsets(
&self,
offsets: HashMap<TopicPartition, OffsetAndMetadata>,
) -> Result<(), Error>
pub async fn commit_async_offsets( &self, offsets: HashMap<TopicPartition, OffsetAndMetadata>, ) -> Result<(), Error>
AsyncKafkaConsumer::commit_async_offsets.
Trait Implementations§
Source§impl Clone for ConsumerHandle
impl Clone for ConsumerHandle
Source§fn clone(&self) -> ConsumerHandle
fn clone(&self) -> ConsumerHandle
1.0.0 · Source§fn clone_from(&mut self, source: &Self)
fn clone_from(&mut self, source: &Self)
source. Read more