laminar_connectors/kafka/
metrics.rs1use prometheus::{IntCounter, Registry};
4
5use crate::prom::reg_or_local;
6
7#[derive(Debug, Clone)]
9pub struct KafkaSourceMetrics {
10 pub records_polled: IntCounter,
12 pub bytes_polled: IntCounter,
14 pub errors: IntCounter,
16 pub batches_polled: IntCounter,
18 pub commits: IntCounter,
20 pub commit_failures: IntCounter,
22 pub rebalances: IntCounter,
24 pub sr_discovery_successes: IntCounter,
26 pub sr_discovery_failures: IntCounter,
28 pub sr_discovery_timeouts: IntCounter,
30}
31
32impl KafkaSourceMetrics {
33 #[must_use]
36 #[allow(clippy::missing_panics_doc, clippy::too_many_lines)]
37 pub fn new(registry: Option<&Registry>) -> Self {
38 let mut local = None;
39 let handle = reg_or_local(registry, &mut local);
40
41 Self {
42 records_polled: handle.counter(
43 "kafka_source_records_polled_total",
44 "Total records polled from Kafka",
45 ),
46 bytes_polled: handle.counter(
47 "kafka_source_bytes_polled_total",
48 "Total bytes polled from Kafka",
49 ),
50 errors: handle.counter("kafka_source_errors_total", "Total Kafka consumer errors"),
51 batches_polled: handle.counter(
52 "kafka_source_batches_polled_total",
53 "Total batches polled from Kafka",
54 ),
55 commits: handle.counter(
56 "kafka_source_commits_total",
57 "Total offset commits to Kafka",
58 ),
59 commit_failures: handle.counter(
60 "kafka_source_commit_failures_total",
61 "Kafka broker progress commit failures",
62 ),
63 rebalances: handle.counter(
64 "kafka_source_rebalances_total",
65 "Total consumer group rebalances",
66 ),
67 sr_discovery_successes: handle.counter(
68 "kafka_source_sr_discovery_successes_total",
69 "Schema Registry discovery successes",
70 ),
71 sr_discovery_failures: handle.counter(
72 "kafka_source_sr_discovery_failures_total",
73 "Schema Registry discovery failures",
74 ),
75 sr_discovery_timeouts: handle.counter(
76 "kafka_source_sr_discovery_timeouts_total",
77 "Schema Registry discovery timeouts",
78 ),
79 }
80 }
81
82 pub fn record_poll(&self, records: u64, bytes: u64) {
84 self.records_polled.inc_by(records);
85 self.bytes_polled.inc_by(bytes);
86 self.batches_polled.inc();
87 }
88
89 pub fn record_error(&self) {
91 self.errors.inc();
92 }
93
94 pub fn record_rebalance(&self) {
96 self.rebalances.inc();
97 }
98
99 pub fn record_sr_discovery_success(&self) {
101 self.sr_discovery_successes.inc();
102 }
103
104 pub fn record_sr_discovery_failure(&self) {
106 self.sr_discovery_failures.inc();
107 }
108
109 pub fn record_sr_discovery_timeout(&self) {
111 self.sr_discovery_timeouts.inc();
112 }
113}
114
115impl Default for KafkaSourceMetrics {
116 fn default() -> Self {
117 Self::new(None)
118 }
119}
120
121#[cfg(test)]
122mod tests {
123 use super::*;
124
125 #[test]
126 fn test_initial_zeros() {
127 let m = KafkaSourceMetrics::new(None);
128 assert_eq!(m.records_polled.get(), 0);
129 assert_eq!(m.bytes_polled.get(), 0);
130 assert_eq!(m.errors.get(), 0);
131 }
132
133 #[test]
134 fn test_record_poll() {
135 let m = KafkaSourceMetrics::new(None);
136 m.record_poll(100, 5000);
137 m.record_poll(200, 10000);
138
139 assert_eq!(m.records_polled.get(), 300);
140 assert_eq!(m.bytes_polled.get(), 15000);
141 }
142
143 #[test]
144 fn test_record_error() {
145 let m = KafkaSourceMetrics::new(None);
146 m.record_error();
147 m.record_error();
148
149 assert_eq!(m.errors.get(), 2);
150 }
151
152 #[test]
153 fn test_record_commit_failure() {
154 let m = KafkaSourceMetrics::new(None);
155 m.commit_failures.inc();
156 assert_eq!(m.commit_failures.get(), 1);
157 }
158
159 #[test]
160 fn test_record_rebalance() {
161 let m = KafkaSourceMetrics::new(None);
162 m.record_rebalance();
163 m.record_rebalance();
164
165 assert_eq!(m.rebalances.get(), 2);
166 }
167
168 #[test]
169 fn test_sr_discovery_counters() {
170 let m = KafkaSourceMetrics::new(None);
171 m.record_sr_discovery_success();
172 m.record_sr_discovery_success();
173 m.record_sr_discovery_failure();
174 m.record_sr_discovery_timeout();
175
176 assert_eq!(m.sr_discovery_successes.get(), 2);
177 assert_eq!(m.sr_discovery_failures.get(), 1);
178 assert_eq!(m.sr_discovery_timeouts.get(), 1);
179 }
180
181 #[test]
182 fn test_registered_on_prometheus_registry() {
183 let reg = Registry::new();
184 let m = KafkaSourceMetrics::new(Some(®));
185 m.record_poll(10, 500);
186 m.record_error();
187
188 let families = reg.gather();
190 let names: Vec<&str> = families
191 .iter()
192 .map(prometheus::proto::MetricFamily::name)
193 .collect();
194 assert!(names.contains(&"kafka_source_records_polled_total"));
195 assert!(names.contains(&"kafka_source_errors_total"));
196 }
197}