Skip to main content

metrics/
metrics.rs

1//! Print producer/consumer metrics as log lines or Prometheus text.
2//!
3//! Broker on `KAFKA_BOOTSTRAP`. Topic `KAFKA_TOPIC` (default `partitionline`).
4//! Set `FORMAT=prom` for a minimal Prometheus exposition (no prom crate).
5
6use partitionline::{Consumer, ProduceRecord, Producer, ProducerConfig};
7
8#[tokio::main]
9async fn main() -> partitionline::Result<()> {
10    let bootstrap = std::env::var("KAFKA_BOOTSTRAP").unwrap_or_else(|_| "127.0.0.1:9092".into());
11    let topic = std::env::var("KAFKA_TOPIC").unwrap_or_else(|_| "partitionline".into());
12    let format = std::env::var("FORMAT").unwrap_or_else(|_| "log".into());
13
14    let producer = Producer::new(
15        ProducerConfig::bootstrap([bootstrap.clone()]).linger(std::time::Duration::ZERO),
16    )
17    .await?;
18    let md = producer
19        .send(ProduceRecord::to(topic.clone()).value(&b"hello"[..]))
20        .await?;
21    let pm = producer.metrics();
22
23    if format == "prom" {
24        println!("# HELP partitionline_produce_records_acked Produce records acked");
25        println!("# TYPE partitionline_produce_records_acked counter");
26        println!("partitionline_produce_records_acked {}", pm.records_acked);
27        println!("# HELP partitionline_produce_ack_p99_seconds Produce ack p99 latency");
28        println!("# TYPE partitionline_produce_ack_p99_seconds gauge");
29        println!(
30            "partitionline_produce_ack_p99_seconds {}",
31            (pm.ack_latency.p99_nanos as f64) / 1_000_000_000.0
32        );
33    } else {
34        println!(
35            "produced {}-{}@{} queued={} acked={} bytes={} ack_us={} p50_us={} p99_us={} topics={}",
36            md.topic,
37            md.partition,
38            md.offset,
39            pm.records_queued,
40            pm.records_acked,
41            pm.bytes_queued,
42            pm.ack_latency.mean_nanos().unwrap_or(0) / 1000,
43            pm.ack_latency.p50_nanos / 1000,
44            pm.ack_latency.p99_nanos / 1000,
45            pm.topics.len()
46        );
47    }
48    producer.close().await?;
49
50    let mut consumer = Consumer::connect(bootstrap).await?;
51    consumer.assign(&topic, 0, 0).await?;
52    let recs = consumer.fetch().await?;
53    let cm = consumer.metrics();
54    if format == "prom" {
55        println!("# HELP partitionline_fetch_rounds Fetch rounds completed");
56        println!("# TYPE partitionline_fetch_rounds counter");
57        println!("partitionline_fetch_rounds {}", cm.fetch_rounds);
58        println!("# HELP partitionline_fetch_p99_seconds Fetch round p99 latency");
59        println!("# TYPE partitionline_fetch_p99_seconds gauge");
60        println!(
61            "partitionline_fetch_p99_seconds {}",
62            (cm.fetch_latency.p99_nanos as f64) / 1_000_000_000.0
63        );
64    } else {
65        println!(
66            "fetched {} records rounds={} bytes={} errors={} fetch_us={} p50_us={} p99_us={} topics={}",
67            recs.len(),
68            cm.fetch_rounds,
69            cm.bytes_fetched,
70            cm.fetch_errors,
71            cm.fetch_latency.mean_nanos().unwrap_or(0) / 1000,
72            cm.fetch_latency.p50_nanos / 1000,
73            cm.fetch_latency.p99_nanos / 1000,
74            cm.topics.len()
75        );
76    }
77    consumer.close().await?;
78    Ok(())
79}