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
impl ConsumerConfig
Sourcepub const GROUP_ID_CONFIG: &'static str = "group.id"
pub const GROUP_ID_CONFIG: &'static str = "group.id"
Config key: group.id.
Sourcepub const GROUP_INSTANCE_ID_CONFIG: &'static str = "group.instance.id"
pub const GROUP_INSTANCE_ID_CONFIG: &'static str = "group.instance.id"
Config key: group.instance.id.
Sourcepub const GROUP_PROTOCOL_CONFIG: &'static str = "group.protocol"
pub const GROUP_PROTOCOL_CONFIG: &'static str = "group.protocol"
Config key: group.protocol.
Sourcepub const DEFAULT_GROUP_PROTOCOL: &'static str = "classic"
pub const DEFAULT_GROUP_PROTOCOL: &'static str = "classic"
Default value of group.protocol. Matches Java’s DEFAULT_GROUP_PROTOCOL.
Sourcepub const GROUP_REMOTE_ASSIGNOR_CONFIG: &'static str = "group.remote.assignor"
pub const GROUP_REMOTE_ASSIGNOR_CONFIG: &'static str = "group.remote.assignor"
Config key: group.remote.assignor.
Sourcepub const MAX_POLL_RECORDS_CONFIG: &'static str = "max.poll.records"
pub const MAX_POLL_RECORDS_CONFIG: &'static str = "max.poll.records"
Config key: max.poll.records.
Sourcepub const DEFAULT_MAX_POLL_RECORDS: i32 = 500
pub const DEFAULT_MAX_POLL_RECORDS: i32 = 500
Default value of max.poll.records.
Sourcepub const MAX_POLL_INTERVAL_MS_CONFIG: &'static str = "max.poll.interval.ms"
pub const MAX_POLL_INTERVAL_MS_CONFIG: &'static str = "max.poll.interval.ms"
Config key: max.poll.interval.ms.
Sourcepub const SESSION_TIMEOUT_MS_CONFIG: &'static str = "session.timeout.ms"
pub const SESSION_TIMEOUT_MS_CONFIG: &'static str = "session.timeout.ms"
Config key: session.timeout.ms.
Sourcepub const HEARTBEAT_INTERVAL_MS_CONFIG: &'static str = "heartbeat.interval.ms"
pub const HEARTBEAT_INTERVAL_MS_CONFIG: &'static str = "heartbeat.interval.ms"
Config key: heartbeat.interval.ms.
Sourcepub const BOOTSTRAP_SERVERS_CONFIG: &'static str = "bootstrap.servers"
pub const BOOTSTRAP_SERVERS_CONFIG: &'static str = "bootstrap.servers"
Config key: bootstrap.servers.
Sourcepub const CLIENT_DNS_LOOKUP_CONFIG: &'static str = CommonClientConfigs::CLIENT_DNS_LOOKUP_CONFIG
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.
Sourcepub const CLIENT_ID_CONFIG: &'static str = "client.id"
pub const CLIENT_ID_CONFIG: &'static str = "client.id"
Config key: client.id.
Sourcepub const CLIENT_RACK_CONFIG: &'static str = "client.rack"
pub const CLIENT_RACK_CONFIG: &'static str = "client.rack"
Config key: client.rack.
Sourcepub const ENABLE_AUTO_COMMIT_CONFIG: &'static str = "enable.auto.commit"
pub const ENABLE_AUTO_COMMIT_CONFIG: &'static str = "enable.auto.commit"
Config key: enable.auto.commit.
Sourcepub const AUTO_COMMIT_INTERVAL_MS_CONFIG: &'static str = "auto.commit.interval.ms"
pub const AUTO_COMMIT_INTERVAL_MS_CONFIG: &'static str = "auto.commit.interval.ms"
Config key: auto.commit.interval.ms.
Sourcepub const PARTITION_ASSIGNMENT_STRATEGY_CONFIG: &'static str = "partition.assignment.strategy"
pub const PARTITION_ASSIGNMENT_STRATEGY_CONFIG: &'static str = "partition.assignment.strategy"
Config key: partition.assignment.strategy.
Sourcepub const AUTO_OFFSET_RESET_CONFIG: &'static str = "auto.offset.reset"
pub const AUTO_OFFSET_RESET_CONFIG: &'static str = "auto.offset.reset"
Config key: auto.offset.reset.
Sourcepub const FETCH_MIN_BYTES_CONFIG: &'static str = "fetch.min.bytes"
pub const FETCH_MIN_BYTES_CONFIG: &'static str = "fetch.min.bytes"
Config key: fetch.min.bytes.
Sourcepub const DEFAULT_FETCH_MIN_BYTES: i32 = 1
pub const DEFAULT_FETCH_MIN_BYTES: i32 = 1
Default value of fetch.min.bytes.
Sourcepub const FETCH_MAX_BYTES_CONFIG: &'static str = "fetch.max.bytes"
pub const FETCH_MAX_BYTES_CONFIG: &'static str = "fetch.max.bytes"
Config key: fetch.max.bytes.
Sourcepub const DEFAULT_FETCH_MAX_BYTES: i32
pub const DEFAULT_FETCH_MAX_BYTES: i32
Default value of fetch.max.bytes.
Sourcepub const FETCH_MAX_WAIT_MS_CONFIG: &'static str = "fetch.max.wait.ms"
pub const FETCH_MAX_WAIT_MS_CONFIG: &'static str = "fetch.max.wait.ms"
Config key: fetch.max.wait.ms.
Sourcepub const DEFAULT_FETCH_MAX_WAIT_MS: i32 = 500
pub const DEFAULT_FETCH_MAX_WAIT_MS: i32 = 500
Default value of fetch.max.wait.ms.
Sourcepub const MAX_PARTITION_FETCH_BYTES_CONFIG: &'static str = "max.partition.fetch.bytes"
pub const MAX_PARTITION_FETCH_BYTES_CONFIG: &'static str = "max.partition.fetch.bytes"
Config key: max.partition.fetch.bytes.
Sourcepub const DEFAULT_MAX_PARTITION_FETCH_BYTES: i32
pub const DEFAULT_MAX_PARTITION_FETCH_BYTES: i32
Default value of max.partition.fetch.bytes.
Sourcepub const CHECK_CRCS_CONFIG: &'static str = "check.crcs"
pub const CHECK_CRCS_CONFIG: &'static str = "check.crcs"
Config key: check.crcs.
Sourcepub const SEND_BUFFER_CONFIG: &'static str = "send.buffer.bytes"
pub const SEND_BUFFER_CONFIG: &'static str = "send.buffer.bytes"
Config key: send.buffer.bytes.
Sourcepub const RECEIVE_BUFFER_CONFIG: &'static str = "receive.buffer.bytes"
pub const RECEIVE_BUFFER_CONFIG: &'static str = "receive.buffer.bytes"
Config key: receive.buffer.bytes.
Sourcepub const SOCKET_CONNECTION_SETUP_TIMEOUT_MS_CONFIG: &'static str = "socket.connection.setup.timeout.ms"
pub const SOCKET_CONNECTION_SETUP_TIMEOUT_MS_CONFIG: &'static str = "socket.connection.setup.timeout.ms"
Config key: socket.connection.setup.timeout.ms.
Sourcepub const SOCKET_CONNECTION_SETUP_TIMEOUT_MAX_MS_CONFIG: &'static str = "socket.connection.setup.timeout.max.ms"
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.
Sourcepub const CONNECTIONS_MAX_IDLE_MS_CONFIG: &'static str = "connections.max.idle.ms"
pub const CONNECTIONS_MAX_IDLE_MS_CONFIG: &'static str = "connections.max.idle.ms"
Config key: connections.max.idle.ms.
Sourcepub const RECONNECT_BACKOFF_MS_CONFIG: &'static str = "reconnect.backoff.ms"
pub const RECONNECT_BACKOFF_MS_CONFIG: &'static str = "reconnect.backoff.ms"
Config key: reconnect.backoff.ms.
Sourcepub const RECONNECT_BACKOFF_MAX_MS_CONFIG: &'static str = "reconnect.backoff.max.ms"
pub const RECONNECT_BACKOFF_MAX_MS_CONFIG: &'static str = "reconnect.backoff.max.ms"
Config key: reconnect.backoff.max.ms.
Sourcepub const RETRY_BACKOFF_MS_CONFIG: &'static str = "retry.backoff.ms"
pub const RETRY_BACKOFF_MS_CONFIG: &'static str = "retry.backoff.ms"
Config key: retry.backoff.ms.
Sourcepub const RETRY_BACKOFF_MAX_MS_CONFIG: &'static str = "retry.backoff.max.ms"
pub const RETRY_BACKOFF_MAX_MS_CONFIG: &'static str = "retry.backoff.max.ms"
Config key: retry.backoff.max.ms.
Sourcepub const REQUEST_TIMEOUT_MS_CONFIG: &'static str = "request.timeout.ms"
pub const REQUEST_TIMEOUT_MS_CONFIG: &'static str = "request.timeout.ms"
Config key: request.timeout.ms.
Sourcepub const DEFAULT_API_TIMEOUT_MS_CONFIG: &'static str = "default.api.timeout.ms"
pub const DEFAULT_API_TIMEOUT_MS_CONFIG: &'static str = "default.api.timeout.ms"
Config key: default.api.timeout.ms.
Sourcepub const METADATA_MAX_AGE_CONFIG: &'static str = "metadata.max.age.ms"
pub const METADATA_MAX_AGE_CONFIG: &'static str = "metadata.max.age.ms"
Config key: metadata.max.age.ms.
Sourcepub const METADATA_RECOVERY_STRATEGY_CONFIG: &'static str = "metadata.recovery.strategy"
pub const METADATA_RECOVERY_STRATEGY_CONFIG: &'static str = "metadata.recovery.strategy"
Config key: metadata.recovery.strategy.
Sourcepub const METADATA_RECOVERY_REBOOTSTRAP_TRIGGER_MS_CONFIG: &'static str = "metadata.recovery.rebootstrap.trigger.ms"
pub const METADATA_RECOVERY_REBOOTSTRAP_TRIGGER_MS_CONFIG: &'static str = "metadata.recovery.rebootstrap.trigger.ms"
Config key: metadata.recovery.rebootstrap.trigger.ms.
Sourcepub const EXCLUDE_INTERNAL_TOPICS_CONFIG: &'static str = "exclude.internal.topics"
pub const EXCLUDE_INTERNAL_TOPICS_CONFIG: &'static str = "exclude.internal.topics"
Config key: exclude.internal.topics.
Sourcepub const DEFAULT_EXCLUDE_INTERNAL_TOPICS: bool = true
pub const DEFAULT_EXCLUDE_INTERNAL_TOPICS: bool = true
Default value of exclude.internal.topics.
Sourcepub const THROW_ON_FETCH_STABLE_OFFSET_UNSUPPORTED: &'static str = "internal.throw.on.fetch.stable.offset.unsupported"
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.
Sourcepub const ISOLATION_LEVEL_CONFIG: &'static str = "isolation.level"
pub const ISOLATION_LEVEL_CONFIG: &'static str = "isolation.level"
Config key: isolation.level.
Sourcepub const ALLOW_AUTO_CREATE_TOPICS_CONFIG: &'static str = "allow.auto.create.topics"
pub const ALLOW_AUTO_CREATE_TOPICS_CONFIG: &'static str = "allow.auto.create.topics"
Config key: allow.auto.create.topics.
Sourcepub const DEFAULT_ALLOW_AUTO_CREATE_TOPICS: bool = true
pub const DEFAULT_ALLOW_AUTO_CREATE_TOPICS: bool = true
Default value of allow.auto.create.topics.
Sourcepub const ENABLE_METRICS_PUSH_CONFIG: &'static str = "enable.metrics.push"
pub const ENABLE_METRICS_PUSH_CONFIG: &'static str = "enable.metrics.push"
Config key: enable.metrics.push.
Sourcepub const METRICS_SAMPLE_WINDOW_MS_CONFIG: &'static str = "metrics.sample.window.ms"
pub const METRICS_SAMPLE_WINDOW_MS_CONFIG: &'static str = "metrics.sample.window.ms"
Config key: metrics.sample.window.ms.
Sourcepub const METRICS_NUM_SAMPLES_CONFIG: &'static str = "metrics.num.samples"
pub const METRICS_NUM_SAMPLES_CONFIG: &'static str = "metrics.num.samples"
Config key: metrics.num.samples.
Sourcepub const METRICS_RECORDING_LEVEL_CONFIG: &'static str = "metrics.recording.level"
pub const METRICS_RECORDING_LEVEL_CONFIG: &'static str = "metrics.recording.level"
Config key: metrics.recording.level.
Sourcepub const SHARE_ACKNOWLEDGEMENT_MODE_CONFIG: &'static str = "share.acknowledgement.mode"
pub const SHARE_ACKNOWLEDGEMENT_MODE_CONFIG: &'static str = "share.acknowledgement.mode"
Config key: share.acknowledgement.mode.
Sourcepub const SHARE_ACQUIRE_MODE_CONFIG: &'static str = "share.acquire.mode"
pub const SHARE_ACQUIRE_MODE_CONFIG: &'static str = "share.acquire.mode"
Config key: share.acquire.mode.
Sourcepub const SECURITY_PROVIDERS_CONFIG: &'static str = "security.providers"
pub const SECURITY_PROVIDERS_CONFIG: &'static str = "security.providers"
Config key: security.providers.
Sourcepub const SECURITY_PROTOCOL_CONFIG: &'static str = "security.protocol"
pub const SECURITY_PROTOCOL_CONFIG: &'static str = "security.protocol"
Config key: security.protocol.
Sourcepub const SASL_MECHANISM_CONFIG: &'static str = SaslConfigs::SASL_MECHANISM
pub const SASL_MECHANISM_CONFIG: &'static str = SaslConfigs::SASL_MECHANISM
Config key: sasl.mechanism.
Sourcepub const SASL_JAAS_CONFIG: &'static str = SaslConfigs::SASL_JAAS_CONFIG
pub const SASL_JAAS_CONFIG: &'static str = SaslConfigs::SASL_JAAS_CONFIG
Config key: sasl.jaas.config.
Sourcepub const CONFIG_PROVIDERS_CONFIG: &'static str = "config.providers"
pub const CONFIG_PROVIDERS_CONFIG: &'static str = "config.providers"
Config key: config.providers.
Sourcepub fn bootstrap_servers(&self) -> &[String]
pub fn bootstrap_servers(&self) -> &[String]
bootstrap.servers.
Sourcepub fn group_instance_id(&self) -> Option<&str>
pub fn group_instance_id(&self) -> Option<&str>
group.instance.id, if any.
Sourcepub fn group_protocol(&self) -> &str
pub fn group_protocol(&self) -> &str
group.protocol.
Sourcepub fn group_remote_assignor(&self) -> Option<&str>
pub fn group_remote_assignor(&self) -> Option<&str>
group.remote.assignor, if any.
Sourcepub fn max_poll_records(&self) -> i32
pub fn max_poll_records(&self) -> i32
max.poll.records.
Sourcepub fn max_poll_interval_ms(&self) -> i32
pub fn max_poll_interval_ms(&self) -> i32
max.poll.interval.ms.
Sourcepub fn session_timeout_ms(&self) -> i32
pub fn session_timeout_ms(&self) -> i32
session.timeout.ms.
Sourcepub fn heartbeat_interval_ms(&self) -> i32
pub fn heartbeat_interval_ms(&self) -> i32
heartbeat.interval.ms.
Sourcepub fn enable_auto_commit(&self) -> bool
pub fn enable_auto_commit(&self) -> bool
enable.auto.commit.
Sourcepub fn auto_commit_interval_ms(&self) -> i32
pub fn auto_commit_interval_ms(&self) -> i32
auto.commit.interval.ms.
Sourcepub fn auto_offset_reset(&self) -> &str
pub fn auto_offset_reset(&self) -> &str
auto.offset.reset.
Sourcepub fn partition_assignment_strategy(&self) -> &[String]
pub fn partition_assignment_strategy(&self) -> &[String]
partition.assignment.strategy.
Sourcepub fn security_protocol(&self) -> &str
pub fn security_protocol(&self) -> &str
security.protocol - the protocol name (e.g. "PLAINTEXT", "SASL_SSL").
Sourcepub fn metadata_recovery_strategy(&self) -> &str
pub fn metadata_recovery_strategy(&self) -> &str
metadata.recovery.strategy.
Sourcepub fn throw_on_fetch_stable_offset_unsupported(&self) -> bool
pub fn throw_on_fetch_stable_offset_unsupported(&self) -> bool
internal.throw.on.fetch.stable.offset.unsupported.
Sourcepub fn request_timeout_ms(&self) -> i32
pub fn request_timeout_ms(&self) -> i32
request.timeout.ms.
Sourcepub fn retry_backoff_ms(&self) -> i64
pub fn retry_backoff_ms(&self) -> i64
retry.backoff.ms.
Sourcepub fn retry_backoff_max_ms(&self) -> i64
pub fn retry_backoff_max_ms(&self) -> i64
retry.backoff.max.ms.
Sourcepub fn set_bootstrap_servers(self, bootstrap_servers: Vec<String>) -> Self
pub fn set_bootstrap_servers(self, bootstrap_servers: Vec<String>) -> Self
Set bootstrap.servers.
Sourcepub fn set_client_id(self, client_id: impl Into<String>) -> Self
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.
Sourcepub fn set_group_id(self, group_id: impl Into<String>) -> Self
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.
Sourcepub fn set_group_protocol(self, protocol: impl Into<String>) -> Self
pub fn set_group_protocol(self, protocol: impl Into<String>) -> Self
Set group.protocol.
Sourcepub fn set_auto_offset_reset(self, value: impl Into<String>) -> Self
pub fn set_auto_offset_reset(self, value: impl Into<String>) -> Self
Set auto.offset.reset.
Sourcepub fn set_enable_auto_commit(self, value: bool) -> Self
pub fn set_enable_auto_commit(self, value: bool) -> Self
Set enable.auto.commit.
Sourcepub fn set_request_timeout_ms(self, value: i32) -> Self
pub fn set_request_timeout_ms(self, value: i32) -> Self
Set request.timeout.ms.
Sourcepub fn set_retry_backoff_ms(self, value: i64) -> Self
pub fn set_retry_backoff_ms(self, value: i64) -> Self
Set retry.backoff.ms.
Sourcepub fn new(props: &HashMap<String, String>) -> Result<Self, Error>
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
impl Clone for ConsumerConfig
Source§fn clone(&self) -> ConsumerConfig
fn clone(&self) -> ConsumerConfig
1.0.0 · Source§fn clone_from(&mut self, source: &Self)
fn clone_from(&mut self, source: &Self)
source. Read more