Skip to main content

ProducerConfig

Struct ProducerConfig 

Source
pub struct ProducerConfig { /* private fields */ }
Expand description

Configuration for the Kafka Producer.

Documentation for these configurations can be found in the Kafka documentation.

Corresponds to org.apache.kafka.clients.producer.ProducerConfig.

Implementations§

Source§

impl ProducerConfig

Source

pub const MAX_IN_FLIGHT_REQUESTS_FOR_IDEMPOTENCE: i32 = MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION_FOR_IDEMPOTENCE

Maximum number of in-flight requests per connection when idempotence is enabled.

Source

pub const BOOTSTRAP_SERVERS_CONFIG: &'static str = "bootstrap.servers"

Config key: bootstrap.servers

Source

pub const CLIENT_DNS_LOOKUP_CONFIG: &'static str = CommonClientConfigs::CLIENT_DNS_LOOKUP_CONFIG

Config key: client.dns.lookup. Java’s ProducerConfig.java declares it as its own public alias of CommonClientConfigs.CLIENT_DNS_LOOKUP_CONFIG.

Source

pub const CLIENT_ID_CONFIG: &'static str = "client.id"

Config key: client.id

Source

pub const BATCH_SIZE_CONFIG: &'static str = "batch.size"

Config key: batch.size

Source

pub const LINGER_MS_CONFIG: &'static str = "linger.ms"

Config key: linger.ms

Source

pub const BUFFER_MEMORY_CONFIG: &'static str = "buffer.memory"

Config key: buffer.memory

Source

pub const MAX_BLOCK_MS_CONFIG: &'static str = "max.block.ms"

Config key: max.block.ms

Source

pub const ACKS_CONFIG: &'static str = "acks"

Config key: acks

Source

pub const RETRIES_CONFIG: &'static str = "retries"

Config key: retries

Source

pub const DELIVERY_TIMEOUT_MS_CONFIG: &'static str = "delivery.timeout.ms"

Config key: delivery.timeout.ms

Source

pub const REQUEST_TIMEOUT_MS_CONFIG: &'static str = "request.timeout.ms"

Config key: request.timeout.ms

Source

pub const ENABLE_IDEMPOTENCE_CONFIG: &'static str = "enable.idempotence"

Config key: enable.idempotence

Source

pub const MAX_REQUEST_SIZE_CONFIG: &'static str = "max.request.size"

Config key: max.request.size

Source

pub const MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION: &'static str = "max.in.flight.requests.per.connection"

Config key: max.in.flight.requests.per.connection

Source

pub const COMPRESSION_TYPE_CONFIG: &'static str = "compression.type"

Config key: compression.type

Source

pub const CONNECTIONS_MAX_IDLE_MS_CONFIG: &'static str = "connections.max.idle.ms"

Config key: connections.max.idle.ms

Source

pub const RECONNECT_BACKOFF_MS_CONFIG: &'static str = "reconnect.backoff.ms"

Config key: reconnect.backoff.ms

Source

pub const RECONNECT_BACKOFF_MAX_MS_CONFIG: &'static str = "reconnect.backoff.max.ms"

Config key: reconnect.backoff.max.ms

Source

pub const RETRY_BACKOFF_MS_CONFIG: &'static str = "retry.backoff.ms"

Config key: retry.backoff.ms

Source

pub const RETRY_BACKOFF_MAX_MS_CONFIG: &'static str = "retry.backoff.max.ms"

Config key: retry.backoff.max.ms

Source

pub const SEND_BUFFER_CONFIG: &'static str = "send.buffer.bytes"

Config key: send.buffer.bytes

Source

pub const RECEIVE_BUFFER_CONFIG: &'static str = "receive.buffer.bytes"

Config key: receive.buffer.bytes

Source

pub const METADATA_MAX_AGE_CONFIG: &'static str = "metadata.max.age.ms"

Config key: metadata.max.age.ms

Source

pub const METADATA_MAX_IDLE_CONFIG: &'static str = "metadata.max.idle.ms"

Config key: metadata.max.idle.ms

Source

pub const PARTITIONER_ADAPTIVE_PARTITIONING_ENABLE_CONFIG: &'static str = "partitioner.adaptive.partitioning.enable"

