demml / demml/scouter

KakfaSettings and KafkaConfig use different env variables

Open
#108 0 comments 0 reactions 0 assignees View on GitHub
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

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.