Skip to main content

DynProducer

Trait DynProducer 

Source
pub trait DynProducer<K, V>:
    Sealed<K, V>
    + Send
    + Sync {
    // Required methods
    fn init_transactions<'a>(
        &'a self,
    ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'a>>
       where K: 'a,
             V: 'a;
    fn begin_transaction(&self) -> Result<(), Error>;
    fn send_offsets_to_transaction<'a>(
        &'a self,
        offsets: HashMap<TopicPartition, OffsetAndMetadata>,
        group_metadata: &'a dyn ConsumerGroupMetadata,
    ) -> 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;
    fn abort_transaction<'a>(
        &'a self,
    ) -> Pin<Box<dyn Future<Output = Result<(), 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;
    fn send_with_callback<'a>(
        &'a self,
        record: ProducerRecord<K, V>,
        callback: Option<Callback>,
    ) -> Pin<Box<dyn Future<Output = Result<KafkaFuture<RecordMetadata>, 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;
    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 metrics(&self) -> HashMap<MetricName, Arc<KafkaMetric>>;
    fn close<'a>(
        &'a self,
    ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'a>>
       where K: 'a,
             V: 'a;
    fn close_with_timeout<'a>(
        &'a self,
        timeout: Duration,
    ) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'a>>
       where K: 'a,
             V: 'a;
}
Expand description

A dyn-compatible Producer, for holding different producer implementations behind one type, e.g. Vec<Box<dyn DynProducer<K, V>>>.

Every Producer implements DynProducer automatically, and the trait is sealed, so it cannot be implemented directly: implement Producer instead. Each method forwards to the Producer method of the same name; asynchronous methods return the future boxed, which costs one allocation per call. dyn DynProducer<K, V> implements Producer in turn, so a boxed producer can be passed as &*boxed to generic code taking P: Producer<K, V> + ?Sized.

With both traits in scope a method call on a dyn DynProducer is ambiguous, since both define it; call it as DynProducer::send(&*producer, record), or import only one of the two traits.

use confluent_kafka::producer::{DynProducer, MockProducer};

let producers: Vec<Box<dyn DynProducer<String, String>>> = vec![
    Box::new(MockProducer::<String, String>::with_auto_complete(true)),
    Box::new(MockProducer::<String, String>::with_auto_complete(false)),
];
assert_eq!(2, producers.len());

Required Methods§

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, ) -> 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<Callback>, ) -> 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,

Trait Implementations§

Source§

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

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. Read more
Source§

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

Should be called before the start of each new transaction. Read more
Source§

async fn send_offsets_to_transaction( &self, offsets: HashMap<TopicPartition, OffsetAndMetadata>, group_metadata: &dyn ConsumerGroupMetadata, ) -> Result<(), Error>

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

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

Commits the ongoing transaction. Read more
Source§

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

Aborts the ongoing transaction. Read more
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). Read more
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. Read more
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. 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§

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

Get the full set of producer metrics maintained by this producer. Read more
Source§

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

Close this producer. This method awaits until all previously sent requests complete. Read more
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. Read more

Implementors§

Source§

impl<K, V, P: Producer<K, V>> DynProducer<K, V> for P