Kafka broker¶
ruststream-rdkafka is the Apache Kafka broker for the
RustStream messaging framework, backed by
rdkafka / librdkafka.
[dependencies]
ruststream = { version = "0.5", features = ["macros", "json"] }
ruststream-rdkafka = "0.5"
serde = { version = "1", features = ["derive"] }
A minimal service is one handler and one app function:
use ruststream::runtime::{App, AppInfo, HandlerResult, RustStream};
use ruststream::subscriber;
use ruststream_rdkafka::KafkaBroker;
use serde::Deserialize;
#[derive(Debug, Deserialize)]
struct Order {
id: u64,
}
#[subscriber("orders")]
async fn handle(order: &Order) -> HandlerResult {
println!("got order {}", order.id);
HandlerResult::Ack
}
#[ruststream::app]
fn app() -> impl App {
RustStream::new(AppInfo::new("orders", "0.1.0")).with_broker(
KafkaBroker::new(["localhost:9092"]).default_group("orders-svc"),
|b| {
b.include(handle);
},
)
}
The transport model¶
- A subscription is one consumer joining one consumer group on one topic.
KafkaTopicdescribes it; the bare-string#[subscriber("orders")]form uses the broker'sdefault_group(Kafka cannot subscribe without a group). - The outgoing message name is the destination topic, and the partition-key header becomes the record's native key, so Kafka itself keeps per-key ordering (see Publishing).
- Settlement follows Kafka's committed-position model: the
Commitmode picks between librdkafka auto-commit (Auto, the default) and precise per-message acknowledgement over a contiguous watermark (Tracked). See Topics and groups. - Configuration delegates to librdkafka: unset options mean librdkafka defaults, and raw
config(key, value)passthroughs on the broker, the producer, and the descriptor reach every property this crate does not surface as a typed option. - Lazy startup:
KafkaBroker::newis synchronous and I/O-free, so a service composes with the synchronous#[ruststream::app]builder; the real network work happens in the idempotent asyncBroker::connect.
Scaffold a service¶
cargo generate --git https://github.com/powersemmi/ruststream-rdkafka templates/kafka-topic --name my-service
The starter wires one Kafka broker with a default consumer group, a tracked-commit subscriber
with a retry/dead-letter pipeline and a published reply, and the #[ruststream::app] entry
point (run / asyncapi gen).
Guides¶
- Topics and groups - descriptors, start offsets, commit modes, keyed lanes.
- Publishing - producing to topics, record keys, delivery guarantees.
- Testing - the in-process test broker and the live-cluster suites.