Skip to main content

MockProducer

Struct MockProducer 

Source
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>

Source

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 - If true, automatically complete all requests successfully. Otherwise the user must call complete_next() or error_next() after send() to complete the call and resolve the returned future (the crate-internal FutureRecordMetadata).

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.

Source

pub fn with_options(options: MockProducerOptions<K, V>) -> Self

Create a mock producer with a custom partitioner and key/value serializers.

§Arguments

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.

Source

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).

Source

pub fn history(&self) -> Vec<ProducerRecord<K, V>>
where K: Clone, V: Clone,

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().

Source

pub fn uncommitted_records(&self) -> Vec<ProducerRecord<K, V>>
where K: Clone, V: Clone,

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).

Source

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).

Source

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.

Source

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.

Source

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).

Source

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().

Source

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).

Source

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).

Source

pub fn transaction_initialized(&self) -> bool

Returns true if init_transactions() has completed successfully.

Corresponds to Java’s MockProducer.transactionInitialized() (MockProducer.java:436).

Source

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).

Source

pub fn transaction_committed(&self) -> bool

Returns true if the most recent transaction was committed.

Corresponds to Java’s MockProducer.transactionCommitted() (MockProducer.java:444).

Source

pub fn transaction_aborted(&self) -> bool

Returns true if the most recent transaction was aborted.

Corresponds to Java’s MockProducer.transactionAborted() (MockProducer.java:448).

Source

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).

Source

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).

Source

pub fn closed(&self) -> bool

Returns true if the producer is closed.

Corresponds to Java’s MockProducer.closed().

Source

pub fn flushed(&self) -> bool

Returns true if there are no pending completions.

Corresponds to Java’s MockProducer.flushed().

Source

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.

Source

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.

Source

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.

Source

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.

Source

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).

Source

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.

Source

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).

Source

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).

Source

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).

Source

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>

Source§

fn default() -> Self

Create a new mock producer with an empty cluster and auto_complete=false.

Corresponds to Java’s no-arg MockProducer() constructor.

Source§

impl<K: Send + Sync, V: Send + Sync> Producer<K, V> for MockProducer<K, V>

Source§

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>

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>

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>

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>

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>>

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>

Asynchronously send a record to a topic. Equivalent to send_with_callback(record, None). Read more
Source§

async fn send_with_callback( &self, record: ProducerRecord<K, V>, callback: Option<Callback>, ) -> Result<KafkaFuture<RecordMetadata>, Error>

Asynchronously send a record to a topic and invoke the provided callback when the send has been acknowledged. Read more
Source§

async fn flush(&self) -> Result<(), Error>

Invoking this method makes all buffered records immediately available to send and awaits the completion of the requests associated with these records. Read more
Source§

async fn partitions_for(&self, topic: &str) -> Result<Vec<PartitionInfo>, Error>

Get the partition metadata for the given topic. Read more
Source§

async fn close(&self) -> Result<(), Error>

Close this producer. This method awaits until all previously sent requests complete. Read more
Source§

async fn close_with_timeout(&self, _timeout: Duration) -> Result<(), Error>

Close this producer, waiting up to the given timeout for pending requests to complete. Read more

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>
where K: Send, V: Send,

§

impl<K, V> Sync for MockProducer<K, V>
where K: Send, V: Send,

§

impl<K, V> Unpin for MockProducer<K, V>
where K: Unpin, V: Unpin,

§

impl<K, V> UnsafeUnpin for MockProducer<K, V>

§

impl<K, V> UnwindSafe for MockProducer<K, V>

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<K, V, P> DynProducer<K, V> for P
where 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,

Source§

fn begin_transaction(&self) -> Result<(), Error>

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,

Source§

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,

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,

Source§

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,

Source§

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>>

Source§

fn close<'a>( &'a self, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'a>>
where K: 'a, V: 'a,

Source§

fn close_with_timeout<'a>( &'a self, timeout: Duration, ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'a>>
where K: 'a, V: 'a,

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

§

fn vzip(self) -> V