pub trait Consumer<K, 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, thecrate::producer::Producer-style traits andConsumergain these methods together.metrics()IS translated (Phase M7) — it is in Java’sConsumerinterface and snapshots the registry that the metrics managers populate. SeeConsumer::metrics.subscribe(Pattern, [ConsumerRebalanceListener])(Java regex): superseded by theSubscriptionPatternvariants which match server-side regex semantics.poll(long timeoutMs),close(Duration timeout)(both@Deprecatedin Java): not translated per CLAUDE.md §3.
Required Methods§
Sourcefn assignment(&self) -> HashSet<TopicPartition>
fn assignment(&self) -> HashSet<TopicPartition>
Translates Java’s Set<TopicPartition> assignment().
Sourcefn subscription(&self) -> HashSet<String>
fn subscription(&self) -> HashSet<String>
Translates Java’s Set<String> subscription().
Sourcefn paused(&self) -> HashSet<TopicPartition>
fn paused(&self) -> HashSet<TopicPartition>
Translates Java’s Set<TopicPartition> paused().
Sourcefn group_metadata(&self) -> Arc<dyn ConsumerGroupMetadata>
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.
Sourcefn current_lag(&self, topic_partition: &TopicPartition) -> Option<i64>
fn current_lag(&self, topic_partition: &TopicPartition) -> Option<i64>
Translates Java’s OptionalLong currentLag(TopicPartition).
Returns Option<i64> — the natural Rust analog.
Sourcefn metrics(&self) -> HashMap<MetricName, Arc<KafkaMetric>>
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.
Sourcefn 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<'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.
Sourcefn 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_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).
Sourcefn 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<'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.
Sourcefn 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 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).
Sourcefn 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 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.
Sourcefn 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 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().
Sourcefn 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 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).
Sourcefn 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<'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().
Sourcefn 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_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).
Sourcefn 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<'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).
Sourcefn 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_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).
Sourcefn 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<'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().
Sourcefn 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_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).
Sourcefn 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 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).
Sourcefn 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<'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).
Sourcefn 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_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).
Sourcefn 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_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>).
Sourcefn 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 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>).
Sourcefn 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<'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).
Sourcefn 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 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).
Sourcefn 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<'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>).
Sourcefn 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 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).
Sourcefn 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<'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).
Sourcefn 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 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).
Sourcefn 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<'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().
Sourcefn 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 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).
Sourcefn 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<'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.
Sourcefn 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 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).
Sourcefn 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<'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>).
Sourcefn 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 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).
Sourcefn 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<'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>).
Sourcefn 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 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).
Sourcefn 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 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).
Sourcefn 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 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).
Sourcefn 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<'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.
Sourcefn 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 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.
Sourcefn 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<'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.
Sourcefn 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 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.
Sourcefn wakeup(&self)
fn wakeup(&self)
Translates Java’s void wakeup(). Sync — callable from any task,
including signal handlers.
Sourcefn handle(&self) -> ConsumerHandle
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.