<a id="rust-client"></a>

# Rust Client for Apache Kafka

## Overview

The Confluent Rust client for Apache Kafka® is a
[native Rust library](https://github.com/confluentinc/kafka-clients/tree/master/rust)
that provides a producer, consumer, and admin client. The client is
compatible with Kafka brokers, Confluent Cloud, and Confluent Platform.



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](https://tokio.rs/) runtime.

<a id="rust-client-requirements"></a>

## Requirements

- [Rust](https://rust-lang.org/tools/install/) 1.95 or later.
  The client uses the Rust 2024 edition.
- The [Tokio](https://tokio.rs/) async runtime.
- Kafka 4.0 or later, Confluent Platform 8.0 or later, or Confluent Cloud.
- Linux `amd64` or macOS `arm64`.

<a id="installation-rust-client"></a>

## Installation

The client is available on [crates.io](https://crates.io/) as the
`confluent-kafka` crate. Add it and Tokio to your project using
`cargo`:

```bash
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:

```toml
[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](https://github.com/rustls/rustls).

<a id="configuration-rust-client"></a>

## Configuration

Configure each client with a `HashMap<String, String>` of configuration
properties, which you pass to `ProducerConfig::new`, `ConsumerConfig::new`,
or `AdminClientConfig::new`:

- [Producer configuration reference](/platform/current/installation/configuration/producer-configs.html)
- [Consumer configuration reference](/platform/current/installation/configuration/consumer-configs.html)
- [Admin client configuration reference](/platform/current/installation/configuration/admin-configs.html)

To connect to a Kafka cluster in Confluent Cloud, provide the cluster’s bootstrap
server and an API key and secret:

```rust
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.

<a id="example-code-rust-client"></a>

## 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`.

<a id="admin-rust-client"></a>

### 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.

```rust
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.

<a id="producer-rust-client"></a>

### 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.

### Send individually

This example waits for acknowledgement before sending the next record, which
limits throughput:

```rust
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(())
}
```

### Send in batches

This example sends records in batches, collecting the futures returned by
`send` and awaiting them after sending:

```rust
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(())
}
```

<a id="consumer-rust-client"></a>

### 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](https://cwiki.apache.org/confluence/display/KAFKA/KIP-848%3A+The+Next+Generation+of+the+Consumer+Rebalance+Protocol).

### Manual commit

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.

```rust
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 commit

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:

```rust
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()
            );
        }
    }
}
```

<a id="api-docs-rust-client"></a>

## API documentation

The Rust client API reference is generated from the source with `cargo doc`.
See [Kafka Rust Client API](api.md#confluent-kafka-rust-api).

For the changes in each release, see [Kafka Rust Client Changelog](changelog.md#rust-changelog).

## Related content

* Confluent Developer: [Apache Kafka 101](https://developer.confluent.io/learn-kafka/apache-kafka/)
* Rust documentation: [The Rust Programming Language](https://doc.rust-lang.org/book/)
* Tokio documentation: [Tokio: an asynchronous runtime for Rust](https://tokio.rs/)
