Skip to main content

KafkaProducer

Struct KafkaProducer 

Source
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>

Source

pub const NETWORK_THREAD_PREFIX: &str = "kafka-producer-network-thread"

Network thread name prefix.

Source

pub const PRODUCER_METRIC_GROUP_NAME: &str = "producer-metrics"

Producer metric group name.

Source

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:

  1. Parses and resolves bootstrap server addresses from the config
  2. Creates ProducerMetadata and bootstraps it with the resolved addresses
  3. Creates a PlaintextChannelBuilder, Selector, and NetworkClient
  4. Creates a BufferPool and RecordAccumulator
  5. Spawns the background sender task via the crate-internal with_client_options
§Arguments
  • config - The producer configuration
  • key_serializer - The key serializer
  • value_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");
Source

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:

  1. Ensures any transactions initiated by previous instances of the producer with the same transactional.id are 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.
  2. 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::LocalIllegalState if no transactional.id has been configured
  • Error::UnsupportedVersion as a fatal error indicating the broker does not support transactions
  • An authorization error indicating that the configured transactional.id is 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::Timeout if initializing the transaction takes longer than max.block.ms
Source

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::LocalIllegalState if no transactional.id has been configured or if Self::init_transactions has not yet been invoked
  • A producer-fenced error if another producer with the same transactional.id is active
  • An invalid-producer-epoch error if the producer has attempted to produce with an old epoch to the partition leader
  • Error::UnsupportedVersion as a fatal error indicating the broker does not support transactions
  • Any previous fatal error the producer has encountered
Source

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::LocalIllegalArgument if group_metadata has a generation id greater than zero but an unknown member id
  • Error::LocalIllegalState if no transactional.id has been configured or no transaction has been started
  • A producer-fenced error if another producer with the same transactional.id is active
  • Error::UnsupportedVersion as 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.id or 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::Timeout if sending the offsets takes longer than max.block.ms
Source

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::LocalIllegalState if no transactional.id has been configured or no transaction has been started
  • A producer-fenced error if another producer with the same transactional.id is active
  • Error::UnsupportedVersion as a fatal error indicating the broker does not support transactions
  • An authorization error indicating that the configured transactional.id is 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::Timeout if committing takes longer than max.block.ms
Source

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::LocalIllegalState if no transactional.id has been configured or no transaction has been started
  • A producer-fenced error if another producer with the same transactional.id is active
  • An invalid-producer-epoch error if the producer has attempted to produce with an old epoch to the partition leader
  • Error::UnsupportedVersion as a fatal error indicating the broker does not support transactions
  • An authorization error indicating that the configured transactional.id is not authorized
  • Any previous fatal error the producer has encountered
  • Error::Timeout if aborting takes longer than max.block.ms
Source§

impl KafkaProducer<Vec<u8>, Vec<u8>>

Source

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>

Source§

fn drop(&mut self)

Executes the destructor for this type. Read more
Source§

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

Source§

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>

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>

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 commit_transaction(&self) -> Result<(), Error>

Commits the ongoing transaction.

Source§

async fn abort_transaction(&self) -> Result<(), Error>

Aborts the ongoing transaction.

Source§

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>

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>

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>>

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>

Get the partition metadata for the given topic.

Source§

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>

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> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<K, V, P> DynProducer<K, V> for P
where 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,

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 + '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,

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

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,

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

§

fn vzip(self) -> V