Skip to main content

Consumer

Trait Consumer 

Source
pub trait Consumer<K, V>: Send + 'static
where K: Send + 'static, V: Send + 'static,
{
Show 47 methods // Required methods fn assignment(&self) -> HashSet<TopicPartition>; fn subscription(&self) -> HashSet<String>; fn paused(&self) -> HashSet<TopicPartition>; fn group_metadata(&self) -> Arc<dyn ConsumerGroupMetadata>; fn client_id(&self) -> &str; fn current_lag(&self, topic_partition: &TopicPartition) -> Option<i64>; fn metrics(&self) -> HashMap<MetricName, Arc<KafkaMetric>>; fn subscribe_with_topics<'life0, 'async_trait>( &'life0 mut self, topics: Vec<String>, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait; fn subscribe_with_topics_listener<'life0, 'async_trait>( &'life0 mut self, topics: Vec<String>, listener: Arc<dyn ConsumerRebalanceListener>, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait; fn subscribe_with_pattern<'life0, 'async_trait>( &'life0 mut self, pattern: SubscriptionPattern, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait; fn subscribe_with_pattern_listener<'life0, 'async_trait>( &'life0 mut self, pattern: SubscriptionPattern, listener: Arc<dyn ConsumerRebalanceListener>, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait; fn assign<'life0, 'async_trait>( &'life0 mut self, partitions: Vec<TopicPartition>, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait; fn unsubscribe<'life0, 'async_trait>( &'life0 mut self, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait; fn poll<'life0, 'async_trait>( &'life0 mut self, timeout: Duration, ) -> Pin<Box<dyn Future<Output = Result<ConsumerRecords<K, V>, Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait; fn commit_sync<'life0, 'async_trait>( &'life0 mut self, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait; fn commit_sync_with_timeout<'life0, 'async_trait>( &'life0 mut self, timeout: Duration, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait; fn commit_sync_with_offsets<'life0, 'async_trait>( &'life0 mut self, offsets: HashMap<TopicPartition, OffsetAndMetadata>, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait; fn commit_sync_with_offsets_timeout<'life0, 'async_trait>( &'life0 mut self, offsets: HashMap<TopicPartition, OffsetAndMetadata>, timeout: Duration, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait; fn commit_async<'life0, 'async_trait>( &'life0 mut self, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait; fn commit_async_with_callback<'life0, 'async_trait>( &'life0 mut self, callback: Arc<dyn OffsetCommitCallback>, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait; fn commit_async_with_offsets_callback<'life0, 'async_trait>( &'life0 mut self, offsets: HashMap<TopicPartition, OffsetAndMetadata>, callback: Arc<dyn OffsetCommitCallback>, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait; fn seek_with_offset<'life0, 'async_trait>( &'life0 mut self, partition: TopicPartition, offset: i64, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait; fn seek_with_offset_and_metadata<'life0, 'async_trait>( &'life0 mut self, partition: TopicPartition, offset_and_metadata: OffsetAndMetadata, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait; fn seek_to_beginning<'life0, 'life1, 'async_trait>( &'life0 mut self, partitions: &'life1 [TopicPartition], ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait; fn seek_to_end<'life0, 'life1, 'async_trait>( &'life0 mut self, partitions: &'life1 [TopicPartition], ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait; fn position<'life0, 'life1, 'async_trait>( &'life0 mut self, partition: &'life1 TopicPartition, ) -> Pin<Box<dyn Future<Output = Result<i64, Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait; fn position_with_timeout<'life0, 'life1, 'async_trait>( &'life0 mut self, partition: &'life1 TopicPartition, timeout: Duration, ) -> Pin<Box<dyn Future<Output = Result<i64, Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait; fn committed<'life0, 'life1, 'async_trait>( &'life0 mut self, partitions: &'life1 [TopicPartition], ) -> Pin<Box<dyn Future<Output = Result<HashMap<TopicPartition, OffsetAndMetadata>, Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait; fn committed_with_timeout<'life0, 'life1, 'async_trait>( &'life0 mut self, partitions: &'life1 [TopicPartition], timeout: Duration, ) -> Pin<Box<dyn Future<Output = Result<HashMap<TopicPartition, OffsetAndMetadata>, Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait; fn partitions_for<'life0, 'life1, 'async_trait>( &'life0 mut self, topic: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<Vec<PartitionInfo>, Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait; fn partitions_for_with_timeout<'life0, 'life1, 'async_trait>( &'life0 mut self, topic: &'life1 str, timeout: Duration, ) -> Pin<Box<dyn Future<Output = Result<Vec<PartitionInfo>, Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait; fn list_topics<'life0, 'async_trait>( &'life0 mut self, ) -> Pin<Box<dyn Future<Output = Result<HashMap<String, Vec<PartitionInfo>>, Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait; fn list_topics_with_timeout<'life0, 'async_trait>( &'life0 mut self, timeout: Duration, ) -> Pin<Box<dyn Future<Output = Result<HashMap<String, Vec<PartitionInfo>>, Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait; fn offsets_for_times<'life0, 'async_trait>( &'life0 mut self, timestamps_to_search: HashMap<TopicPartition, i64>, ) -> Pin<Box<dyn Future<Output = Result<HashMap<TopicPartition, OffsetAndTimestamp>, Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait; fn offsets_for_times_with_timeout<'life0, 'async_trait>( &'life0 mut self, timestamps_to_search: HashMap<TopicPartition, i64>, timeout: Duration, ) -> Pin<Box<dyn Future<Output = Result<HashMap<TopicPartition, OffsetAndTimestamp>, Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait; fn beginning_offsets<'life0, 'life1, 'async_trait>( &'life0 mut self, partitions: &'life1 [TopicPartition], ) -> Pin<Box<dyn Future<Output = Result<HashMap<TopicPartition, i64>, Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait; fn beginning_offsets_with_timeout<'life0, 'life1, 'async_trait>( &'life0 mut self, partitions: &'life1 [TopicPartition], timeout: Duration, ) -> Pin<Box<dyn Future<Output = Result<HashMap<TopicPartition, i64>, Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait; fn end_offsets<'life0, 'life1, 'async_trait>( &'life0 mut self, partitions: &'life1 [TopicPartition], ) -> Pin<Box<dyn Future<Output = Result<HashMap<TopicPartition, i64>, Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait; fn end_offsets_with_timeout<'life0, 'life1, 'async_trait>( &'life0 mut self, partitions: &'life1 [TopicPartition], timeout: Duration, ) -> Pin<Box<dyn Future<Output = Result<HashMap<TopicPartition, i64>, Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait; fn pause<'life0, 'life1, 'async_trait>( &'life0 mut self, partitions: &'life1 [TopicPartition], ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait; fn resume<'life0, 'life1, 'async_trait>( &'life0 mut self, partitions: &'life1 [TopicPartition], ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait; fn enforce_rebalance<'life0, 'async_trait>( &'life0 mut self, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait; fn enforce_rebalance_with_reason<'life0, 'life1, 'async_trait>( &'life0 mut self, reason: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait; fn close<'life0, 'async_trait>( &'life0 mut self, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait; fn close_with_options<'life0, 'async_trait>( &'life0 mut self, options: CloseOptions, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait; fn wakeup(&self); fn handle(&self) -> ConsumerHandle;
}
Expand description

