A pure Rust Kafka client library built on Tokio async runtime. Supports SASL authentication (PLAIN, SCRAM-SHA-256, SCRAM-SHA-512, GSSAPI/Kerberos) and TLS encryption.
- Pure Rust - No C bindings or external dependencies required
- Async/Await - Built on Tokio for modern async Rust
- SASL Authentication - Supports PLAIN, SCRAM-SHA-256, SCRAM-SHA-512, and Kerberos (via GSSAPI)
- TLS Encryption - Secure connections with configurable TLS
- Layered Architecture - Clean separation between transport, protocol, and API layers
- High-Level API - Easy-to-use Producer, Consumer, and Admin interfaces
- Low-Level Access - Direct protocol access for advanced use cases
- Connection Pooling - Efficient broker connection management
- Metadata Caching - Automatic cluster metadata refresh
Add to your Cargo.toml:
[dependencies]
kafka_client = "0.5"use kafka_client::Client;
let client = Client::builder(vec!["127.0.0.1:9092".parse()?])
.with_client_id("my-app")
.build()
.await?;
// Query cluster metadata
let brokers = client.metadata().get_all_brokers().await;
println!("Discovered {} brokers", brokers.len());use kafka_client::{Client, ProducerConfig, ProducerRecord};
use bytes::Bytes;
let client = Client::builder(vec!["127.0.0.1:9092".parse()?])
.build()
.await?;
// Default config
let producer = client.producer_default().await;
// Or with custom config
let producer = client.producer(
ProducerConfig::new().with_acks(-1).with_retries(3)
).await;
// Send messages
let record = ProducerRecord::new("my-topic", Bytes::from("hello world"));
let metadata = producer.send(record).await?;
println!("Sent to partition {} at offset {}", metadata.partition, metadata.offset);
// Ensure all messages are delivered
producer.flush().await?;use kafka_client::{Client, ConsumerConfig};
let client = Client::builder(vec!["127.0.0.1:9092".parse()?])
.build()
.await?;
// Default group consumer
let mut consumer = client.consumer_default();
// Or with custom config
let mut consumer = client.consumer(
ConsumerConfig::new("my-consumer-group")
.with_earliest()
);
consumer.subscribe(vec!["my-topic".to_string()]).await?;
// Poll for messages
let records = consumer.poll_timeout(Duration::from_millis(5000)).await?;
for record in records {
println!("Received: {:?}", record.value);
}use kafka_client::{Client, admin::NewTopic};
let client = Client::builder(vec!["127.0.0.1:9092".parse()?])
.build()
.await?;
let admin = client.admin();
// Create a topic
admin.create_topic(&NewTopic::new("orders", 3, 3)).await?;
// List all topics
let topics = admin.list_topics().await?;
for t in &topics { println!("{}", t.name); }
// Describe cluster
let info = admin.describe_cluster().await?;
println!("{} brokers, controller: {:?}", info.brokers.len(), info.controller_id);use kafka_client::{Client, SaslMechanismType};
// PLAIN authentication
let client = Client::builder(vec!["127.0.0.1:9092".parse()?])
.with_sasl(SaslMechanismType::Plain, "username", "password")
.build()
.await?;
// SCRAM-SHA-256 authentication
let client = Client::builder(vec!["127.0.0.1:9092".parse()?])
.with_sasl(SaslMechanismType::ScramSha256, "username", "password")
.build()
.await?;Kafka calls this mechanism "GSSAPI" on the wire, but it is backed by Kerberos — the client acquires
a service ticket from a KDC and exchanges it with the broker via the GSS-API framework.
Uses a pure Rust Kerberos implementation — no system libkrb5 needed.
use kafka_client::{Client, KerberosCredentials};
let client = Client::builder(vec!["127.0.0.1:9096".parse()?])
.with_kerberos(
KerberosCredentials::new("client@EXAMPLE.COM")
.with_keytab("/path/to/client.keytab")
)
.with_kdc("kdc.example.com", 88)
.with_broker_hostname("kafka-broker.example.com")
.build()
.await?;Supported Kerberos encryption types: AES-128/256-CTS-HMAC-SHA1-96 (RFC 3962) and AES-128/256-CTS-HMAC-SHA256/384 (RFC 8009).
use kafka_client::{Client, TlsConfig};
let tls = TlsConfig {
domain: "kafka.example.com".to_string(),
..Default::default()
};
// TLS only
let client = Client::builder(vec!["127.0.0.1:9093".parse()?])
.with_tls_config(tls)
.build()
.await?;
// TLS + SASL
let client = Client::builder(vec!["127.0.0.1:9093".parse()?])
.with_sasl_tls(tls, SaslMechanismType::ScramSha256, "username", "password")
.build()
.await?;The library uses a layered architecture for clean separation of concerns:
┌─────────────────────────────────────────┐
│ High-Level API Layer │
│ (Client, Producer, Consumer, Admin) │
├─────────────────────────────────────────┤
│ Cluster Layer │
│ (ClusterClient, BrokerManager, │
│ MetadataCache) [crate-internal] │
├─────────────────────────────────────────┤
│ Connection Layer │
│ (Connection, ConnectionPool) │
├─────────────────────────────────────────┤
│ Protocol Layer │
│ (Wire Codec, Request/Response) │
├─────────────────────────────────────────┤
│ Transport Layer │
│ (TCP, TLS, SASL) │
└─────────────────────────────────────────┘
See the examples directory for complete working examples:
basic_connect.rs- Simple connection and metadata queryproduce_consume.rs- Complete produce/consume workflowsasl_auth.rs- SASL authentication with different mechanismstls_connect.rs- TLS encryption and TLS+SASLadmin_operations.rs- Topic management operationsraw_connection.rs- Low-level Connection APIkerberos.rs- Kerberos authentication via GSSAPI (requires KDC)
Run an example:
# Basic connection
cargo run --example basic_connect
# SASL authentication
SASL_MECHANISM=SCRAM-SHA-256 \
SASL_USERNAME=user \
SASL_PASSWORD=pass \
cargo run --example sasl_auth
# TLS connection
KAFKA_DOMAIN=kafka.example.com \
cargo run --example tls_connect- Rust 1.92 or later
- Kafka 1.0.0 or later (3-broker tests verified on Kafka 4.x)
- GSSAPI/Kerberos requires SaslHandshake v1 (Kafka 1.0.0+). Pre-1.0.0 "bare TCP" GSSAPI and SaslHandshake v0 are not supported.
Integration tests require a running Kafka broker:
# Using Docker
docker run -d --name kafka -p 9092:9092 apache/kafka:latest
# Run tests
cargo test --features integration_testsLicensed under either of:
- Apache License, Version 2.0 (LICENSE-APACHE or http://www.apache.org/licenses/LICENSE-2.0)
- MIT License (LICENSE-MIT or http://opensource.org/licenses/MIT)
at your option.
Contributions are welcome! Please feel free to submit a Pull Request.
- Fork the repository
- Create your feature branch (
git checkout -b feature/amazing-feature) - Commit your changes (
git commit -m 'Add amazing feature') - Push to the branch (
git push origin feature/amazing-feature) - Open a Pull Request
- zzzdong - GitHub - kuwater@163.com