Skip to main content

Producer

Trait Producer 

Source
pub trait Producer<K, V>: Send + Sync {
    // Required methods
    fn init_transactions(
        &self,
    ) -> impl Future<Output = Result<(), Error>> + Send;
    fn begin_transaction(&self) -> Result<(), Error>;
    fn send_offsets_to_transaction(
        &self,
        offsets: HashMap<TopicPartition, OffsetAndMetadata>,
        group_metadata: &dyn ConsumerGroupMetadata,
    ) -> impl Future<Output = Result<(), Error>> + Send;
    fn commit_transaction(
        &self,
    ) -> impl Future<Output = Result<(), Error>> + Send;
    fn abort_transaction(
        &self,
    ) -> impl Future<Output = Result<(), Error>> + Send;
    fn send(
        &self,
        record: ProducerRecord<K, V>,
    ) -> impl Future<Output = Result<KafkaFuture<RecordMetadata>, Error>> + Send;
    fn send_with_callback(
        &self,
        record: ProducerRecord<K, V>,
        callback: Option<Callback>,
    ) -> impl Future<Output = Result<KafkaFuture<RecordMetadata>, Error>> + Send;
    fn flush(&self) -> impl Future<Output = Result<(), Error>> + Send;
    fn partitions_for(
        &self,
        topic: &str,
    ) -> impl Future<Output = Result<Vec<PartitionInfo>, Error>> + Send;
    fn metrics(&self) -> HashMap<MetricName, Arc<KafkaMetric>>;
    fn close(&self) -> impl Future<Output = Result<(), Error>> + Send;
    fn close_with_timeout(
        &self,
        timeout: Duration,
    ) -> impl Future<Output = Result<(), Error>> + Send;
}
Expand description

The interface for the KafkaProducer.

Translated from org.apache.kafka.clients.producer.Producer.

Every asynchronous method returns impl Future + Send rather than being declared async fn, so that generic code can move the returned futures across tasks and box them as dyn Future + Send. That guarantee is what lets DynProducer be blanket-implemented for every Producer. Implementations may still be written with async fn.

Required Methods§

Source

fn init_transactions(&self) -> impl Future<Output = Result<(), Error>> + Send

Needs to be called before any other method when the transactional.id is set in the configuration.

See KafkaProducer::init_transactions.

§Errors

Returns Err if:

  • No transactional.id has been configured (LocalIllegalState)
  • The broker does not support transactions (UnsupportedVersion)
  • The configured transactional.id is not authorized, or the idempotent producer id is unavailable
  • The producer has encountered a previous fatal error
  • Initialization does not complete within max.block.ms (Timeout)
Source

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

Should be called before the start of each new transaction.

See KafkaProducer::begin_transaction.

Stays synchronous because Java’s beginTransaction (KafkaProducer.java:674-681) is a pure state transition and never blocks.

§Errors

Returns Err if no transactional.id has been configured, if init_transactions has not yet been invoked, if another producer with the same transactional.id has fenced this one, or if the producer has encountered a previous fatal error.

Source

fn send_offsets_to_transaction( &self, offsets: HashMap<TopicPartition, OffsetAndMetadata>, group_metadata: &dyn ConsumerGroupMetadata, ) -> impl Future<Output = Result<(), Error>> + Send

Sends a list of specified offsets to the consumer group coordinator, and also marks those offsets as part of the current transaction.

See KafkaProducer::send_offsets_to_transaction.

§Errors

Returns Err if no transactional.id has been configured or no transaction has been started, if group_metadata is invalid, if the commit failed and cannot be retried, or if the offsets are not sent within max.block.ms (Timeout).

Source

fn commit_transaction(&self) -> impl Future<Output = Result<(), Error>> + Send

Commits the ongoing transaction.

See KafkaProducer::commit_transaction.

§Errors

Returns Err if no transactional.id has been configured or no transaction has been started, if the producer has encountered a previous fatal or abortable error, or if the commit does not complete within max.block.ms (Timeout).

Source

fn abort_transaction(&self) -> impl Future<Output = Result<(), Error>> + Send

Aborts the ongoing transaction.

See KafkaProducer::abort_transaction.

§Errors

Returns Err if no transactional.id has been configured or no transaction has been started, if the producer has encountered a previous fatal error, or if the abort does not complete within max.block.ms (Timeout).

Source

fn send( &self, record: ProducerRecord<K, V>, ) -> impl Future<Output = Result<KafkaFuture<RecordMetadata>, Error>> + Send

Asynchronously send a record to a topic. Equivalent to send_with_callback(record, None).

See send_with_callback for details.

Source

fn send_with_callback( &self, record: ProducerRecord<K, V>, callback: Option<Callback>, ) -> impl Future<Output = Result<KafkaFuture<RecordMetadata>, Error>> + Send

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

The send is asynchronous and this method will return immediately once the record has been stored in the buffer of records waiting to be sent. It may block waiting for metadata or buffer space.

§Arguments
  • record - The record to send
  • callback - A user-supplied callback to execute when the record has been acknowledged by the server (None indicates no callback)
§Errors

Returns Err if:

Source

fn flush(&self) -> impl Future<Output = Result<(), Error>> + Send

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

§Errors

Returns Err if an error occurs during flushing.

Source

fn partitions_for( &self, topic: &str, ) -> impl Future<Output = Result<Vec<PartitionInfo>, Error>> + Send

Get the partition metadata for the given topic.

This can be used for custom partitioning.

§Errors

Returns Err if:

  • The topic cannot be found within max.block.ms (Timeout)
  • The producer has been closed
Source

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

Get the full set of producer metrics maintained by this producer.

Translated from Producer.metrics(). The returned map is keyed by MetricName; the value type is Arc<KafkaMetric> — KafkaMetric is the concrete registry entry (Java’s Metric interface). This method does not block in Java, so it stays a synchronous fn.

The returned HashMap is a snapshot clone of Arc<KafkaMetric> handles; mutating it does not affect the registry (Java’s Collections.unmodifiableMap analog).

Source

fn close(&self) -> impl Future<Output = Result<(), Error>> + Send

Close this producer. This method awaits until all previously sent requests complete.

§Errors

Returns Err if an error occurs during closing.

Source

fn close_with_timeout( &self, timeout: Duration, ) -> impl Future<Output = Result<(), Error>> + Send

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

If the producer is unable to complete all requests before the timeout expires, this method will fail any unsent and unacknowledged records immediately.

§Errors

Returns Err if an error occurs during closing.

Dyn Compatibility§

This trait is not dyn compatible.

In older versions of Rust, dyn compatibility was called "object safety", so this trait is not object safe.

Implementors§

Source§

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

Source§

impl<K, V> Producer<K, V> for dyn DynProducer<K, V> + '_

Source§

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