The single dispatch trait that MockConsumer (Phase 3) and AsyncKafkaConsumer (Phase 11) both implement. Translates Java’s org.apache.kafka.clients.consumer.Consumer<K, V> interface.

§Parameter conventions

The trait uses three argument shapes deliberately:

  • Owned collections (Vec<T>, HashMap<K, V>): the implementation stores or forwards the input long-term (subscription state, request payload). Ownership transfer avoids a per-element clone.

  • Borrowed slices / maps (&[T], &HashMap<K, V>): the implementation iterates but does not retain the input. Callers can pass &Vec<T>, &[T; N], or any slice without conversion.

  • Borrowed scalars (&TopicPartition, &str): read-only access to a single value.

Methods that take Vec<T> are the ones that consume the input; methods that take &[T] only iterate it. This rule is mechanical: if the impl retains, it owns; if the impl reads, it borrows.

§Bounds: Send + 'static, NOT Send + Sync

The trait is Send + 'static so it can be stored as Box<dyn Consumer<K, V>> and moved between tokio tasks (required for multi-thread runtime support). Sync is intentionally NOT required because the API is &mut self — only one task can call methods at a time, no shared &Consumer reference exists.

Users who need cross-task sharing wrap in Arc<Mutex<dyn Consumer>>, which works without Sync on the trait itself. Dropping Sync lets users plug in K/V types that are Send but not Sync (e.g. types containing Cell) without artificial restrictions.

