KakfaSettings and KafkaConfig use different env variables
- Dominant language
- Rust
- Stars
- 13
- Forks
- 1
- PR merge metrics
- No merged PRs in 30d
Description
```rust
#[derive(Clone, Serialize)]
pub struct KafkaSettings {
pub brokers: String,
pub num_workers: usize,
pub topics: Vec,
pub group_id: String,
pub username: Option,
pub password: Option,
pub security_protocol: String,
pub sasl_mechanism: String,
pub offset_reset: String,
pub cert_location: Option,
}
impl KafkaSettings {
pub fn __str__(&self) -> String {
// serialize the struct to a string
ProfileFuncs::__str__(self)
}
}
impl std::fmt::Debug for KafkaSettings {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("KafkaSettings")
.field("brokers", &self.brokers)
.field("num_workers", &self.num_workers)
.field("topics", &self.topics)
.field("group_id", &self.group_id)
.field("username", &self.username)
.field("password", &self.password.as_ref().map(|_| "***"))
.field("security_protocol", &self.security_protocol)
.field("offset_reset", &self.offset_reset)
.field("sasl_mechanism", &self.sasl_mechanism)
.finish()
}
}
impl Default for KafkaSettings {
fn default() -> Self {
let brokers =
std::env::var("KAFKA_BROKERS").unwrap_or_else(|_| "localhost:9092".to_string());
let num_workers = std::env::var("KAFKA_WORKER_COUNT")
.unwrap_or_else(|_| "3".to_string())
.parse::()
.unwrap();
let topics = std::env::var("KAFKA_TOPIC")
.unwrap_or_else(|_| "scouter_monitoring".to_string())
.split(',')
.map(|s| s.to_string())
.collect();
let group_id = std::env::var("KAFKA_GROUP").unwrap_or("scouter".to_string());
let offset_reset = std::env::var("KAFKA_OFFSET_RESET")
.unwrap_or_else(|_| "earliest".to_string())
.to_string();
let username: Option = std::env::var("KAFKA_USERNAME").ok();
let password: Option = std::env::var("KAFKA_PASSWORD").ok();
let security_protocol = std::env::var("KAFKA_SECURITY_PROTOCOL")
.ok()
.unwrap_or_else(|| "SASL_SSL".to_string());
let sasl_mechanism = std::env::var("KAFKA_SASL_MECHANISM")
.ok()
.unwrap_or_else(|| "PLAIN".to_string());
let cert_location = std::env::var("KAFKA_CERT_LOCATION").ok();
Self {
brokers,
num_workers,
topics,
group_id,
username,
password,
security_protocol,
sasl_mechanism,
offset_reset,
cert_location,
}
}
}
fn add_kafka_security(config: &mut HashMap) -> Result<(), PyEventError> {
if !config.contains_key("sasl.username") || !config.contains_key("sasl.password") {
if let (Ok(sasl_username), Ok(sasl_password)) = (
env::var("KAFKA_SASL_USERNAME"),
env::var("KAFKA_SASL_PASSWORD"),
) {
config.insert("sasl.username".to_string(), sasl_username);
config.insert("sasl.password".to_string(), sasl_password);
config.insert("security.protocol".to_string(), "SASL_SSL".to_string());
config.insert("sasl.mechanism".to_string(), "PLAIN".to_string());
}
}
Ok(())
}
fn add_kafka_args(
brokers: String,
compression: CompressionType,
message_timeout: u64,
message_max_bytes: i32,
config: &mut HashMap,
) -> Result<(), PyEventError> {
config.insert("bootstrap.servers".to_string(), brokers);
config.insert("compression.type".to_string(), compression.to_string());
config.insert(
"message.timeout.ms".to_string(),
message_timeout.to_string(),
);
config.insert(
"message.max.bytes".to_string(),
message_max_bytes.to_string(),
);
Ok(())
}
#[pyclass]
#[derive(Clone, Debug)]
pub struct KafkaConfig {
#[pyo3(get, set)]
pub brokers: String,
#[pyo3(get, set)]
pub topic: String,
#[pyo3(get, set)]
pub compression_type: CompressionType,
#[pyo3(get, set)]
pub message_timeout_ms: u64,
#[pyo3(get, set)]
pub message_max_bytes: i32,
#[pyo3(get, set)]
pub log_level: LogLevel,
#[pyo3(get, set)]
pub config: HashMap,
#[pyo3(get, set)]
pub max_retries: i32,
#[pyo3(get)]
pub transport_type: TransportType,
}
#[pymethods]
#[allow(clippy::too_many_arguments)]
impl KafkaConfig {
#[new]
#[pyo3(signature = (brokers=None, topic=None, compression_type=CompressionType::Gzip.to_string(), message_timeout_ms=600000, message_max_bytes=2097164, log_level=LogLevel::Info, config=None, max_retries=3))]
pub fn new(
brokers: Option,
topic: Option,
compression_type: Option,
message_timeout_ms: Option,
message_max_bytes: Option,
log_level: Option,
config: Option<&Bound<'_, PyDict>>,
max_retries: Option,
) -> Result {
let brokers = brokers.unwrap_or_else(|| {
env::var("KAFKA_BROKERS").unwrap_or_else(|_| "localhost:9092".to_string())
});
let topic = topic.unwrap_or_else(|| {
env::var("KAFKA_TOPIC").unwrap_or_else(|_| "scouter_monitoring".to_string())
});
let compression_type =
CompressionType::from_str(&compression_type.unwrap_or("gzip".to_string()))?;
let message_timeout_ms = message_timeout_ms.unwrap_or(600_000);
let message_max_bytes = message_max_bytes.unwrap_or(2097164);
let mut config = match config {
Some(config) => config
.iter()
.map(|(k, v)| (k.to_string(), v.to_string()))
.collect(),
None => HashMap::new(),
};
add_kafka_security(&mut config)?;
add_kafka_args(
brokers.clone(),
compression_type.clone(),
message_timeout_ms,
message_max_bytes,
&mut config,
)?;
let log_level = if let Some(level) = log_level {
level
} else {
let env_var = env::var("LOG_LEVEL")
.unwrap_or_else(|_| "info".to_string())
.to_lowercase();
match env_var.as_str() {
"info" => LogLevel::Info,
"debug" => LogLevel::Debug,
"error" => LogLevel::Error,
"warn" => LogLevel::Warn,
"trace" => LogLevel::Trace,
_ => LogLevel::Info,
}
};
Ok(KafkaConfig {
brokers,
topic,
compression_type,
message_timeout_ms,
message_max_bytes,
log_level,
config,
max_retries: max_retries.unwrap_or(3),
transport_type: TransportType::Kafka,
})
}
}
```
As the title says, the env vars should align between settings (server) and kafka config (used in client)
Contributor guide
Assessment
This issue has not been assessed yet.