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§
Sourcefn init_transactions(&self) -> impl Future<Output = Result<(), Error>> + Send
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.idhas been configured (LocalIllegalState) - The broker does not support transactions
(
UnsupportedVersion) - The configured
transactional.idis 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)
Sourcefn begin_transaction(&self) -> Result<(), Error>
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.
Sourcefn send_offsets_to_transaction(
&self,
offsets: HashMap<TopicPartition, OffsetAndMetadata>,
group_metadata: &dyn ConsumerGroupMetadata,
) -> impl Future<Output = Result<(), Error>> + Send
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).
Sourcefn commit_transaction(&self) -> impl Future<Output = Result<(), Error>> + Send
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).
Sourcefn abort_transaction(&self) -> impl Future<Output = Result<(), Error>> + Send
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).
Sourcefn send(
&self,
record: ProducerRecord<K, V>,
) -> impl Future<Output = Result<KafkaFuture<RecordMetadata>, Error>> + Send
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.
Sourcefn send_with_callback(
&self,
record: ProducerRecord<K, V>,
callback: Option<Callback>,
) -> 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
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 sendcallback- A user-supplied callback to execute when the record has been acknowledged by the server (Noneindicates no callback)
§Errors
Returns Err if:
- The producer has already been closed (
LocalIllegalState) - The key or value cannot be serialized (
Serialization) - A Kafka-related error occurs
Sourcefn flush(&self) -> impl Future<Output = Result<(), Error>> + Send
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.
Sourcefn partitions_for(
&self,
topic: &str,
) -> impl Future<Output = Result<Vec<PartitionInfo>, Error>> + Send
fn partitions_for( &self, topic: &str, ) -> impl Future<Output = Result<Vec<PartitionInfo>, Error>> + Send
Sourcefn metrics(&self) -> HashMap<MetricName, Arc<KafkaMetric>>
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).
Sourcefn close(&self) -> impl Future<Output = Result<(), Error>> + Send
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.
Sourcefn close_with_timeout(
&self,
timeout: Duration,
) -> impl Future<Output = Result<(), Error>> + Send
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.