Rust Client for Apache Kafka
Overview
The Confluent Rust client for Apache Kafka® is a native Rust library that provides a producer, consumer, and admin client. The client is compatible with Kafka brokers, Confluent Cloud, and Confluent Platform.
Ready to get started?
Sign up for Confluent Cloud, the fully managed cloud-native service for Apache Kafka® and get started for free using the Cloud quick start.
Download Confluent Platform, the self managed, enterprise-grade distribution of Apache Kafka and get started using the Confluent Platform quick start.
The Rust client is a native implementation of the Kafka protocol. Unlike the Python, Go, and .NET clients, which are bindings on top of the C/C++ client, the Rust client doesn’t have a C library dependency or a foreign function interface (FFI) layer.
Note
The Rust client is a Preview release.
A Preview feature is a Rust client component that is being introduced to gain early feedback from developers. Preview features can be used for evaluation and non-production testing purposes or to provide feedback to Confluent. The warranty, SLA, and Support Services provisions of your agreement with Confluent do not apply to Preview features. Confluent may discontinue providing preview releases of the Preview features at any time in Confluent’s’ sole discretion.
The Rust client follows the design of the Kafka Java Client:
Configuration properties use the same names as the Java Client, such as
bootstrap.servers,linger.ms, anddelivery.timeout.ms.The API mirrors the Java client, with producer, consumer, and admin types such as
KafkaProducer,KafkaConsumer, andKafkaAdminClient.Operations that perform network I/O are
asyncfunctions that run on the Tokio runtime.
Requirements
Installation
The client is available on crates.io as the
confluent-kafka crate. Add it and Tokio to your project using
cargo:
cargo add confluent-kafka@=0.1.0
cargo add tokio --features macros,rt-multi-thread
Because the API can change between Preview releases, pin the exact
version. You can also add the dependencies directly to your
Cargo.toml file:
[dependencies]
confluent-kafka = "=0.1.0"
tokio = { version = "1", features = ["macros", "rt-multi-thread"] }
In your code, import the crate as confluent_kafka.
The client includes built-in TLS support that uses rustls.
Configuration
Configure each client with a HashMap<String, String> of configuration
properties, which you pass to ProducerConfig::new, ConsumerConfig::new,
or AdminClientConfig::new:
To connect to a Kafka cluster in Confluent Cloud, provide the cluster’s bootstrap server and an API key and secret:
let props = HashMap::from([
("bootstrap.servers".to_string(), "<bootstrap_server>".to_string()),
("security.protocol".to_string(), "SASL_SSL".to_string()),
("sasl.mechanism".to_string(), "PLAIN".to_string()),
(
"sasl.jaas.config".to_string(),
"org.apache.kafka.common.security.plain.PlainLoginModule required \
username=\"<api_key>\" password=\"<api_secret>\";"
.to_string(),
),
]);
The client supports only SASL/PLAIN authentication.The client doesn’t support SASL/SCRAM, OAuth, and AWS IAM aren’t supported.
Example code
The following examples show how to create a topic, produce messages, and
consume messages. Each example is a complete main.rs file that you can
run against a local Kafka broker at localhost:9092.
Create a topic
To create a topic, create a KafkaAdminClient and call create_topics.
Admin operations return a result object that holds a future for each item in
the request. Call all to wait for every item to complete, as the following
example does, or check the per-item futures to handle each result separately.
use std::collections::HashMap;
use confluent_kafka::admin::{Admin, AdminClientConfig, KafkaAdminClient, NewTopic};
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let props = HashMap::from([
("bootstrap.servers".to_string(), "localhost:9092".to_string()),
]);
let config = AdminClientConfig::new(&props)?;
let admin = KafkaAdminClient::new(config)?;
let topic = NewTopic::with_num_partitions_replication_factor("my-topic", Some(3), Some(1));
admin.create_topics(&[topic]).all().get().await?;
println!("Created topic my-topic");
admin.close().await?;
Ok(())
}
Note
Confluent Cloud requires a replication factor of 3. If you are connecting to Confluent Cloud,
update the NewTopic call to use Some(3) for the replication factor.
Produce messages
Create a KafkaProducer with a serializer for the key and the value, then
call send, which returns a KafkaFuture for the delivery. Call get
on the future to wait for the broker to acknowledge the message and return its
metadata. Call close before exiting to deliver any buffered messages.
This example waits for acknowledgement before sending the next record, which limits throughput:
use std::collections::HashMap;
use confluent_kafka::common::serialization::StringSerializer;
use confluent_kafka::producer::{KafkaProducer, Producer, ProducerConfig, ProducerRecord};
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let props = HashMap::from([
("bootstrap.servers".to_string(), "localhost:9092".to_string()),
("acks".to_string(), "all".to_string()),
]);
let config = ProducerConfig::new(&props)?;
let producer = KafkaProducer::<String, String>::new(
config,
Box::new(StringSerializer),
Box::new(StringSerializer),
)?;
for i in 0..10 {
let record = ProducerRecord::with_key(
"my-topic".to_string(),
Some(format!("key-{i}")),
Some(format!("value-{i}")),
);
let ack = producer.send(record).await?;
let metadata = ack.get().await?;
println!(
"Delivered to {} [{}] at offset {}",
metadata.topic(),
metadata.partition(),
metadata.offset()
);
}
producer.close().await?;
Ok(())
}
This example sends records in batches, collecting the futures returned by
send and awaiting them after sending:
use std::collections::HashMap;
use confluent_kafka::common::serialization::StringSerializer;
use confluent_kafka::producer::{KafkaProducer, Producer, ProducerConfig, ProducerRecord};
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let props = HashMap::from([
("bootstrap.servers".to_string(), "localhost:9092".to_string()),
("acks".to_string(), "all".to_string()),
]);
let config = ProducerConfig::new(&props)?;
let producer = KafkaProducer::<String, String>::new(
config,
Box::new(StringSerializer),
Box::new(StringSerializer),
)?;
let mut i = 0;
for _batch in 0..10 {
let mut futures = Vec::new();
for _ in 0..10 {
let record = ProducerRecord::with_key(
"my-topic".to_string(),
Some(format!("key-{i}")),
Some(format!("value-{i}")),
);
futures.push(producer.send(record).await?);
i += 1;
}
for ack in futures {
let metadata = ack.get().await?;
println!(
"Delivered to {} [{}] at offset {}",
metadata.topic(),
metadata.partition(),
metadata.offset()
);
}
}
producer.close().await?;
Ok(())
}
Consume messages
Create a consumer with KafkaConsumer::new, subscribe to one or more topics,
and call poll in a loop. Each call to poll returns a batch of records.
The Rust client supports only the consumer group protocol, so you must set
group.protocol to consumer. If group.protocol is classic, the
client returns an error. If the broker doesn’t support the consumer group
protocol, the consumer hangs without a clear error. Use Kafka 4.0 or later,
Confluent Platform 8.0 or later, or Confluent Cloud. For more details, see KIP-848.
This example turns off the default automatic offset commits and calls commit_sync
after processing each batch, which gives at-least-once delivery. The consumer also
provides commit_async and variants that commit specific offsets.
use std::collections::HashMap;
use std::time::Duration;
use confluent_kafka::common::serialization::StringDeserializer;
use confluent_kafka::consumer::{ConsumerConfig, KafkaConsumer};
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let props = HashMap::from([
("bootstrap.servers".to_string(), "localhost:9092".to_string()),
("group.id".to_string(), "my-group".to_string()),
("group.protocol".to_string(), "consumer".to_string()),
("auto.offset.reset".to_string(), "earliest".to_string()),
("enable.auto.commit".to_string(), "false".to_string()),
]);
let config = ConsumerConfig::new(&props)?;
let mut consumer = KafkaConsumer::new::<String, String>(
config,
Box::new(StringDeserializer),
Box::new(StringDeserializer),
)?;
consumer.subscribe_with_topics(vec!["my-topic".to_string()]).await?;
loop {
let records = consumer.poll(Duration::from_millis(500)).await?;
for record in &records {
println!(
"{} [{}] at offset {}: key={:?} value={:?}",
record.topic(),
record.partition(),
record.offset(),
record.key(),
record.value()
);
}
if !records.is_empty() {
consumer.commit_sync().await?;
}
}
}
Automatic offset commits are on by default. To use them, leave
enable.auto.commit unset or set it to true, and omit the explicit
commit calls:
use std::collections::HashMap;
use std::time::Duration;
use confluent_kafka::common::serialization::StringDeserializer;
use confluent_kafka::consumer::{ConsumerConfig, KafkaConsumer};
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let props = HashMap::from([
("bootstrap.servers".to_string(), "localhost:9092".to_string()),
("group.id".to_string(), "my-group".to_string()),
("group.protocol".to_string(), "consumer".to_string()),
("auto.offset.reset".to_string(), "earliest".to_string()),
]);
let config = ConsumerConfig::new(&props)?;
let mut consumer = KafkaConsumer::new::<String, String>(
config,
Box::new(StringDeserializer),
Box::new(StringDeserializer),
)?;
consumer.subscribe_with_topics(vec!["my-topic".to_string()]).await?;
loop {
let records = consumer.poll(Duration::from_millis(500)).await?;
for record in &records {
println!(
"{} [{}] at offset {}: key={:?} value={:?}",
record.topic(),
record.partition(),
record.offset(),
record.key(),
record.value()
);
}
}
}
API documentation
The Rust client API reference is generated from the source with cargo doc.
See Kafka Rust Client API.
For the changes in each release, see Kafka Rust Client Changelog.