use std::collections::HashMap;
use std::time::Duration;
use rdkafka::ClientConfig;
use rdkafka::message::{Header, OwnedHeaders};
use rdkafka::producer::{FutureProducer, FutureRecord};
use crate::config::KafkaAuthConfig;
use crate::connector::lru_cache::LruCache;
use crate::errors::OrionError;
pub struct KafkaProducer {
producer: FutureProducer,
}
pub struct KafkaProducerCache {
global_brokers: String,
global: std::sync::Arc<KafkaProducer>,
auth: KafkaAuthConfig,
extra_config: HashMap<String, String>,
cache: LruCache<std::sync::Arc<KafkaProducer>>,
}
impl KafkaProducerCache {
pub fn new(
global_brokers: String,
global: std::sync::Arc<KafkaProducer>,
auth: KafkaAuthConfig,
extra_config: HashMap<String, String>,
max_entries: usize,
) -> Self {
Self {
global_brokers: normalize_brokers(&global_brokers),
global,
auth,
extra_config,
cache: LruCache::new(max_entries, "kafka_producer"),
}
}
pub async fn for_brokers(
&self,
brokers: &[String],
) -> Result<std::sync::Arc<KafkaProducer>, OrionError> {
if brokers.is_empty() {
return Ok(self.global.clone());
}
let key = normalize_brokers(&brokers.join(","));
if key == self.global_brokers {
return Ok(self.global.clone());
}
self.cache
.get_or_create(&key, || async {
KafkaProducer::new(&key, &self.auth, &self.extra_config).map(std::sync::Arc::new)
})
.await
}
}
fn normalize_brokers(brokers: &str) -> String {
let mut parts: Vec<&str> = brokers
.split(',')
.map(str::trim)
.filter(|s| !s.is_empty())
.collect();
parts.sort_unstable();
parts.join(",")
}
impl KafkaProducer {
pub fn new(
brokers: &str,
auth: &KafkaAuthConfig,
extra_config: &HashMap<String, String>,
) -> Result<Self, OrionError> {
let mut client_config = ClientConfig::new();
client_config
.set("bootstrap.servers", brokers)
.set("message.timeout.ms", "30000")
.set("enable.idempotence", "true")
.set("delivery.timeout.ms", "120000")
.set("acks", "all")
.set("compression.type", "lz4")
.set("linger.ms", "5")
.set("batch.size", "65536");
super::apply_client_auth(&mut client_config, auth, extra_config);
let producer: FutureProducer =
client_config.create().map_err(|e| OrionError::Internal {
context: "Failed to create Kafka producer".to_string(),
source: Some(Box::new(e)),
})?;
Ok(Self { producer })
}
pub async fn send(
&self,
topic: &str,
key: Option<&str>,
payload: &[u8],
) -> Result<(), OrionError> {
let mut record = FutureRecord::to(topic).payload(payload);
if let Some(k) = key {
record = record.key(k);
}
let headers = {
let mut trace_headers = std::collections::HashMap::new();
crate::server::trace_context::inject_trace_context(&mut trace_headers);
let mut kafka_headers = OwnedHeaders::new();
for (k, v) in &trace_headers {
kafka_headers = kafka_headers.insert(Header {
key: k,
value: Some(v.as_bytes()),
});
}
kafka_headers
};
let record = record.headers(headers);
self.producer
.send(record, Duration::from_secs(30))
.await
.map_err(|(e, _)| OrionError::Internal {
context: format!("Kafka send to '{topic}' failed"),
source: Some(Box::new(e)),
})?;
Ok(())
}
}