1use 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}