Skip to main content

ConsumerConfig

Struct ConsumerConfig 

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

Configuration for the Kafka Consumer.

Documentation for the underlying keys can be found in the Kafka documentation.

Corresponds to org.apache.kafka.clients.consumer.ConsumerConfig.

Implementations§

Source§

impl ConsumerConfig

Source

pub const GROUP_ID_CONFIG: &'static str = "group.id"

Config key: group.id.

Source

pub const GROUP_INSTANCE_ID_CONFIG: &'static str = "group.instance.id"

Config key: group.instance.id.

Source

pub const GROUP_PROTOCOL_CONFIG: &'static str = "group.protocol"

Config key: group.protocol.

Source

pub const DEFAULT_GROUP_PROTOCOL: &'static str = "classic"

Default value of group.protocol. Matches Java’s DEFAULT_GROUP_PROTOCOL.

Source

pub const GROUP_REMOTE_ASSIGNOR_CONFIG: &'static str = "group.remote.assignor"

Config key: group.remote.assignor.

Source

pub const MAX_POLL_RECORDS_CONFIG: &'static str = "max.poll.records"

Config key: max.poll.records.

Source

pub const DEFAULT_MAX_POLL_RECORDS: i32 = 500

Default value of max.poll.records.

Source

pub const MAX_POLL_INTERVAL_MS_CONFIG: &'static str = "max.poll.interval.ms"

Config key: max.poll.interval.ms.

Source

pub const SESSION_TIMEOUT_MS_CONFIG: &'static str = "session.timeout.ms"

Config key: session.timeout.ms.

Source

pub const HEARTBEAT_INTERVAL_MS_CONFIG: &'static str = "heartbeat.interval.ms"

Config key: heartbeat.interval.ms.

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 ConsumerConfig.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 CLIENT_RACK_CONFIG: &'static str = "client.rack"

Config key: client.rack.

Source

pub const ENABLE_AUTO_COMMIT_CONFIG: &'static str = "enable.auto.commit"

Config key: enable.auto.commit.

Source

pub const AUTO_COMMIT_INTERVAL_MS_CONFIG: &'static str = "auto.commit.interval.ms"

Config key: auto.commit.interval.ms.

Source

pub const PARTITION_ASSIGNMENT_STRATEGY_CONFIG: &'static str = "partition.assignment.strategy"

Config key: partition.assignment.strategy.

Source

pub const AUTO_OFFSET_RESET_CONFIG: &'static str = "auto.offset.reset"

Config key: auto.offset.reset.

Source

pub const FETCH_MIN_BYTES_CONFIG: &'static str = "fetch.min.bytes"

Config key: fetch.min.bytes.

Source

pub const DEFAULT_FETCH_MIN_BYTES: i32 = 1

Default value of fetch.min.bytes.

Source

pub const FETCH_MAX_BYTES_CONFIG: &'static str = "fetch.max.bytes"

Config key: fetch.max.bytes.

Source

pub const DEFAULT_FETCH_MAX_BYTES: i32

Default value of fetch.max.bytes.

Source

pub const FETCH_MAX_WAIT_MS_CONFIG: &'static str = "fetch.max.wait.ms"

Config key: fetch.max.wait.ms.

Source

pub const DEFAULT_FETCH_MAX_WAIT_MS: i32 = 500

Default value of fetch.max.wait.ms.

Source

pub const MAX_PARTITION_FETCH_BYTES_CONFIG: &'static str = "max.partition.fetch.bytes"

Config key: max.partition.fetch.bytes.

Source

pub const DEFAULT_MAX_PARTITION_FETCH_BYTES: i32

Default value of max.partition.fetch.bytes.

Source

pub const CHECK_CRCS_CONFIG: &'static str = "check.crcs"

Config key: check.crcs.

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 SOCKET_CONNECTION_SETUP_TIMEOUT_MS_CONFIG: &'static str = "socket.connection.setup.timeout.ms"

Config key: socket.connection.setup.timeout.ms.

Source

pub const SOCKET_CONNECTION_SETUP_TIMEOUT_MAX_MS_CONFIG: &'static str = "socket.connection.setup.timeout.max.ms"

Config key: socket.connection.setup.timeout.max.ms.

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 REQUEST_TIMEOUT_MS_CONFIG: &'static str = "request.timeout.ms"

Config key: request.timeout.ms.

Source

pub const DEFAULT_API_TIMEOUT_MS_CONFIG: &'static str = "default.api.timeout.ms"

Config key: default.api.timeout.ms.

Source

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

Config key: metadata.max.age.ms.

Source

pub const METADATA_RECOVERY_STRATEGY_CONFIG: &'static str = "metadata.recovery.strategy"

Config key: metadata.recovery.strategy.

Source

pub const METADATA_RECOVERY_REBOOTSTRAP_TRIGGER_MS_CONFIG: &'static str = "metadata.recovery.rebootstrap.trigger.ms"

Config key: metadata.recovery.rebootstrap.trigger.ms.

Source

pub const EXCLUDE_INTERNAL_TOPICS_CONFIG: &'static str = "exclude.internal.topics"

Config key: exclude.internal.topics.

Source

pub const DEFAULT_EXCLUDE_INTERNAL_TOPICS: bool = true

Default value of exclude.internal.topics.

Source

pub const THROW_ON_FETCH_STABLE_OFFSET_UNSUPPORTED: &'static str = "internal.throw.on.fetch.stable.offset.unsupported"

Config key: internal.throw.on.fetch.stable.offset.unsupported.

Source

pub const ISOLATION_LEVEL_CONFIG: &'static str = "isolation.level"

Config key: isolation.level.

Source

pub const ALLOW_AUTO_CREATE_TOPICS_CONFIG: &'static str = "allow.auto.create.topics"

