mz_test_util/kafka/
kafka_client.rs1use std::time::Duration;
13
14use anyhow::Context;
15use mz_kafka_util::admin::EnsureTopicConfig;
16use mz_kafka_util::client::{
17 MzClientContext, create, create_new_client_config_simple, create_with_context,
18};
19use rdkafka::admin::{AdminClient, AdminOptions, NewTopic, TopicReplication};
20use rdkafka::error::KafkaError;
21use rdkafka::producer::{DeliveryFuture, FutureProducer, FutureRecord};
22
23pub struct KafkaClient {
24 producer: FutureProducer<MzClientContext>,
25 kafka_url: String,
26}
27
28impl KafkaClient {
29 pub fn new(
30 kafka_url: &str,
31 group_id: &str,
32 configs: &[(&str, &str)],
33 ) -> Result<KafkaClient, anyhow::Error> {
34 let mut config = create_new_client_config_simple();
35 config.set("bootstrap.servers", kafka_url);
36 config.set("group.id", group_id);
37 for (key, val) in configs {
38 config.set(*key, *val);
39 }
40
41 let producer = create_with_context(&config, MzClientContext::default())?;
42
43 Ok(KafkaClient {
44 producer,
45 kafka_url: kafka_url.to_string(),
46 })
47 }
48
49 pub async fn create_topic(
50 &self,
51 topic_name: &str,
52 partitions: i32,
53 replication: i32,
54 configs: &[(&str, &str)],
55 timeout: Option<Duration>,
56 ) -> Result<(), anyhow::Error> {
57 let mut config = create_new_client_config_simple();
58 config.set("bootstrap.servers", &self.kafka_url);
59
60 let client = create::<AdminClient<_>>(&config).expect("creating admin kafka client failed");
61
62 let admin_opts = AdminOptions::new().request_timeout(timeout);
63
64 let mut topic = NewTopic::new(topic_name, partitions, TopicReplication::Fixed(replication));
65 for (key, val) in configs {
66 topic = topic.set(key, val);
67 }
68
69 mz_kafka_util::admin::ensure_topic(&client, &admin_opts, &topic, EnsureTopicConfig::Check)
70 .await
71 .context(format!("creating Kafka topic: {}", topic_name))?;
72
73 Ok(())
74 }
75
76 pub fn send(&self, topic_name: &str, message: &[u8]) -> Result<DeliveryFuture, KafkaError> {
77 let record: FutureRecord<&Vec<u8>, _> = FutureRecord::to(topic_name)
78 .payload(message)
79 .timestamp(chrono::Utc::now().timestamp_millis());
80 self.producer.send_result(record).map_err(|(e, _message)| e)
81 }
82
83 pub fn send_key_value(
84 &self,
85 topic_name: &str,
86 key: &[u8],
87 message: Option<Vec<u8>>,
88 ) -> Result<DeliveryFuture, KafkaError> {
89 let mut record: FutureRecord<_, _> = FutureRecord::to(topic_name)
90 .key(key)
91 .timestamp(chrono::Utc::now().timestamp_millis());
92 if let Some(message) = &message {
93 record = record.payload(message);
94 }
95 self.producer.send_result(record).map_err(|(e, _message)| e)
96 }
97}