laminar_connectors/kafka/sink_metrics/
mod.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;