pub struct KafkaProducer<K, V> { /* private fields */ }Expand description
A Kafka client that publishes records to the Kafka cluster.
The producer is thread-safe and sharing a single producer instance across tasks will generally be faster than having multiple instances.
The producer consists of a pool of buffer space that holds records that haven’t yet been transmitted to the server, as well as a background I/O task that is responsible for turning these records into requests and transmitting them to the cluster. Failure to close the producer after use will leak these resources.
The send method is asynchronous. When called, it
adds the record to a buffer of pending record sends and immediately returns.
This allows the producer to batch together individual records for efficiency.
Translated from org.apache.kafka.clients.producer.KafkaProducer.
Implementations§
Source§impl<K, V> KafkaProducer<K, V>
impl<K, V> KafkaProducer<K, V>
Sourcepub const NETWORK_THREAD_PREFIX: &str = "kafka-producer-network-thread"
pub const NETWORK_THREAD_PREFIX: &str = "kafka-producer-network-thread"
Network thread name prefix.
Sourcepub const PRODUCER_METRIC_GROUP_NAME: &str = "producer-metrics"
pub const PRODUCER_METRIC_GROUP_NAME: &str = "producer-metrics"
Producer metric group name.
Sourcepub fn new(
config: ProducerConfig,
key_serializer: Box<dyn Serializer<K> + Send + Sync>,
value_serializer: Box<dyn Serializer<V> + Send + Sync>,
) -> Result<Self, Error>where
K: 'static,
V: 'static,
pub fn new(
config: ProducerConfig,
key_serializer: Box<dyn Serializer<K> + Send + Sync>,
value_serializer: Box<dyn Serializer<V> + Send + Sync>,
) -> Result<Self, Error>where
K: 'static,
V: 'static,
Creates a KafkaProducer from configuration, serializers, and an optional
compression override.
This is the primary public factory method, mirroring Java’s
new KafkaProducer(Properties, Serializer, Serializer) constructor
(KafkaProducer.java:339) and its Map twin (:312).
§Java’s no-serializer constructors are deliberately not translated
Java has two further public constructors, KafkaProducer(Map) (:295)
and KafkaProducer(Properties) (:324), which both delegate to
this(configs, null, null). They exist only because the private
constructor can fill a null serializer in reflectively:
if (keySerializer == null) {
keySerializer = config.getConfiguredInstance(KEY_SERIALIZER_CLASS_CONFIG, Serializer.class);(KafkaProducer.java:391-392)
key.serializer is a Type.CLASS entry (ProducerConfig.java:479-482),
so honouring those constructors means loading and instantiating a class
named by a string at run time. Rust has no reflection, and — unlike
partitioner.type, where the built-in names can be mapped to concrete
types by
ProducerConfig::resolve_partitioner —
a serializer is typed in the producer’s own K / V, so no such mapping
can be written for an arbitrary K. The serializers are therefore
always supplied here as instances, and
ProducerConfig deliberately carries
no key.serializer / value.serializer state at all. A key.serializer
entry in the property map is ignored as an unknown key.
Consequence for CLAUDE.md §2: the Java constructor group’s parameter
intersection is {configs}, and Java does have an overload with exactly
that (:295) — but it is untranslatable, so there is no Rust
constructor that could hold the plain name on its behalf. Rather than
leave new permanently unused and rename the only general-purpose
constructor after a sibling that can never exist,
the plain name stays here. The consumer side takes the
identical decision — see
AsyncKafkaConsumer::new.
It internally wires up all infrastructure components:
- Parses and resolves bootstrap server addresses from the config
- Creates
ProducerMetadataand bootstraps it with the resolved addresses - Creates a
PlaintextChannelBuilder,Selector, andNetworkClient - Creates a
BufferPoolandRecordAccumulator - Spawns the background sender task via the crate-internal
with_client_options
§Arguments
config- The producer configurationkey_serializer- The key serializervalue_serializer- The value serializer
§Errors
Returns Error::LocalIllegalArgument if no valid bootstrap server addresses
can be resolved from config.bootstrap_servers.
§Examples
use std::collections::HashMap;
use confluent_kafka::producer::KafkaProducer;
use confluent_kafka::producer::ProducerConfig;
use confluent_kafka::common::serialization::StringSerializer;
let props = HashMap::from([
("bootstrap.servers".to_string(), "localhost:9092".to_string()),
("client.id".to_string(), "my-producer".to_string()),
]);
let config = ProducerConfig::new(&props)
.expect("Invalid config");
let producer = KafkaProducer::<String, String>::new(
config,
Box::new(StringSerializer::default()),
Box::new(StringSerializer::default()),
).expect("Failed to create producer");Sourcepub async fn init_transactions(&self) -> Result<(), Error>
pub async fn init_transactions(&self) -> Result<(), Error>
Needs to be called before any other method when the transactional.id is
set in the configuration.
Translated from KafkaProducer.initTransactions()
(KafkaProducer.java:648-659). This method does the following:
- Ensures any transactions initiated by previous instances of the producer
with the same
transactional.idare completed. If the previous instance had failed with a transaction in progress, it will be aborted. If the last transaction had begun completion, but not yet finished, this method awaits its completion. - Gets the internal producer id and epoch, used in all future transactional messages issued by the producer.
Java blocks on result.await(maxBlockTimeMs, ..), so this is async
(CLAUDE.md §11.1) and returns Error::Timeout when the transactional
state cannot be initialized before max.block.ms expires. It is safe to
retry in that case, but once the transactional state has been successfully
initialized this method should no longer be used.
Java’s InterruptException path has no Rust analogue — a task is not
interrupted, it is dropped.
§Errors
Error::LocalIllegalStateif notransactional.idhas been configuredError::UnsupportedVersionas a fatal error indicating the broker does not support transactions- An authorization error indicating that the configured
transactional.idis not authorized, or the idempotent producer id is unavailable; the user may retry after fixing the permission - Any previous fatal error the producer has encountered
Error::Timeoutif initializing the transaction takes longer thanmax.block.ms
Sourcepub fn begin_transaction(&self) -> Result<(), Error>
pub fn begin_transaction(&self) -> Result<(), Error>
Should be called before the start of each new transaction. Note that prior
to the first invocation of this method, Self::init_transactions must be
invoked exactly one time.
Translated from KafkaProducer.beginTransaction()
(KafkaProducer.java:674-681). Stays synchronous: Java’s body is a pure
state transition with no wait, so CLAUDE.md §11.1 does not apply.
§Errors
Error::LocalIllegalStateif notransactional.idhas been configured or ifSelf::init_transactionshas not yet been invoked- A producer-fenced error if another producer with the same
transactional.idis active - An invalid-producer-epoch error if the producer has attempted to produce with an old epoch to the partition leader
Error::UnsupportedVersionas a fatal error indicating the broker does not support transactions- Any previous fatal error the producer has encountered
Sourcepub async fn send_offsets_to_transaction(
&self,
offsets: HashMap<TopicPartition, OffsetAndMetadata>,
group_metadata: &dyn ConsumerGroupMetadata,
) -> Result<(), Error>
pub 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. These offsets will be considered committed only if the transaction is committed successfully.
Translated from
KafkaProducer.sendOffsetsToTransaction(Map, ConsumerGroupMetadata)
(KafkaProducer.java:733-746).
The committed offset should be the next message the application will
consume, i.e. next_record_to_be_processed.offset(). The leader epoch
should also be added as commit metadata.
This method should be used when consumed and produced messages need to be
batched together, typically in a consume-transform-produce pattern. Thus
group_metadata should be obtained from the consumer’s group_metadata()
to leverage consumer group metadata, which provides stronger fencing than
a metadata carrying only the group id.
Java blocks until the request has been received and acknowledged by the
consumer group coordinator; the offsets are not considered committed until
the transaction itself is successfully committed via
Self::commit_transaction.
Note that the consumer should have enable.auto.commit=false and should
also not commit offsets manually.
offsets is taken by value because the transaction manager moves it into
the AddOffsetsToTxn handler that carries it to the coordinator — the same
convention AsyncKafkaConsumer::commit_sync_with_offsets already uses for an
offsets map. group_metadata is borrowed: ConsumerGroupMetadata is a
trait, and the handler keeps a copy of its four fields.
§Errors
Error::LocalIllegalArgumentifgroup_metadatahas a generation id greater than zero but an unknown member idError::LocalIllegalStateif notransactional.idhas been configured or no transaction has been started- A producer-fenced error if another producer with the same
transactional.idis active Error::UnsupportedVersionas a fatal error indicating the broker does not support transactions, or does not support the latest version of the transactional API with all consumer group metadata- An authorization error indicating that the configured
transactional.idor the consumer group id is not authorized - A commit-failed error if the commit cannot be retried (e.g. the consumer has been kicked out of the group); users should handle this by aborting the transaction
Error::Timeoutif sending the offsets takes longer thanmax.block.ms
Sourcepub async fn commit_transaction(&self) -> Result<(), Error>
pub async fn commit_transaction(&self) -> Result<(), Error>
Commits the ongoing transaction. This method will flush any unsent records before actually committing the transaction.
Translated from KafkaProducer.commitTransaction()
(KafkaProducer.java:779-786).
If any of the send calls which were part of the transaction
hit irrecoverable errors, this method returns the last received error
immediately and the transaction is not committed. So all send calls in a
transaction must succeed in order for this method to succeed.
If the transaction is committed successfully and this method returns
Ok(()), it is guaranteed that all callbacks for records in the
transaction will have been invoked and completed. Note that errors returned
by callbacks are ignored; the producer proceeds to commit the transaction in
any case.
A Error::Timeout does not mean the request did not reach the
broker — only that the acknowledgement did not arrive in time, so it is up
to the application to decide how to handle it. It is safe to retry, but it
is not possible to attempt a different operation (such as
Self::abort_transaction) since the commit may already be in the process
of completing. If not retrying, the only option is to close the producer.
§Errors
Error::LocalIllegalStateif notransactional.idhas been configured or no transaction has been started- A producer-fenced error if another producer with the same
transactional.idis active Error::UnsupportedVersionas a fatal error indicating the broker does not support transactions- An authorization error indicating that the configured
transactional.idis not authorized - An invalid-producer-epoch error if the producer has attempted to produce with an old epoch to the partition leader
- Any previous fatal or abortable error the producer has encountered
Error::Timeoutif committing takes longer thanmax.block.ms
Sourcepub async fn abort_transaction(&self) -> Result<(), Error>
pub async fn abort_transaction(&self) -> Result<(), Error>
Aborts the ongoing transaction. Any unflushed produce messages will be aborted when this call is made.
Translated from KafkaProducer.abortTransaction()
(KafkaProducer.java:813-821).
This call returns an error immediately if any prior send
call failed with a producer-fenced or an authorization error.
A Error::Timeout does not mean the request did not reach the
broker — see Self::commit_transaction for the full note; it is safe to
retry, but not to attempt a different operation.
§Errors
Error::LocalIllegalStateif notransactional.idhas been configured or no transaction has been started- A producer-fenced error if another producer with the same
transactional.idis active - An invalid-producer-epoch error if the producer has attempted to produce with an old epoch to the partition leader
Error::UnsupportedVersionas a fatal error indicating the broker does not support transactions- An authorization error indicating that the configured
transactional.idis not authorized - Any previous fatal error the producer has encountered
Error::Timeoutif aborting takes longer thanmax.block.ms
Source§impl KafkaProducer<Vec<u8>, Vec<u8>>
impl KafkaProducer<Vec<u8>, Vec<u8>>
Sourcepub async fn send(
&self,
record: ProducerRecord<&[u8], &[u8]>,
callback: Option<Callback>,
) -> Result<KafkaFuture<RecordMetadata>, Error>
pub async fn send( &self, record: ProducerRecord<&[u8], &[u8]>, callback: Option<Callback>, ) -> Result<KafkaFuture<RecordMetadata>, Error>
Send a record with borrowed byte-slice key/value, bypassing serialization.
This is the zero-copy path for callers that already have &[u8] data
(e.g. the C FFI layer). The slices are passed directly through to the
accumulator’s batch buffer without any intermediate allocation.
Trait Implementations§
Source§impl<K, V> Drop for KafkaProducer<K, V>
impl<K, V> Drop for KafkaProducer<K, V>
Source§impl<K, V> Producer<K, V> for KafkaProducer<K, V>
impl<K, V> Producer<K, V> for KafkaProducer<K, V>
Source§async fn init_transactions(&self) -> Result<(), Error>
async fn init_transactions(&self) -> Result<(), Error>
Needs to be called before any other method when the transactional.id is
set in the configuration.
Source§fn begin_transaction(&self) -> Result<(), Error>
fn begin_transaction(&self) -> Result<(), Error>
Should be called before the start of each new transaction.
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>
Sends a list of specified offsets to the consumer group coordinator, and also marks those offsets as part of the current transaction.
Source§async fn send(
&self,
record: ProducerRecord<K, V>,
) -> Result<KafkaFuture<RecordMetadata>, Error>
async fn send( &self, record: ProducerRecord<K, V>, ) -> Result<KafkaFuture<RecordMetadata>, Error>
Asynchronously send a record to a topic.
See send_with_callback for details.
Source§async fn send_with_callback(
&self,
record: ProducerRecord<K, V>,
callback: Option<Callback>,
) -> Result<KafkaFuture<RecordMetadata>, Error>
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.
Source§async fn flush(&self) -> Result<(), Error>
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.
Translated from KafkaProducer.flush().
Source§fn 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 KafkaProducer.metrics() — a snapshot of the registry.
Source§async fn partitions_for(&self, topic: &str) -> Result<Vec<PartitionInfo>, Error>
async fn partitions_for(&self, topic: &str) -> Result<Vec<PartitionInfo>, Error>
Get the partition metadata for the given topic.
Source§async fn close(&self) -> Result<(), Error>
async fn close(&self) -> Result<(), Error>
Close this producer. This method awaits until all previously sent requests complete.
Source§async fn close_with_timeout(&self, timeout: Duration) -> Result<(), Error>
async fn close_with_timeout(&self, timeout: Duration) -> Result<(), Error>
Close this producer, waiting up to the given timeout for pending requests to complete.
Translated from KafkaProducer.close(Duration timeout).
If timeout > 0: initiates a graceful close and awaits the sender task up
to the remaining time. If the sender task is still alive after the timeout,
it is force-closed and awaited indefinitely.
If timeout == 0: force-closes immediately without draining.
Note: Rust’s Duration is unsigned, so the negative-timeout check from
Java is omitted (impossible to construct a negative Duration).
Auto Trait Implementations§
impl<K, V> !Freeze for KafkaProducer<K, V>
impl<K, V> !RefUnwindSafe for KafkaProducer<K, V>
impl<K, V> Send for KafkaProducer<K, V>
impl<K, V> Sync for KafkaProducer<K, V>
impl<K, V> Unpin for KafkaProducer<K, V>
impl<K, V> UnsafeUnpin for KafkaProducer<K, V>
impl<K, V> !UnwindSafe for KafkaProducer<K, V>
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Source§impl<K, V, P> DynProducer<K, V> for Pwhere
P: Producer<K, V>,
impl<K, V, P> DynProducer<K, V> for Pwhere
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,
fn init_transactions<'a>(
&'a self,
) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'a>>where
K: 'a,
V: 'a,
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,
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,
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,
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,
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,
Producer::send. Read moreSource§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,
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,
fn flush<'a>(
&'a self,
) -> Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'a>>where
K: 'a,
V: 'a,
Producer::flush. Read moreSource§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 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>>
fn metrics(&self) -> HashMap<MetricName, Arc<KafkaMetric>>
Producer::metrics.