Config key: partitioner.adaptive.partitioning.enable

Source

pub const PARTITIONER_AVAILABILITY_TIMEOUT_MS_CONFIG: &'static str = "partitioner.availability.timeout.ms"

Config key: partitioner.availability.timeout.ms

Source

pub const PARTITIONER_IGNORE_KEYS_CONFIG: &'static str = "partitioner.ignore.keys"

Config key: partitioner.ignore.keys

Source

pub const PARTITIONER_TYPE_CONFIG: &'static str = "partitioner.type"

Config key: partitioner.type

Source

pub const CONSISTENT_RANDOM_PARTITIONER: &'static str = "ConsistentRandomPartitioner"

Accepted partitioner.type value selecting the CRC-32 key hash (the crate-internal KeyHasher::Crc32) — the default, librdkafka consistent_random parity.

Source

pub const MURMUR2_RANDOM_PARTITIONER: &'static str = "Murmur2RandomPartitioner"

Accepted partitioner.type value selecting the murmur2 key hash (the crate-internal KeyHasher::Murmur2) — exact Java-client parity.

Source

pub const ROUND_ROBIN_PARTITIONER: &'static str = "RoundRobinPartitioner"

Accepted partitioner.type value selecting the RoundRobinPartitioner.

Source

pub const TRANSACTIONAL_ID_CONFIG: &'static str = "transactional.id"

Config key: transactional.id

Source

pub const TRANSACTION_TIMEOUT_CONFIG: &'static str = "transaction.timeout.ms"

Config key: transaction.timeout.ms

Source

pub const METRICS_SAMPLE_WINDOW_MS_CONFIG: &'static str = "metrics.sample.window.ms"

Config key: metrics.sample.window.ms

Source

pub const METRICS_NUM_SAMPLES_CONFIG: &'static str = "metrics.num.samples"

Config key: metrics.num.samples

Source

pub const METRICS_RECORDING_LEVEL_CONFIG: &'static str = "metrics.recording.level"

Config key: metrics.recording.level

Source

pub const TRANSACTION_TWO_PHASE_COMMIT_ENABLE_CONFIG: &'static str = "transaction.two.phase.commit.enable"

Config key: transaction.two.phase.commit.enable

Source

pub const SECURITY_PROTOCOL_CONFIG: &'static str = CommonClientConfigs::SECURITY_PROTOCOL_CONFIG

Config key: security.protocol

Source

pub const SASL_MECHANISM_CONFIG: &'static str = SaslConfigs::SASL_MECHANISM

Config key: sasl.mechanism

Source

pub const SASL_JAAS_CONFIG: &'static str = SaslConfigs::SASL_JAAS_CONFIG

Config key: sasl.jaas.config

Source

pub fn new(props: &HashMap<String, String>) -> Result<Self, Error>

Creates a ProducerConfig from a map of string key-value pairs.

This is the Rust equivalent of Java’s new ProducerConfig(Map<String, Object>) or new ProducerConfig(Properties). Starts with default values and overrides each field that has a matching entry in the map.

Unknown keys are logged as warnings and ignored, matching Java’s behavior.

§Errors

Returns Error::LocalIllegalArgument if a value cannot be parsed for its expected type (e.g., "abc" for an integer field).

Source

pub fn set_partitioner<K: 'static, V: 'static>( self, partitioner: Box<dyn Partitioner<K, V>>, ) -> Self

Sets the Partitioner that determines which partition each record goes to: the value of Java’s partitioner.class.

Java names a class, which the producer instantiates by reflection; here the partitioner itself is the value. KafkaProducer::new configures it with the user configs plus the resolved client.id, as Java does (KafkaProducer.java:381-388), and closes it on close. As with any configured partitioner, adaptive partitioning is then disabled. The producer built from this config takes ownership of the partitioner.

This replaces a partitioner.type given in the properties passed to new, as a later value for the same key would.

The partitioner’s K / V must be the record types of the producer built from this config; otherwise KafkaProducer::new fails with “<partitioner> is not an instance of <expected>”.

Trait Implementations§

Source§

impl Debug for ProducerConfig

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more
Source§

impl Default for ProducerConfig

Source§

fn default() -> Self

Returns the “default value” for a type. Read more

Auto Trait Implementations§

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