laminar_connectors/kafka/
sink_metrics.rs1use prometheus::{IntCounter, IntGauge, Registry};
4
5use crate::prom::reg_or_local;
6
7#[derive(Debug, Clone)]
9pub struct KafkaSinkMetrics {
10 pub records_written: IntCounter,
12 pub bytes_written: IntCounter,
14 pub errors_total: IntCounter,
16 pub dlq_records: IntCounter,
18 pub serialization_errors: IntCounter,
20 pub produce_latency_sum_us: IntCounter,
22 pub produce_latency_max_us: IntGauge,
24 pub produce_latency_count: IntCounter,
26}
27
28impl KafkaSinkMetrics {
29 #[must_use]
31 #[allow(clippy::missing_panics_doc)]
32 pub fn new(registry: Option<&Registry>) -> Self {
33 let mut local = None;
34 let reg = reg_or_local(registry, &mut local);
35
36 Self {
37 records_written: reg.counter(
38 "kafka_sink_records_written_total",
39 "Records written to Kafka",
40 ),
41 bytes_written: reg.counter("kafka_sink_bytes_written_total", "Bytes written to Kafka"),
42 errors_total: reg.counter("kafka_sink_errors_total", "Kafka sink errors"),
43 dlq_records: reg.counter("kafka_sink_dlq_records_total", "Records routed to DLQ"),
44 serialization_errors: reg.counter(
45 "kafka_sink_serialization_errors_total",
46 "Serialization errors",
47 ),
48 produce_latency_sum_us: reg.counter(
49 "kafka_sink_produce_latency_sum_us",
50 "Sum of produce latencies (us)",
51 ),
52 produce_latency_count: reg.counter(
53 "kafka_sink_produce_latency_count",
54 "Produce latency samples",
55 ),
56 produce_latency_max_us: reg.gauge(
57 "kafka_sink_produce_latency_max_us",
58 "Max produce delivery latency (us)",
59 ),
60 }
61 }
62
63 pub fn record_write(&self, records: u64, bytes: u64) {
65 self.records_written.inc_by(records);
66 self.bytes_written.inc_by(bytes);
67 }
68
69 pub fn record_error(&self) {
71 self.errors_total.inc();
72 }
73
74 pub fn record_dlq(&self) {
76 self.dlq_records.inc();
77 }
78
79 pub fn record_serialization_error(&self) {
81 self.serialization_errors.inc();
82 }
83
84 #[allow(clippy::cast_possible_wrap)]
86 pub fn record_produce_latency(&self, latency_us: u64) {
87 self.produce_latency_sum_us.inc_by(latency_us);
88 self.produce_latency_count.inc();
89 if latency_us as i64 > self.produce_latency_max_us.get() {
90 self.produce_latency_max_us.set(latency_us as i64);
91 }
92 }
93}
94
95impl Default for KafkaSinkMetrics {
96 fn default() -> Self {
97 Self::new(None)
98 }
99}
100
101#[cfg(test)]
102mod tests {
103 use super::*;
104
105 #[test]
106 fn test_initial_zeros() {
107 let m = KafkaSinkMetrics::new(None);
108 assert_eq!(m.records_written.get(), 0);
109 assert_eq!(m.bytes_written.get(), 0);
110 assert_eq!(m.errors_total.get(), 0);
111 }
112
113 #[test]
114 fn test_record_write() {
115 let m = KafkaSinkMetrics::new(None);
116 m.record_write(100, 5000);
117 m.record_write(200, 10000);
118 assert_eq!(m.records_written.get(), 300);
119 assert_eq!(m.bytes_written.get(), 15000);
120 }
121
122 #[test]
123 fn test_produce_latency() {
124 let m = KafkaSinkMetrics::new(None);
125 m.record_produce_latency(100);
126 m.record_produce_latency(300);
127 m.record_produce_latency(50);
128
129 assert_eq!(m.produce_latency_count.get(), 3);
130 assert_eq!(m.produce_latency_sum_us.get(), 450);
131 assert_eq!(m.produce_latency_max_us.get(), 300);
132 }
133}