Implementation:Fede1024 Rust rdkafka Statistics Structs
| Knowledge Sources | |
|---|---|
| Domains | Monitoring, Observability, Serialization |
| Last Updated | 2026-02-07 19:00 GMT |
Overview
Strongly-typed Rust data structures for deserializing the JSON statistics emitted by librdkafka, covering client metrics, per-broker stats, per-topic/partition stats, consumer group state, and exactly-once semantics information.
Description
The statistics module defines a hierarchy of serde-annotated structs that mirror the JSON schema documented in librdkafka's STATISTICS.md. The top-level Statistics struct contains scalar fields for client identity and aggregate counters, plus HashMap-based collections for per-broker (Broker) and per-topic (Topic) statistics. Each Broker contains connection metrics, request counters, latency histograms (Window structs with percentiles), and partition assignments. Each Topic contains per-partition (Partition) offset tracking, consumer lag, and message counts. Optional sections include ConsumerGroup for rebalance tracking and ExactlyOnceSemantics for transactional producer state.
Usage
Use these types in a ClientContext::stats callback implementation to receive and process librdkafka statistics. Enable statistics by setting statistics.interval.ms in ClientConfig. The JSON string from the callback can be deserialized into the Statistics struct using serde_json::from_str. Use for monitoring consumer lag, broker health, producer queue depth, and transactional state.
Code Reference
Source Location
- Repository: Fede1024_Rust_rdkafka
- File: src/statistics.rs
- Lines: 1-725
Signature
#[derive(Deserialize, Serialize, Debug, Default, Clone)]
pub struct Statistics {
pub name: String,
pub client_id: String,
pub client_type: String,
pub ts: i64,
pub time: i64,
pub age: i64,
pub replyq: i64,
pub msg_cnt: u64,
pub msg_size: u64,
pub msg_max: u64,
pub msg_size_max: u64,
pub tx: i64,
pub tx_bytes: i64,
pub rx: i64,
pub rx_bytes: i64,
pub txmsgs: i64,
pub txmsg_bytes: i64,
pub rxmsgs: i64,
pub rxmsg_bytes: i64,
pub simple_cnt: i64,
pub metadata_cache_cnt: i64,
pub brokers: HashMap<String, Broker>,
pub topics: HashMap<String, Topic>,
pub cgrp: Option<ConsumerGroup>,
pub eos: Option<ExactlyOnceSemantics>,
}
#[derive(Deserialize, Serialize, Debug, Default, Clone)]
pub struct Broker {
pub name: String,
pub nodeid: i32,
pub state: String,
pub stateage: i64,
pub outbuf_cnt: i64,
pub waitresp_cnt: i64,
pub tx: u64,
pub rx: u64,
pub txerrs: u64,
pub rxerrs: u64,
pub int_latency: Option<Window>,
pub outbuf_latency: Option<Window>,
pub rtt: Option<Window>,
pub throttle: Option<Window>,
pub toppars: HashMap<String, TopicPartition>,
// ... additional fields
}
#[derive(Deserialize, Serialize, Debug, Default, Clone)]
pub struct Window {
pub min: i64,
pub max: i64,
pub avg: i64,
pub sum: i64,
pub cnt: i64,
pub stddev: i64,
pub p50: i64,
pub p75: i64,
pub p90: i64,
pub p95: i64,
pub p99: i64,
pub p99_99: i64,
pub outofrange: i64,
pub hdrsize: i64,
}
#[derive(Deserialize, Serialize, Debug, Default, Clone)]
pub struct Partition {
pub partition: i32,
pub broker: i32,
pub leader: i32,
pub consumer_lag: i64,
pub consumer_lag_stored: i64,
pub fetch_state: String,
pub committed_offset: i64,
pub stored_offset: i64,
pub hi_offset: i64,
pub lo_offset: i64,
// ... additional fields
}
#[derive(Deserialize, Serialize, Debug, Default, Clone)]
pub struct ConsumerGroup {
pub state: String,
pub join_state: String,
pub rebalance_age: i64,
pub rebalance_cnt: i64,
pub rebalance_reason: String,
pub assignment_size: i32,
// ... additional fields
}
#[derive(Deserialize, Serialize, Debug, Default, Clone)]
pub struct ExactlyOnceSemantics {
pub idemp_state: String,
pub txn_state: String,
pub txn_may_enq: bool,
pub producer_id: i64,
pub producer_epoch: i64,
// ... additional fields
}
Import
use rdkafka::statistics::Statistics;
// or from crate root:
use rdkafka::Statistics;
I/O Contract
Inputs
| Name | Type | Required | Description |
|---|---|---|---|
| JSON string | &str | Yes | Statistics JSON emitted by librdkafka callback |
Outputs
| Name | Type | Description |
|---|---|---|
| Statistics | Statistics struct | Deserialized top-level statistics with all nested data |
| Broker stats | HashMap<String, Broker> | Per-broker connection, latency, and throughput metrics |
| Topic/Partition stats | HashMap<String, Topic> | Per-topic and per-partition offset and lag metrics |
| Consumer group | Option<ConsumerGroup> | Rebalance state and assignment info (consumers only) |
| EOS state | Option<ExactlyOnceSemantics> | Transactional producer state (transactional producers only) |
Usage Examples
Deserializing Statistics in a Context Callback
use rdkafka::client::ClientContext;
use rdkafka::statistics::Statistics;
use rdkafka::config::ClientConfig;
use rdkafka::consumer::StreamConsumer;
struct StatsContext;
impl ClientContext for StatsContext {
fn stats(&self, statistics: Statistics) {
// Access aggregate metrics
println!("Messages produced: {}", statistics.txmsgs);
println!("Messages consumed: {}", statistics.rxmsgs);
// Check per-broker health
for (name, broker) in &statistics.brokers {
println!("Broker {}: state={}, rtt_avg={}us",
name, broker.state,
broker.rtt.as_ref().map_or(0, |w| w.avg));
}
// Monitor consumer lag
for (topic_name, topic) in &statistics.topics {
for (partition_id, partition) in &topic.partitions {
if partition.consumer_lag >= 0 {
println!("{}-{}: lag={}",
topic_name, partition_id, partition.consumer_lag);
}
}
}
// Check consumer group state
if let Some(cgrp) = &statistics.cgrp {
println!("Group state: {}, rebalances: {}",
cgrp.state, cgrp.rebalance_cnt);
}
}
}
fn create_monitored_consumer() -> StreamConsumer<StatsContext> {
ClientConfig::new()
.set("bootstrap.servers", "localhost:9092")
.set("group.id", "my-group")
.set("statistics.interval.ms", "5000")
.create_with_context(StatsContext)
.expect("Consumer creation failed")
}