Config key: allow.auto.create.topics.

Source

pub const DEFAULT_ALLOW_AUTO_CREATE_TOPICS: bool = true

Default value of allow.auto.create.topics.

Source

pub const ENABLE_METRICS_PUSH_CONFIG: &'static str = "enable.metrics.push"

Config key: enable.metrics.push.

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 SHARE_ACKNOWLEDGEMENT_MODE_CONFIG: &'static str = "share.acknowledgement.mode"

Config key: share.acknowledgement.mode.

Source

pub const SHARE_ACQUIRE_MODE_CONFIG: &'static str = "share.acquire.mode"

Config key: share.acquire.mode.

Source

pub const SECURITY_PROVIDERS_CONFIG: &'static str = "security.providers"

Config key: security.providers.

Source

pub const SECURITY_PROTOCOL_CONFIG: &'static str = "security.protocol"

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 const CONFIG_PROVIDERS_CONFIG: &'static str = "config.providers"

Config key: config.providers.

Source

pub fn bootstrap_servers(&self) -> &[String]

bootstrap.servers.

Source

pub fn client_id(&self) -> &str

client.id.

Source

pub fn group_id(&self) -> Option<&str>

group.id, if any.

Source

pub fn group_instance_id(&self) -> Option<&str>

group.instance.id, if any.

Source

pub fn group_protocol(&self) -> &str

group.protocol.

Source

pub fn group_remote_assignor(&self) -> Option<&str>

group.remote.assignor, if any.

Source

pub fn max_poll_records(&self) -> i32

max.poll.records.

Source

pub fn max_poll_interval_ms(&self) -> i32

max.poll.interval.ms.

Source

pub fn session_timeout_ms(&self) -> i32

session.timeout.ms.

Source

pub fn heartbeat_interval_ms(&self) -> i32

heartbeat.interval.ms.

Source

pub fn enable_auto_commit(&self) -> bool

enable.auto.commit.

Source

pub fn auto_commit_interval_ms(&self) -> i32

auto.commit.interval.ms.

Source

pub fn auto_offset_reset(&self) -> &str

auto.offset.reset.

Source

pub fn partition_assignment_strategy(&self) -> &[String]

partition.assignment.strategy.

Source

pub fn security_protocol(&self) -> &str

security.protocol - the protocol name (e.g. "PLAINTEXT", "SASL_SSL").

Source

pub fn metadata_recovery_strategy(&self) -> &str

metadata.recovery.strategy.

Source

pub fn throw_on_fetch_stable_offset_unsupported(&self) -> bool

internal.throw.on.fetch.stable.offset.unsupported.

Source

pub fn request_timeout_ms(&self) -> i32

request.timeout.ms.

Source

pub fn retry_backoff_ms(&self) -> i64

retry.backoff.ms.

Source

pub fn retry_backoff_max_ms(&self) -> i64

retry.backoff.max.ms.

Source

pub fn set_bootstrap_servers(self, bootstrap_servers: Vec<String>) -> Self

Set bootstrap.servers.

Source

pub fn set_client_id(self, client_id: impl Into<String>) -> Self

Set client.id, trimmed as ConsumerConfig::new trims it (Java’s ConfigDef.parseType), so a blank id is generated at construction.

Source

pub fn set_group_id(self, group_id: impl Into<String>) -> Self

Set group.id.

A client.id generated by ConsumerConfig::new already embeds the group id, so after changing the group either set client.id explicitly or clear it with set_client_id(""), which the consumer constructor then regenerates.

Trimmed as ConsumerConfig::new trims it (Java’s ConfigDef.parseType), so a blank id is rejected at construction.

Source

pub fn set_group_protocol(self, protocol: impl Into<String>) -> Self

Set group.protocol.

Source

pub fn set_auto_offset_reset(self, value: impl Into<String>) -> Self

Set auto.offset.reset.

Source

pub fn set_enable_auto_commit(self, value: bool) -> Self

Set enable.auto.commit.

Source

pub fn set_request_timeout_ms(self, value: i32) -> Self

Set request.timeout.ms.

Source

pub fn set_retry_backoff_ms(self, value: i64) -> Self

Set retry.backoff.ms.

Source

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

Parses a string-typed property map into a typed ConsumerConfig.

Mirrors Java’s new ConsumerConfig(Map<String, String>). Unknown keys are logged as warnings and ignored, matching Java’s behavior.

Construction-time validation matches the Java ConfigDef validators that we have translated so far. Cross-field validation (e.g. partition.assignment.strategy being forbidden under group.protocol=consumer) is deferred until a later phase per consumer-threading.md §20.

§Errors

Returns Error::LocalIllegalArgument if a value cannot be parsed for its expected type, or fails its validator.

Trait Implementations§

Source§

impl Clone for ConsumerConfig

Source§

fn clone(&self) -> ConsumerConfig

Returns a duplicate of the value. Read more
1.0.0 · Source§

fn clone_from(&mut self, source: &Self)

Performs copy-assignment from source. Read more
Source§

impl Debug for ConsumerConfig

Source§

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

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

impl Default for ConsumerConfig

Source§

fn default() -> Self

Default values match Java’s ConfigDef.define(...) second argument for every key.

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> CloneToUninit for T
where T: Clone,

Source§

unsafe fn clone_to_uninit(&self, dest: *mut u8)

🔬This is a nightly-only experimental API. (clone_to_uninit)
Performs copy-assignment from self to dest. 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> ToOwned for T
where T: Clone,

Source§

type Owned = T

The resulting type after obtaining ownership.
Source§

fn to_owned(&self) -> T

Creates owned data from borrowed data, usually by cloning. Read more
Source§

fn clone_into(&self, target: &mut T)

Uses borrowed data to replace owned data, usually by cloning. Read more
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