pub struct MockConsumer<K, V> { /* private fields */ }Expand description
A mock of the Consumer interface, intended
for testing code that uses Kafka.
Translated from
org.apache.kafka.clients.consumer.MockConsumer.
This struct is NOT thread-safe. However, you can use
MockConsumer::schedule_poll_task to write multi-task tests where one
task waits for Consumer::poll to be
called by another task and can safely perform operations during a
callback.
§Java mapping
Java’s MockConsumer(String offsetResetStrategy) is MockConsumer::new.
The deprecated MockConsumer(OffsetResetStrategy) and the deprecated
OffsetResetStrategy enum are not translated (CLAUDE.md §3); that
constructor only forwards offsetResetStrategy.toString(), so Rust callers
pass the auto.offset.reset string ("earliest", "latest", "none") to
MockConsumer::new instead.
Implementations§
Source§impl<K, V> MockConsumer<K, V>
impl<K, V> MockConsumer<K, V>
Sourcepub fn new(offset_reset_strategy: &str) -> Result<Self, Error>
pub fn new(offset_reset_strategy: &str) -> Result<Self, Error>
A mock consumer is instantiated by providing
ConsumerConfig::AUTO_OFFSET_RESET_CONFIG value as the input.
Returns an IllegalArgumentError when offset_reset_strategy is not
a valid auto.offset.reset value, as Java’s
AutoOffsetResetStrategy.fromString throws.
Sourcepub fn add_record(&mut self, record: ConsumerRecord<K, V>) -> Result<(), Error>
pub fn add_record(&mut self, record: ConsumerRecord<K, V>) -> Result<(), Error>
Add a record to the buffer that the next Consumer::poll
call will return. The record’s (topic, partition) must already be
assigned to the consumer.
Translates Java’s addRecord(ConsumerRecord<K, V>)
(MockConsumer.java:322). Returns
Error::LocalIllegalState if the partition is not assigned —
Java throws IllegalStateException in that case.
Sourcepub fn update_beginning_offsets(
&mut self,
offsets: HashMap<TopicPartition, i64>,
)
pub fn update_beginning_offsets( &mut self, offsets: HashMap<TopicPartition, i64>, )
Update the beginning offsets used for seekToBeginning resets.
Translates Java’s updateBeginningOffsets(Map<TopicPartition, Long>).
Sourcepub fn update_end_offsets(&mut self, offsets: HashMap<TopicPartition, i64>)
pub fn update_end_offsets(&mut self, offsets: HashMap<TopicPartition, i64>)
Update the end offsets used for seekToEnd resets and for
Consumer::end_offsets /
Consumer::current_lag.
Translates Java’s updateEndOffsets(Map<TopicPartition, Long>).
Sourcepub fn update_duration_offsets(&mut self, offsets: HashMap<TopicPartition, i64>)
pub fn update_duration_offsets(&mut self, offsets: HashMap<TopicPartition, i64>)
Update the duration-based reset offsets used when the reset strategy
is by_duration:<ISO-8601>.
Translates Java’s updateDurationOffsets(Map<TopicPartition, Long>).
Sourcepub fn update_partitions(
&mut self,
topic: &str,
partitions: Vec<PartitionInfo>,
) -> Result<(), Error>
pub fn update_partitions( &mut self, topic: &str, partitions: Vec<PartitionInfo>, ) -> Result<(), Error>
Configure the PartitionInfo list for a topic, used by
Consumer::partitions_for
and Consumer::list_topics.
Translates Java’s updatePartitions(String, List<PartitionInfo>).
Sourcepub fn set_poll_error(&mut self, error: Error)
pub fn set_poll_error(&mut self, error: Error)
Inject an exception to be returned by the next
Consumer::poll call. The
exception is taken (cleared) on use.
Translates Java’s setPollException(KafkaException).
Sourcepub fn set_offsets_error(&mut self, error: Error)
pub fn set_offsets_error(&mut self, error: Error)
Inject an exception to be returned by the next
Consumer::beginning_offsets /
Consumer::end_offsets
call. The exception is taken (cleared) on use.
Translates Java’s setOffsetsException(KafkaException).
Sourcepub fn set_max_poll_records(
&mut self,
max_poll_records: i64,
) -> Result<(), Error>
pub fn set_max_poll_records( &mut self, max_poll_records: i64, ) -> Result<(), Error>
Set the maximum number of records returned in a single
Consumer::poll call.
Translates Java’s setMaxPollRecords(long). Returns
Error::LocalIllegalArgument when max_poll_records < 1, matching
Java’s IllegalArgumentException.
Sourcepub async fn rebalance(
&mut self,
new_assignment: &[TopicPartition],
) -> Result<(), Error>
pub async fn rebalance( &mut self, new_assignment: &[TopicPartition], ) -> Result<(), Error>
Simulate a rebalance event: compute revoked / added partitions
against the current assignment, invoke any registered
ConsumerRebalanceListener
callbacks, and replace the assignment.
Translates Java’s rebalance(Collection<TopicPartition>). Async
because the Rust listener methods are async fn (per
consumer-threading.md §31 / Phase 2); Java’s listener methods
are sync.
Sourcepub fn schedule_poll_task(
&mut self,
task: Box<dyn FnOnce(&mut MockConsumer<K, V>) + Send>,
)
pub fn schedule_poll_task( &mut self, task: Box<dyn FnOnce(&mut MockConsumer<K, V>) + Send>, )
Schedule a task to run on the next
Consumer::poll call. One task
is consumed per poll invocation, in FIFO order.
Translates Java’s schedulePollTask(Runnable). Java’s Runnable
captures MockConsumer implicitly via closure; Rust closures cannot
capture the outer struct safely, so the task is invoked with
&mut MockConsumer<K, V> explicitly. Callers write
consumer.schedule_poll_task(Box::new(|c| { c.add_record(...).unwrap(); }));.
Sourcepub fn closed(&self) -> bool
pub fn closed(&self) -> bool
Return whether the consumer has been closed.
Translates Java’s closed().
Sourcepub fn should_rebalance(&self) -> bool
pub fn should_rebalance(&self) -> bool
Return whether an enforceRebalance request is pending.
Translates Java’s shouldRebalance().
Sourcepub fn reset_should_rebalance(&mut self)
pub fn reset_should_rebalance(&mut self)
Reset the rebalance-pending flag after handling it.
Translates Java’s resetShouldRebalance().
Sourcepub fn last_poll_timeout(&self) -> Option<Duration>
pub fn last_poll_timeout(&self) -> Option<Duration>
The timeout passed to the most recent
Consumer::poll call, or None
if poll has not been called yet.
Translates Java’s lastPollTimeout().
Trait Implementations§
Source§impl<K, V> Consumer<K, V> for MockConsumer<K, V>
impl<K, V> Consumer<K, V> for MockConsumer<K, V>
Source§fn metrics(&self) -> HashMap<MetricName, Arc<KafkaMetric>>
fn metrics(&self) -> HashMap<MetricName, Arc<KafkaMetric>>
Translates Java’s
synchronized Map<MetricName, ? extends Metric> metrics()
(MockConsumer.java:496-499).
Java: ensureNotClosed(); return Collections.emptyMap();. The mock
has no metrics registry, so the Rust port returns an empty map,
matching Java. (The Rust mock’s sync accessors do not panic on a
closed consumer — see the other accessors — so the ensureNotClosed
guard is not replicated here; the result is the same empty map either
way.)
Source§fn current_lag(&self, topic_partition: &TopicPartition) -> Option<i64>
fn current_lag(&self, topic_partition: &TopicPartition) -> Option<i64>
Translates Java’s
OptionalLong currentLag(TopicPartition)
(MockConsumer.java:681-688).
Behavior summary:
- Returns
Noneif the partition is not assigned to this consumer. (Java throwsIllegalArgumentExceptionfrom the innerposition(tp)call in that case; Rust’s&selfAPI can’t throw, soNoneis the closest analog and lets the caller distinguish “unassigned” from a real lag of0.) - Returns
Some(0)if assigned but no end offset is known (Java’s “caught up” model —endOffsetshas no entry for the partition). - Returns
Some(0)if assigned with end offset known but no position has been set yet. Java’sposition(tp)would callupdateFetchPositionto seed a position; the Rust&selfAPI cannot mutate, soSome(0)is the least-surprising fallback for the caught-up model. - Returns
Some(end - position)otherwise.
Divergence from Java: Java throws on unassigned partitions; Rust
returns None. Callers needing the strict Java behavior must
pre-check assignment().contains(tp) themselves.
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,
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 poll(Duration) (MockConsumer.java:249-320).
Source§fn assignment(&self) -> HashSet<TopicPartition>
fn assignment(&self) -> HashSet<TopicPartition>
Set<TopicPartition> assignment().Source§fn subscription(&self) -> HashSet<String>
fn subscription(&self) -> HashSet<String>
Set<String> subscription().Source§fn paused(&self) -> HashSet<TopicPartition>
fn paused(&self) -> HashSet<TopicPartition>
Set<TopicPartition> paused().Source§fn group_metadata(&self) -> Arc<dyn ConsumerGroupMetadata>
fn group_metadata(&self) -> Arc<dyn ConsumerGroupMetadata>
ConsumerGroupMetadata groupMetadata(). Read moreSource§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<'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,
void subscribe(Collection<String> topics). Read moreSource§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_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,
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,
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,
void subscribe(SubscriptionPattern pattern). Read moreSource§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 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,
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,
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,
void assign(Collection<TopicPartition>). Read moreSource§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 unsubscribe<'life0, 'async_trait>(
&'life0 mut self,
) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
void unsubscribe().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,
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,
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,
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,
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,
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,
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,
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,
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,
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,
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,
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,
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,
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,
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,
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,
void seek(TopicPartition partition, long offset). Read moreSource§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_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,
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,
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,
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,
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,
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,
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,
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,
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,
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,
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,
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,
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,
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,
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,
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,
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,
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,
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,
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,
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,
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,
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,
Map<TopicPartition, OffsetAndTimestamp> offsetsForTimes( Map<TopicPartition, Long>). Read moreSource§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 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,
Map<TopicPartition, OffsetAndTimestamp> offsetsForTimes( Map<TopicPartition, Long>, Duration). Read moreSource§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<'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,
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,
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,
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,
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,
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,
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,
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,
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,
void pause(Collection<TopicPartition>). Read moreSource§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 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,
void resume(Collection<TopicPartition>). Read moreSource§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<'life0, 'async_trait>(
&'life0 mut self,
) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
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,
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,
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,
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,
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,
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,
void close(CloseOptions option). Read moreSource§fn wakeup(&self)
fn wakeup(&self)
void wakeup(). Sync — callable from any task,
including signal handlers.Source§fn handle(&self) -> ConsumerHandle
fn handle(&self) -> ConsumerHandle
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. Read more