§Async surface

Per consumer-threading.md §1: methods that block in Java become async fn; non-blocking accessors stay fn. Per CLAUDE.md §13, the per-call Box<Future> cost of #[async_trait] is amortized over many records (typical poll() granularity is ≤ ~100 calls/sec) and is acceptable on this top-level dispatch trait.

§Methods NOT translated

  • registerMetricForSubscription, unregisterMetricFromSubscription, clientInstanceId(Duration): KIP-714 broker-push telemetry, deferred to its own milestone (see Milestone-9 metrics plan, “Out of scope”). They are omitted from the trait (no stub); when KIP-714 telemetry lands, the crate::producer::Producer-style traits and Consumer gain these methods together.
  • metrics() IS translated (Phase M7) — it is in Java’s Consumer interface and snapshots the registry that the metrics managers populate. See Consumer::metrics.
  • subscribe(Pattern, [ConsumerRebalanceListener]) (Java regex): superseded by the SubscriptionPattern variants which match server-side regex semantics.
  • poll(long timeoutMs), close(Duration timeout) (both @Deprecated in Java): not translated per CLAUDE.md §3.

Required Methods§

Source

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

Translates Java’s Set<TopicPartition> assignment().

Source

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

Translates Java’s Set<String> subscription().

Source

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

Translates Java’s Set<TopicPartition> paused().

Source

fn group_metadata(&self) -> Arc<dyn ConsumerGroupMetadata>

Translates Java’s ConsumerGroupMetadata groupMetadata().

ConsumerGroupMetadata is a trait (see its docs), so the consumer returns its own implementation behind an Arc, which is cheap to clone and to hand to Producer::send_offsets_to_transaction.

Source

fn client_id(&self) -> &str

Returns the consumer’s client.id. Borrowed per CLAUDE.md §14.

Source

fn current_lag(&self, topic_partition: &TopicPartition) -> Option<i64>

Translates Java’s OptionalLong currentLag(TopicPartition).

Returns Option<i64> — the natural Rust analog.

Source

fn metrics(&self) -> HashMap<MetricName, Arc<KafkaMetric>>

Translates Java’s Map<MetricName, ? extends Metric> metrics().

Returns a snapshot of all metrics maintained by the consumer, keyed by MetricName. The value type is Arc<KafkaMetric> — KafkaMetric implements the Metric read interface, mirroring Java’s ? extends Metric wildcard. Read each metric’s name via Metric::metric_name and its current value via Metric::metric_value.

Sync — Java’s metrics() does not block. The returned map is a point-in-time snapshot taken under the registry lock (a cold, monitoring-frequency call), not a live view.

This trait method has no default: it is in Java’s Consumer interface, so every implementation provides it, and adding it without a default is acceptable for this pre-1.0 dispatch trait (consumer-threading.md §2). It MATCHES Java’s public surface — it is not a Rust-only addition.

Source

