Skip to main content

mz_test_util/kafka/
kafka_client.rs

1// Copyright Materialize, Inc. and contributors. All rights reserved.
2//
3// Use of this software is governed by the Business Source License
4// included in the LICENSE file.
5//
6// As of the Change Date specified in that file, in accordance with
7// the Business Source License, use of this software will be governed
8// by the Apache License, Version 2.0.
9
10//! Kafka topic management
11
12use 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}