Skip to main content

laminar_connectors/kafka/
metrics.rs

1//! Prometheus-backed Kafka source metrics.
2
3use prometheus::{IntCounter, Registry};
4
5use crate::prom::reg_or_local;
6
7/// Prometheus-backed counters/gauges for Kafka source connector statistics.
8#[derive(Debug, Clone)]
9pub struct KafkaSourceMetrics {
10    /// Total records polled from Kafka.
11    pub records_polled: IntCounter,
12    /// Total bytes polled from Kafka.
13    pub bytes_polled: IntCounter,
14    /// Total deserialization or consumer errors.
15    pub errors: IntCounter,
16    /// Total batches returned from `poll_batch()`.
17    pub batches_polled: IntCounter,
18    /// Total offset commits to Kafka.
19    pub commits: IntCounter,
20    /// Broker progress commits that failed locally or remotely.
21    pub commit_failures: IntCounter,
22    /// Total consumer group rebalances.
23    pub rebalances: IntCounter,
24    /// Count of successful Schema Registry discoveries at DDL time.
25    pub sr_discovery_successes: IntCounter,
26    /// Count of Schema Registry discovery failures (HTTP error, parse error).
27    pub sr_discovery_failures: IntCounter,
28    /// Count of Schema Registry discovery timeouts.
29    pub sr_discovery_timeouts: IntCounter,
30}
31
32impl KafkaSourceMetrics {
33    /// If `registry` is `Some`, counters are registered there (visible
34    /// in the Prometheus scrape); otherwise a throwaway registry is used.
35    #[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    /// Records a successful poll of `records` records totaling `bytes`.
83    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    /// Records a consumer or deserialization error.
90    pub fn record_error(&self) {
91        self.errors.inc();
92    }
93
94    /// Records a consumer group rebalance event.
95    pub fn record_rebalance(&self) {
96        self.rebalances.inc();
97    }
98
99    /// Records a successful Schema Registry discovery at DDL time.
100    pub fn record_sr_discovery_success(&self) {
101        self.sr_discovery_successes.inc();
102    }
103
104    /// Records a Schema Registry discovery failure.
105    pub fn record_sr_discovery_failure(&self) {
106        self.sr_discovery_failures.inc();
107    }
108
109    /// Records a Schema Registry discovery timeout.
110    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(&reg));
185        m.record_poll(10, 500);
186        m.record_error();
187
188        // Verify the metrics are registered on the registry.
189        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}