fn subscribe_with_topics<'life0, 'async_trait>( &'life0 mut self, topics: Vec<String>, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Translates Java’s void subscribe(Collection<String> topics).

Takes Vec<String> because the impl moves the elements into SubscriptionState.

Source

fn subscribe_with_topics_listener<'life0, 'async_trait>( &'life0 mut self, topics: Vec<String>, listener: Arc<dyn ConsumerRebalanceListener>, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Translates Java’s void subscribe(Collection<String> topics, ConsumerRebalanceListener).

Source

fn subscribe_with_pattern<'life0, 'async_trait>( &'life0 mut self, pattern: SubscriptionPattern, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Translates Java’s void subscribe(SubscriptionPattern pattern).

Server-side regex subscription (KIP-848 RE2/J): the pattern is sent to the group coordinator, which evaluates it. Java’s javadoc notes that no validation of the pattern is performed by the client.

Source

fn subscribe_with_pattern_listener<'life0, 'async_trait>( &'life0 mut self, pattern: SubscriptionPattern, listener: Arc<dyn ConsumerRebalanceListener>, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Translates Java’s void subscribe(SubscriptionPattern pattern, ConsumerRebalanceListener).

Source

fn assign<'life0, 'async_trait>( &'life0 mut self, partitions: Vec<TopicPartition>, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Translates Java’s void assign(Collection<TopicPartition>).

Async because Java’s assign calls applicationEventHandler.addAndGet(new AssignmentChangeEvent(...)) which blocks (AsyncKafkaConsumer.java:1819). The Rust translation .awaits the event handle.

Source

fn unsubscribe<'life0, 'async_trait>( &'life0 mut self, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Translates Java’s void unsubscribe().

Source

fn poll<'life0, 'async_trait>( &'life0 mut self, timeout: Duration, ) -> Pin<Box<dyn Future<Output = Result<ConsumerRecords<K, V>, Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Translates Java’s ConsumerRecords<K, V> poll(Duration timeout).

Source

fn commit_sync<'life0, 'async_trait>( &'life0 mut self, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Translates Java’s void commitSync().

Source

fn commit_sync_with_timeout<'life0, 'async_trait>( &'life0 mut self, timeout: Duration, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Translates Java’s void commitSync(Duration timeout).

Source

fn commit_sync_with_offsets<'life0, 'async_trait>( &'life0 mut self, offsets: HashMap<TopicPartition, OffsetAndMetadata>, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Translates Java’s void commitSync(Map<TopicPartition, OffsetAndMetadata> offsets).

Source

fn commit_sync_with_offsets_timeout<'life0, 'async_trait>( &'life0 mut self, offsets: HashMap<TopicPartition, OffsetAndMetadata>, timeout: Duration, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Translates Java’s void commitSync(Map<TopicPartition, OffsetAndMetadata> offsets, Duration timeout).

Source

fn commit_async<'life0, 'async_trait>( &'life0 mut self, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Translates Java’s void commitAsync().

Source

fn commit_async_with_callback<'life0, 'async_trait>( &'life0 mut self, callback: Arc<dyn OffsetCommitCallback>, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Translates Java’s void commitAsync(OffsetCommitCallback).

Source

fn commit_async_with_offsets_callback<'life0, 'async_trait>( &'life0 mut self, offsets: HashMap<TopicPartition, OffsetAndMetadata>, callback: Arc<dyn OffsetCommitCallback>, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Translates Java’s void commitAsync(Map<TopicPartition, OffsetAndMetadata>, OffsetCommitCallback).

Source

fn seek_with_offset<'life0, 'async_trait>( &'life0 mut self, partition: TopicPartition, offset: i64, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Translates Java’s void seek(TopicPartition partition, long offset).

Returns Result because Java throws IllegalArgumentException / IllegalStateException on invalid input. Async because Java’s seek calls applicationEventHandler.addAndGet(new SeekUnvalidatedEvent(...)) which blocks (AsyncKafkaConsumer.java:1068).

Source

fn seek_with_offset_and_metadata<'life0, 'async_trait>( &'life0 mut self, partition: TopicPartition, offset_and_metadata: OffsetAndMetadata, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Translates Java’s void seek(TopicPartition partition, OffsetAndMetadata).

Source

fn seek_to_beginning<'life0, 'life1, 'async_trait>( &'life0 mut self, partitions: &'life1 [TopicPartition], ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Translates Java’s void seekToBeginning(Collection<TopicPartition>).

Source

fn seek_to_end<'life0, 'life1, 'async_trait>( &'life0 mut self, partitions: &'life1 [TopicPartition], ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Translates Java’s void seekToEnd(Collection<TopicPartition>).

Source

fn position<'life0, 'life1, 'async_trait>( &'life0 mut self, partition: &'life1 TopicPartition, ) -> Pin<Box<dyn Future<Output = Result<i64, Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Translates Java’s long position(TopicPartition).

Source

fn position_with_timeout<'life0, 'life1, 'async_trait>( &'life0 mut self, partition: &'life1 TopicPartition, timeout: Duration, ) -> Pin<Box<dyn Future<Output = Result<i64, Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Translates Java’s long position(TopicPartition, Duration).

Source

fn committed<'life0, 'life1, 'async_trait>( &'life0 mut self, partitions: &'life1 [TopicPartition], ) -> Pin<Box<dyn Future<Output = Result<HashMap<TopicPartition, OffsetAndMetadata>, Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Translates Java’s Map<TopicPartition, OffsetAndMetadata> committed(Set<TopicPartition>).

Source

fn committed_with_timeout<'life0, 'life1, 'async_trait>( &'life0 mut self, partitions: &'life1 [TopicPartition], timeout: Duration, ) -> Pin<Box<dyn Future<Output = Result<HashMap<TopicPartition, OffsetAndMetadata>, Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Translates Java’s Map<TopicPartition, OffsetAndMetadata> committed(Set<TopicPartition>, Duration).

Source

fn partitions_for<'life0, 'life1, 'async_trait>( &'life0 mut self, topic: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<Vec<PartitionInfo>, Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Translates Java’s List<PartitionInfo> partitionsFor(String topic).

Source

fn partitions_for_with_timeout<'life0, 'life1, 'async_trait>( &'life0 mut self, topic: &'life1 str, timeout: Duration, ) -> Pin<Box<dyn Future<Output = Result<Vec<PartitionInfo>, Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Translates Java’s List<PartitionInfo> partitionsFor(String topic, Duration).

Source

fn list_topics<'life0, 'async_trait>( &'life0 mut self, ) -> Pin<Box<dyn Future<Output = Result<HashMap<String, Vec<PartitionInfo>>, Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Translates Java’s Map<String, List<PartitionInfo>> listTopics().

Source

fn list_topics_with_timeout<'life0, 'async_trait>( &'life0 mut self, timeout: Duration, ) -> Pin<Box<dyn Future<Output = Result<HashMap<String, Vec<PartitionInfo>>, Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Translates Java’s Map<String, List<PartitionInfo>> listTopics(Duration).

Source

fn offsets_for_times<'life0, 'async_trait>( &'life0 mut self, timestamps_to_search: HashMap<TopicPartition, i64>, ) -> Pin<Box<dyn Future<Output = Result<HashMap<TopicPartition, OffsetAndTimestamp>, Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Translates Java’s Map<TopicPartition, OffsetAndTimestamp> offsetsForTimes( Map<TopicPartition, Long>).

Contract note (deviation from Java): Java returns a map whose value is nullable, so an unresolved partition (queried but no offset at/after the target time) is conveyed as key present, value null, and result.keySet() always contains every queried partition. The Rust OffsetAndTimestamp value is non-nullable, so an unresolved partition is omitted entirely (key absent) rather than present-with-null. Callers porting Java code that iterates result.keySet() expecting every queried key back must instead treat an absent key as “no offset”. See OffsetAndTimestamp.

Source

fn offsets_for_times_with_timeout<'life0, 'async_trait>( &'life0 mut self, timestamps_to_search: HashMap<TopicPartition, i64>, timeout: Duration, ) -> Pin<Box<dyn Future<Output = Result<HashMap<TopicPartition, OffsetAndTimestamp>, Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Translates Java’s Map<TopicPartition, OffsetAndTimestamp> offsetsForTimes( Map<TopicPartition, Long>, Duration).

See Self::offsets_for_times for the unresolved-partition contract note (unresolved partitions are omitted, not present-with-null, unlike Java).

Source

fn beginning_offsets<'life0, 'life1, 'async_trait>( &'life0 mut self, partitions: &'life1 [TopicPartition], ) -> Pin<Box<dyn Future<Output = Result<HashMap<TopicPartition, i64>, Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Translates Java’s Map<TopicPartition, Long> beginningOffsets(Collection<TopicPartition>).

Source

fn beginning_offsets_with_timeout<'life0, 'life1, 'async_trait>( &'life0 mut self, partitions: &'life1 [TopicPartition], timeout: Duration, ) -> Pin<Box<dyn Future<Output = Result<HashMap<TopicPartition, i64>, Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Translates Java’s Map<TopicPartition, Long> beginningOffsets(Collection<TopicPartition>, Duration).

Source

fn end_offsets<'life0, 'life1, 'async_trait>( &'life0 mut self, partitions: &'life1 [TopicPartition], ) -> Pin<Box<dyn Future<Output = Result<HashMap<TopicPartition, i64>, Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Translates Java’s Map<TopicPartition, Long> endOffsets(Collection<TopicPartition>).

Source

fn end_offsets_with_timeout<'life0, 'life1, 'async_trait>( &'life0 mut self, partitions: &'life1 [TopicPartition], timeout: Duration, ) -> Pin<Box<dyn Future<Output = Result<HashMap<TopicPartition, i64>, Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Translates Java’s Map<TopicPartition, Long> endOffsets(Collection<TopicPartition>, Duration).

Source

fn pause<'life0, 'life1, 'async_trait>( &'life0 mut self, partitions: &'life1 [TopicPartition], ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Translates Java’s void pause(Collection<TopicPartition>).

Async because Java’s pause calls applicationEventHandler.addAndGet(new PausePartitionsEvent(...)) which blocks (AsyncKafkaConsumer.java:1279).

Source

fn resume<'life0, 'life1, 'async_trait>( &'life0 mut self, partitions: &'life1 [TopicPartition], ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Translates Java’s void resume(Collection<TopicPartition>).

Async because Java’s resume calls applicationEventHandler.addAndGet(new ResumePartitionsEvent(...)) which blocks (AsyncKafkaConsumer.java:1292).

Source

fn enforce_rebalance<'life0, 'async_trait>( &'life0 mut self, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Translates Java’s void enforceRebalance() (Consumer.java:267).

Java’s javadoc says this method is classic-protocol-only; under the KIP-848 protocol it returns an unsupported-version error. Match Java behavior.

Source

fn enforce_rebalance_with_reason<'life0, 'life1, 'async_trait>( &'life0 mut self, reason: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Translates Java’s void enforceRebalance(String reason) (Consumer.java:272).

The parameter intersection across Java’s two overloads is empty, so under CLAUDE.md §2 the no-arg form keeps the plain name and this one carries the reason parameter-name suffix.

Source

fn close<'life0, 'async_trait>( &'life0 mut self, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Translates Java’s void close(). Closes the consumer with default timeout.

Source

fn close_with_options<'life0, 'async_trait>( &'life0 mut self, options: CloseOptions, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Translates Java’s void close(CloseOptions option).

Java’s @Deprecated void close(Duration timeout) is not translated (CLAUDE.md §3); pass CloseOptions::new_timeout(timeout) here instead, which is what the Java overload’s body does.

Source

fn wakeup(&self)

Translates Java’s void wakeup(). Sync — callable from any task, including signal handlers.

Source

fn handle(&self) -> ConsumerHandle

Returns a Clone + Send + Sync ConsumerHandle exposing Consumer::wakeup and the reentrant-safe consumer operations, callable from a task / thread other than the one that owns the consumer.

No Java method counterpart — it recovers a Java capability. Java’s Consumer reference is itself shareable across threads, so (1) consumer.wakeup() can be called from another thread while the owning thread blocks in poll() / position(), and (2) a ConsumerRebalanceListener can call back into the consumer (assign/seek/pause/position/…) from inside a callback by capturing the consumer variable. Rust borrows the consumer as &mut self for the duration of a blocking call and Box<dyn Consumer> is not Clone, so neither is expressible with a bare reference; capture a ConsumerHandle instead. See ConsumerHandle.

Implementors§

Source§

impl<K, V> Consumer<K, V> for MockConsumer<K, V>
where K: Send + 'static, V: Send + 'static,