Skip to main content

laminar_connectors/kafka/
sink_metrics.rs

1//! Kafka sink connector metrics.
2
3use prometheus::{IntCounter, IntGauge, Registry};
4
5use crate::prom::reg_or_local;
6
7/// Prometheus-backed counters for Kafka sink connector statistics.
8#[derive(Debug, Clone)]
9pub struct KafkaSinkMetrics {
10    /// Records written to Kafka.
11    pub records_written: IntCounter,
12    /// Bytes written to Kafka (payload only).
13    pub bytes_written: IntCounter,
14    /// Errors encountered.
15    pub errors_total: IntCounter,
16    /// Records routed to dead letter queue.
17    pub dlq_records: IntCounter,
18    /// Serialization errors.
19    pub serialization_errors: IntCounter,
20    /// Sum of produce delivery latencies in microseconds.
21    pub produce_latency_sum_us: IntCounter,
22    /// Maximum produce delivery latency in microseconds.
23    pub produce_latency_max_us: IntGauge,
24    /// Number of produce delivery latency samples.
25    pub produce_latency_count: IntCounter,
26}
27
28impl KafkaSinkMetrics {
29    /// All counters start at zero. Registers on `registry` if provided.
30    #[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    /// Records a successful write of `records` records totaling `bytes`.
64    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    /// Records a production error.
70    pub fn record_error(&self) {
71        self.errors_total.inc();
72    }
73
74    /// Records a DLQ routing event.
75    pub fn record_dlq(&self) {
76        self.dlq_records.inc();
77    }
78
79    /// Records a serialization error.
80    pub fn record_serialization_error(&self) {
81        self.serialization_errors.inc();
82    }
83
84    /// Records a produce delivery latency sample in microseconds.
85    #[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}