Skip to main content

laminar_connectors/postgres/
sink_metrics.rs

1//! `PostgreSQL` sink connector metrics.
2//!
3//! [`PostgresSinkMetrics`] provides prometheus-backed counters for
4//! tracking write statistics.
5
6use prometheus::{IntCounter, Registry};
7
8use crate::prom::reg_or_local;
9
10/// Prometheus-backed counters for `PostgreSQL` sink connector statistics.
11#[derive(Debug, Clone)]
12pub struct PostgresSinkMetrics {
13    /// Total records written to `PostgreSQL`.
14    pub records_written: IntCounter,
15
16    /// Total bytes written (estimated from `RecordBatch` sizes).
17    pub bytes_written: IntCounter,
18
19    /// Total errors encountered.
20    pub errors_total: IntCounter,
21
22    /// Total batches flushed.
23    pub batches_flushed: IntCounter,
24
25    /// Total COPY BINARY operations (append mode).
26    pub copy_operations: IntCounter,
27
28    /// Total upsert operations (upsert mode).
29    pub upsert_operations: IntCounter,
30
31    /// Total changelog deletes applied (Z-set weight -1).
32    pub changelog_deletes: IntCounter,
33}
34
35impl PostgresSinkMetrics {
36    /// Creates a new metrics instance with all counters at zero.
37    #[must_use]
38    #[allow(clippy::missing_panics_doc)]
39    pub fn new(registry: Option<&Registry>) -> Self {
40        let mut local = None;
41        let reg = reg_or_local(registry, &mut local);
42
43        Self {
44            records_written: reg.counter(
45                "postgres_sink_records_written_total",
46                "Total records written to PostgreSQL",
47            ),
48            bytes_written: reg.counter(
49                "postgres_sink_bytes_written_total",
50                "Total bytes written to PostgreSQL",
51            ),
52            errors_total: reg.counter("postgres_sink_errors_total", "Total PostgreSQL sink errors"),
53            batches_flushed: reg.counter(
54                "postgres_sink_batches_flushed_total",
55                "Total batches flushed",
56            ),
57            copy_operations: reg.counter(
58                "postgres_sink_copy_operations_total",
59                "Total COPY BINARY operations",
60            ),
61            upsert_operations: reg.counter(
62                "postgres_sink_upsert_operations_total",
63                "Total upsert operations",
64            ),
65            changelog_deletes: reg.counter(
66                "postgres_sink_changelog_deletes_total",
67                "Total changelog deletes applied",
68            ),
69        }
70    }
71
72    /// Records a successful write of `records` records totaling `bytes`.
73    pub fn record_write(&self, records: u64, bytes: u64) {
74        self.records_written.inc_by(records);
75        self.bytes_written.inc_by(bytes);
76    }
77
78    /// Records a successful batch flush.
79    pub fn record_flush(&self) {
80        self.batches_flushed.inc();
81    }
82
83    /// Records a COPY BINARY operation.
84    pub fn record_copy(&self) {
85        self.copy_operations.inc();
86    }
87
88    /// Records an upsert operation.
89    pub fn record_upsert(&self) {
90        self.upsert_operations.inc();
91    }
92
93    /// Records a write or connection error.
94    pub fn record_error(&self) {
95        self.errors_total.inc();
96    }
97
98    /// Records changelog DELETE operations.
99    pub fn record_deletes(&self, count: u64) {
100        self.changelog_deletes.inc_by(count);
101    }
102}
103
104impl Default for PostgresSinkMetrics {
105    fn default() -> Self {
106        Self::new(None)
107    }
108}
109
110#[cfg(test)]
111mod tests {
112    use super::*;
113
114    #[test]
115    fn test_initial_zeros() {
116        let m = PostgresSinkMetrics::new(None);
117        assert_eq!(m.records_written.get(), 0);
118        assert_eq!(m.bytes_written.get(), 0);
119        assert_eq!(m.errors_total.get(), 0);
120    }
121
122    #[test]
123    fn test_record_write() {
124        let m = PostgresSinkMetrics::new(None);
125        m.record_write(100, 5000);
126        m.record_write(200, 10_000);
127
128        assert_eq!(m.records_written.get(), 300);
129        assert_eq!(m.bytes_written.get(), 15_000);
130    }
131
132    #[test]
133    fn test_flush_and_copy_metrics() {
134        let m = PostgresSinkMetrics::new(None);
135        m.record_flush();
136        m.record_flush();
137        m.record_copy();
138
139        assert_eq!(m.batches_flushed.get(), 2);
140        assert_eq!(m.copy_operations.get(), 1);
141    }
142
143    #[test]
144    fn test_changelog_deletes() {
145        let m = PostgresSinkMetrics::new(None);
146        m.record_deletes(50);
147        m.record_deletes(30);
148
149        assert_eq!(m.changelog_deletes.get(), 80);
150    }
151
152    #[test]
153    fn test_error_counting() {
154        let m = PostgresSinkMetrics::new(None);
155        m.record_error();
156        m.record_error();
157        m.record_error();
158
159        assert_eq!(m.errors_total.get(), 3);
160    }
161
162    #[test]
163    fn test_upsert_metric() {
164        let m = PostgresSinkMetrics::new(None);
165        m.record_upsert();
166
167        assert_eq!(m.upsert_operations.get(), 1);
168    }
169}