Jump to content

Connect SuperML | Leeroopedia MCP: Equip your AI agents with best practices, code verification, and debugging knowledge. Powered by Leeroo — building Organizational Superintelligence. Contact us at founders@leeroo.com.

Implementation:Fede1024 Rust rdkafka Statistics Structs

From Leeroopedia


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

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")
}

Related Pages

Page Connections

Double-click a node to navigate. Hold to expand connections.
Principle
Implementation
Heuristic
Environment