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§
Sourcefn 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,
Sourcefn begin_transaction(&self) -> Result<(), Error>
fn begin_transaction(&self) -> Result<(), Error>
Sourcefn 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 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,
Sourcefn 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,
Sourcefn 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,
Sourcefn 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,
Sourcefn 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 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,
Sourcefn 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,
Sourcefn 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,
Sourcefn metrics(&self) -> HashMap<MetricName, Arc<KafkaMetric>>
fn metrics(&self) -> HashMap<MetricName, Arc<KafkaMetric>>
See Producer::metrics.
Sourcefn close<'a>(
&'a self,
) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'a>>where
K: 'a,
V: 'a,
fn close<'a>(
&'a self,
) -> 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> + '_
impl<K, V> Producer<K, V> for dyn DynProducer<K, V> + '_
Source§fn init_transactions(&self) -> impl Future<Output = Result<(), Error>> + Send
fn init_transactions(&self) -> impl Future<Output = Result<(), Error>> + Send
transactional.id is
set in the configuration. Read moreSource§fn begin_transaction(&self) -> Result<(), Error>
fn begin_transaction(&self) -> Result<(), 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>
Source§fn commit_transaction(&self) -> impl Future<Output = Result<(), Error>> + Send
fn commit_transaction(&self) -> impl Future<Output = Result<(), Error>> + Send
Source§fn abort_transaction(&self) -> impl Future<Output = Result<(), Error>> + Send
fn abort_transaction(&self) -> impl Future<Output = Result<(), Error>> + Send
Source§fn 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
send_with_callback(record, None). Read more