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?

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, and delivery.timeout.ms.

  • The API mirrors the Java client, with producer, consumer, and admin types such as KafkaProducer, KafkaConsumer, and KafkaAdminClient.

  • Operations that perform network I/O are async functions that run on the Tokio runtime.

Requirements

  • Rust 1.95 or later. The client uses the Rust 2024 edition.

  • The Tokio async runtime.

  • Kafka 4.0 or later, Confluent Platform 8.0 or later, or Confluent Cloud.

  • Linux amd64 or macOS arm64.

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.