pub struct MockProducer<K, V> { /* private fields */ }Expand description
A mock of the producer interface for testing code that uses Kafka.
By default this mock will synchronously complete each send call successfully. However it can be configured to allow the user to control the completion of the call and supply an optional error for the producer to throw.
Corresponds to Java’s org.apache.kafka.clients.producer.MockProducer.
§Thread Safety
All methods use interior mutability via Mutex, matching Java’s
synchronized methods. The struct is Send + Sync so it can be shared
via Arc.
Implementations§
Source§impl<K, V> MockProducer<K, V>
impl<K, V> MockProducer<K, V>
Sourcepub fn with_cluster_auto_complete(cluster: Cluster, auto_complete: bool) -> Self
pub fn with_cluster_auto_complete(cluster: Cluster, auto_complete: bool) -> Self
Create a mock producer.
§Arguments
cluster- The cluster holding metadata for this producer.auto_complete- Iftrue, automatically complete all requests successfully. Otherwise the user must callcomplete_next()orerror_next()aftersend()to complete the call and resolve the returned future (the crate-internalFutureRecordMetadata).
Corresponds to Java’s MockProducer(Cluster, boolean, Partitioner, Serializer, Serializer) constructor invoked with a null partitioner and
null serializers (MockProducer.java:113). Delegates to
with_options.
Sourcepub fn with_options(options: MockProducerOptions<K, V>) -> Self
pub fn with_options(options: MockProducerOptions<K, V>) -> Self
Create a mock producer with a custom partitioner and key/value serializers.
§Arguments
options- Every Java parameter, built throughMockProducerOptionsBuilder.MockProducerOptionsis this method’s only parameter because the derived name would list five parameters, past CLAUDE.md §2’s cap of three.
Corresponds to Java’s MockProducer(Cluster, boolean, Partitioner, Serializer, Serializer) constructor (MockProducer.java:113), and — with
cluster left unset — to MockProducer(boolean, Partitioner, Serializer, Serializer) (:137), which passes Cluster.empty() itself.
Sourcepub fn with_auto_complete(auto_complete: bool) -> Self
pub fn with_auto_complete(auto_complete: bool) -> Self
Create a new mock producer with an empty cluster and the given
auto_complete setting.
Equivalent to MockProducer::with_cluster_auto_complete(Cluster::empty(), auto_complete).
Corresponds to Java’s MockProducer(boolean, Partitioner, Serializer, Serializer).
Sourcepub fn history(&self) -> Vec<ProducerRecord<K, V>>
pub fn history(&self) -> Vec<ProducerRecord<K, V>>
Get the list of sent records since the last call to clear().
Returns a clone of the internal sent list.
Corresponds to Java’s MockProducer.history().
Sourcepub fn uncommitted_records(&self) -> Vec<ProducerRecord<K, V>>
pub fn uncommitted_records(&self) -> Vec<ProducerRecord<K, V>>
Get the list of records sent inside the in-flight transaction and not yet committed.
Returns a clone of the internal uncommitted-sends list.
Corresponds to Java’s MockProducer.uncommittedRecords()
(MockProducer.java:471).
Sourcepub fn consumer_group_offsets_history(
&self,
) -> Vec<HashMap<String, HashMap<TopicPartition, OffsetAndMetadata>>>
pub fn consumer_group_offsets_history( &self, ) -> Vec<HashMap<String, HashMap<TopicPartition, OffsetAndMetadata>>>
Get the list of committed consumer group offsets since the last call to
clear() — one entry per committed transaction that
carried offsets.
Corresponds to Java’s MockProducer.consumerGroupOffsetsHistory()
(MockProducer.java:479).
Sourcepub fn committed_offset(
&self,
group: &str,
topic_partition: &TopicPartition,
) -> Option<OffsetAndMetadata>
pub fn committed_offset( &self, group: &str, topic_partition: &TopicPartition, ) -> Option<OffsetAndMetadata>
Look up the offset a committed transaction staged for group /
topic_partition, newest transaction first.
A targeted lookup over the same data as
consumer_group_offsets_history().
It exists because that accessor deep-clones the whole history — a
Vec<HashMap<String, HashMap<TopicPartition, OffsetAndMetadata>>> — which
is wasteful for a caller that wants one entry, and unreasonably so for the
C FFI probe that does it on every call. Here the scan happens under the
lock and only the matching entry is cloned. No Java counterpart; Java
callers index the returned map directly.
Sourcepub fn uncommitted_offsets(
&self,
) -> HashMap<String, HashMap<TopicPartition, OffsetAndMetadata>>
pub fn uncommitted_offsets( &self, ) -> HashMap<String, HashMap<TopicPartition, OffsetAndMetadata>>
Get the offsets staged by the in-flight transaction and not yet committed.
Corresponds to Java’s MockProducer.uncommittedOffsets()
(MockProducer.java:483). Java hands back the live map; behind the mutex
that is not expressible, so this returns a snapshot clone. No Java caller
mutates the returned map.
Sourcepub fn clear(&self)
pub fn clear(&self)
Clear the stored history of sent records and consumer group offsets.
Note: per-topic-partition offset counters are intentionally preserved
across clear() calls, matching Java’s MockProducer.clear() which
does not reset the offsets map. Nor does it reset the transaction
flags — only sentOffsets.
Corresponds to Java’s MockProducer.clear() (MockProducer.java:490).
Sourcepub fn complete_next(&self) -> bool
pub fn complete_next(&self) -> bool
Complete the earliest uncompleted call successfully.
Returns true if there was an uncompleted call to complete.
Corresponds to Java’s MockProducer.completeNext().
Sourcepub fn error_next(&self, error: Error) -> bool
pub fn error_next(&self, error: Error) -> bool
Complete the earliest uncompleted call with the given error.
Returns true if there was an uncompleted call to complete.
Corresponds to Java’s MockProducer.errorNext(RuntimeException).
Sourcepub fn fence_producer(&self) -> Result<(), Error>
pub fn fence_producer(&self) -> Result<(), Error>
Mark this producer as fenced by another producer with the same
transactional.id. Every subsequent transactional call and every
send() then fails with
Errors::ProducerFenced.
Corresponds to Java’s MockProducer.fenceProducer()
(MockProducer.java:429).
§Errors
Returns Err if the producer is closed, is already fenced, or was never
initialized for transactions (Error::local_illegal_state for the first
and last, Errors::ProducerFenced for the second).
Sourcepub fn transaction_initialized(&self) -> bool
pub fn transaction_initialized(&self) -> bool
Returns true if init_transactions() has
completed successfully.
Corresponds to Java’s MockProducer.transactionInitialized()
(MockProducer.java:436).
Sourcepub fn transaction_in_flight(&self) -> bool
pub fn transaction_in_flight(&self) -> bool
Returns true if a transaction has been begun and neither committed nor
aborted.
Corresponds to Java’s MockProducer.transactionInFlight()
(MockProducer.java:440).
Sourcepub fn transaction_committed(&self) -> bool
pub fn transaction_committed(&self) -> bool
Returns true if the most recent transaction was committed.
Corresponds to Java’s MockProducer.transactionCommitted()
(MockProducer.java:444).
Sourcepub fn transaction_aborted(&self) -> bool
pub fn transaction_aborted(&self) -> bool
Returns true if the most recent transaction was aborted.
Corresponds to Java’s MockProducer.transactionAborted()
(MockProducer.java:448).
Sourcepub fn sent_offsets(&self) -> bool
pub fn sent_offsets(&self) -> bool
Returns true if offsets were sent to the current or most recent
transaction. Reset only by
begin_transaction() and
clear() — not by a commit.
Corresponds to Java’s MockProducer.sentOffsets()
(MockProducer.java:456).
Sourcepub fn commit_count(&self) -> i64
pub fn commit_count(&self) -> i64
The number of transactions committed so far. Aborted transactions are not counted.
Corresponds to Java’s MockProducer.commitCount()
(MockProducer.java:460).
Sourcepub fn closed(&self) -> bool
pub fn closed(&self) -> bool
Returns true if the producer is closed.
Corresponds to Java’s MockProducer.closed().
Sourcepub fn flushed(&self) -> bool
pub fn flushed(&self) -> bool
Returns true if there are no pending completions.
Corresponds to Java’s MockProducer.flushed().
Sourcepub fn set_send_error(&self, error: Option<Error>)
pub fn set_send_error(&self, error: Option<Error>)
Set an error to be returned on every send() call
until cleared.
The error persists across calls until explicitly cleared with None,
matching Java’s MockProducer.sendException field semantics.
Pass None to clear a previously set error.
Sourcepub fn set_flush_error(&self, error: Option<Error>)
pub fn set_flush_error(&self, error: Option<Error>)
Set an error to be returned on every flush() call
until cleared.
The error persists across calls until explicitly cleared with None,
matching Java’s MockProducer.flushException field semantics.
Pass None to clear a previously set error.
Sourcepub fn set_partitions_for_error(&self, error: Option<Error>)
pub fn set_partitions_for_error(&self, error: Option<Error>)
Set an error to be returned on every
partitions_for() call until cleared.
The error persists across calls until explicitly cleared with None,
matching Java’s MockProducer.partitionsForException field semantics.
Pass None to clear a previously set error.
Sourcepub fn set_close_error(&self, error: Option<Error>)
pub fn set_close_error(&self, error: Option<Error>)
Set an error to be returned on every close() call
until cleared.
The error persists across calls until explicitly cleared with None,
matching Java’s MockProducer.closeException field semantics.
Pass None to clear a previously set error.
Sourcepub fn set_mock_metrics(&self, name: MetricName, metric: Arc<KafkaMetric>)
pub fn set_mock_metrics(&self, name: MetricName, metric: Arc<KafkaMetric>)
Seed a metric returned by metrics().
Corresponds to Java’s MockProducer.setMockMetrics(MetricName name, Metric metric).
Sourcepub fn set_init_transaction_error(&self, error: Option<Error>)
pub fn set_init_transaction_error(&self, error: Option<Error>)
Set an error to be returned on every
init_transactions() call until cleared.
Matches Java’s public MockProducer.initTransactionException field
(MockProducer.java:79), which likewise persists until set back to null.
Sourcepub fn set_begin_transaction_error(&self, error: Option<Error>)
pub fn set_begin_transaction_error(&self, error: Option<Error>)
Set an error to be returned on every
begin_transaction() call until cleared.
Matches Java’s public MockProducer.beginTransactionException field
(MockProducer.java:80).
Sourcepub fn set_send_offsets_to_transaction_error(&self, error: Option<Error>)
pub fn set_send_offsets_to_transaction_error(&self, error: Option<Error>)
Set an error to be returned on every
send_offsets_to_transaction()
call until cleared.
Matches Java’s public MockProducer.sendOffsetsToTransactionException
field (MockProducer.java:81).
Sourcepub fn set_commit_transaction_error(&self, error: Option<Error>)
pub fn set_commit_transaction_error(&self, error: Option<Error>)
Set an error to be returned on every
commit_transaction() call until cleared.
Matches Java’s public MockProducer.commitTransactionException field
(MockProducer.java:82).
Sourcepub fn set_abort_transaction_error(&self, error: Option<Error>)
pub fn set_abort_transaction_error(&self, error: Option<Error>)
Set an error to be returned on every
abort_transaction() call until cleared.
Matches Java’s public MockProducer.abortTransactionException field
(MockProducer.java:83).
Trait Implementations§
Source§impl<K, V> Default for MockProducer<K, V>
impl<K, V> Default for MockProducer<K, V>
Source§impl<K: Send + Sync, V: Send + Sync> Producer<K, V> for MockProducer<K, V>
impl<K: Send + Sync, V: Send + Sync> Producer<K, V> for MockProducer<K, V>
Source§async fn init_transactions(&self) -> Result<(), Error>
async fn init_transactions(&self) -> Result<(), Error>
Initialize this mock for transactions.
Corresponds to Java’s MockProducer.initTransactions()
(MockProducer.java:145). Stays async because the trait declares it so
(Java’s KafkaProducer.initTransactions blocks); the mock never awaits.
§Errors
Returns Err if the producer is closed, is fenced, has already been
initialized, or an error was installed with
set_init_transaction_error.
Source§fn begin_transaction(&self) -> Result<(), Error>
fn begin_transaction(&self) -> Result<(), Error>
Begin a transaction.
Corresponds to Java’s MockProducer.beginTransaction()
(MockProducer.java:162).
§Errors
Returns Err if the producer is closed, is fenced, was not initialized for
transactions, a transaction is already in flight, or an error was installed
with set_begin_transaction_error.
Source§async fn send_offsets_to_transaction(
&self,
offsets: HashMap<TopicPartition, OffsetAndMetadata>,
group_metadata: &dyn ConsumerGroupMetadata,
) -> Result<(), Error>
async fn send_offsets_to_transaction( &self, offsets: HashMap<TopicPartition, OffsetAndMetadata>, group_metadata: &dyn ConsumerGroupMetadata, ) -> Result<(), Error>
Stage consumer group offsets as part of the in-flight transaction.
Corresponds to Java’s MockProducer.sendOffsetsToTransaction(Map, ConsumerGroupMetadata) (MockProducer.java:182). Java’s
Objects.requireNonNull(groupMetadata) (:184) has no counterpart: the
parameter is a reference and not an Option, so a missing metadata is not
expressible.
An empty offsets map is ignored and leaves
sent_offsets() false (Java :194-196).
§Errors
Returns Err if the producer is closed, is fenced, was not initialized for
transactions, has no open transaction, or an error was installed with
set_send_offsets_to_transaction_error.
Source§async fn commit_transaction(&self) -> Result<(), Error>
async fn commit_transaction(&self) -> Result<(), Error>
Commit the in-flight transaction, publishing its records and offsets.
Corresponds to Java’s MockProducer.commitTransaction()
(MockProducer.java:204).
§Errors
Returns Err if the producer is closed, is fenced, was not initialized for
transactions, has no open transaction, an error was installed with
set_commit_transaction_error,
or the flush() this performs fails.
Source§async fn abort_transaction(&self) -> Result<(), Error>
async fn abort_transaction(&self) -> Result<(), Error>
Abort the in-flight transaction, discarding its records and offsets.
Corresponds to Java’s MockProducer.abortTransaction()
(MockProducer.java:230).
§Errors
Returns Err if the producer is closed, is fenced, was not initialized for
transactions, has no open transaction, an error was installed with
set_abort_transaction_error,
or the flush() this performs fails.
Source§fn metrics(&self) -> HashMap<MetricName, Arc<KafkaMetric>>
fn metrics(&self) -> HashMap<MetricName, Arc<KafkaMetric>>
Return the mock metrics. Corresponds to Java’s MockProducer.metrics()
returning the mockMetrics map.
Source§async fn send(
&self,
record: ProducerRecord<K, V>,
) -> Result<KafkaFuture<RecordMetadata>, Error>
async fn send( &self, record: ProducerRecord<K, V>, ) -> Result<KafkaFuture<RecordMetadata>, Error>
send_with_callback(record, None). Read moreSource§async fn send_with_callback(
&self,
record: ProducerRecord<K, V>,
callback: Option<Callback>,
) -> Result<KafkaFuture<RecordMetadata>, Error>
async fn send_with_callback( &self, record: ProducerRecord<K, V>, callback: Option<Callback>, ) -> Result<KafkaFuture<RecordMetadata>, Error>
Source§async fn flush(&self) -> Result<(), Error>
async fn flush(&self) -> Result<(), Error>
Source§async fn partitions_for(&self, topic: &str) -> Result<Vec<PartitionInfo>, Error>
async fn partitions_for(&self, topic: &str) -> Result<Vec<PartitionInfo>, Error>
Auto Trait Implementations§
impl<K, V> !Freeze for MockProducer<K, V>
impl<K, V> RefUnwindSafe for MockProducer<K, V>
impl<K, V> Send for MockProducer<K, V>
impl<K, V> Sync for MockProducer<K, V>
impl<K, V> Unpin for MockProducer<K, V>
impl<K, V> UnsafeUnpin for MockProducer<K, V>
impl<K, V> UnwindSafe for MockProducer<K, V>
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Source§impl<K, V, P> DynProducer<K, V> for Pwhere
P: Producer<K, V>,
impl<K, V, P> DynProducer<K, V> for Pwhere
P: Producer<K, V>,
Source§fn init_transactions<'a>(
&'a self,
) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'a>>where
K: 'a,
V: 'a,
fn init_transactions<'a>(
&'a self,
) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'a>>where
K: 'a,
V: 'a,
Source§fn send_offsets_to_transaction<'a>(
&'a self,
offsets: HashMap<TopicPartition, OffsetAndMetadata>,
group_metadata: &'a (dyn ConsumerGroupMetadata + 'static),
) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'a>>where
K: 'a,
V: 'a,
fn send_offsets_to_transaction<'a>(
&'a self,
offsets: HashMap<TopicPartition, OffsetAndMetadata>,
group_metadata: &'a (dyn ConsumerGroupMetadata + 'static),
) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'a>>where
K: 'a,
V: 'a,
Source§fn commit_transaction<'a>(
&'a self,
) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'a>>where
K: 'a,
V: 'a,
fn commit_transaction<'a>(
&'a self,
) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'a>>where
K: 'a,
V: 'a,
Source§fn abort_transaction<'a>(
&'a self,
) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'a>>where
K: 'a,
V: 'a,
fn abort_transaction<'a>(
&'a self,
) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'a>>where
K: 'a,
V: 'a,
Source§fn send<'a>(
&'a self,
record: ProducerRecord<K, V>,
) -> Pin<Box<dyn Future<Output = Result<KafkaFuture<RecordMetadata>, Error>> + Send + 'a>>where
K: 'a,
V: 'a,
fn send<'a>(
&'a self,
record: ProducerRecord<K, V>,
) -> Pin<Box<dyn Future<Output = Result<KafkaFuture<RecordMetadata>, Error>> + Send + 'a>>where
K: 'a,
V: 'a,
Producer::send. Read moreSource§fn send_with_callback<'a>(
&'a self,
record: ProducerRecord<K, V>,
callback: Option<Box<dyn FnOnce(Option<&RecordMetadata>, Option<&Error>) + Sync + Send>>,
) -> Pin<Box<dyn Future<Output = Result<KafkaFuture<RecordMetadata>, Error>> + Send + 'a>>where
K: 'a,
V: 'a,
fn send_with_callback<'a>(
&'a self,
record: ProducerRecord<K, V>,
callback: Option<Box<dyn FnOnce(Option<&RecordMetadata>, Option<&Error>) + Sync + Send>>,
) -> Pin<Box<dyn Future<Output = Result<KafkaFuture<RecordMetadata>, Error>> + Send + 'a>>where
K: 'a,
V: 'a,
Source§fn flush<'a>(
&'a self,
) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'a>>where
K: 'a,
V: 'a,
fn flush<'a>(
&'a self,
) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'a>>where
K: 'a,
V: 'a,
Producer::flush. Read moreSource§fn partitions_for<'a>(
&'a self,
topic: &'a str,
) -> Pin<Box<dyn Future<Output = Result<Vec<PartitionInfo>, Error>> + Send + 'a>>where
K: 'a,
V: 'a,
fn partitions_for<'a>(
&'a self,
topic: &'a str,
) -> Pin<Box<dyn Future<Output = Result<Vec<PartitionInfo>, Error>> + Send + 'a>>where
K: 'a,
V: 'a,
Source§fn metrics(&self) -> HashMap<MetricName, Arc<KafkaMetric>>
fn metrics(&self) -> HashMap<MetricName, Arc<KafkaMetric>>
Producer::metrics.