oxigdal_streaming/windowing/
watermark.rs1use crate::core::stream::StreamElement;
4use crate::error::Result;
5use chrono::{DateTime, Duration, Utc};
6use serde::{Deserialize, Serialize};
7use std::collections::BTreeMap;
8use std::sync::Arc;
9use tokio::sync::RwLock;
10
11#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
13pub struct Watermark {
14 pub timestamp: DateTime<Utc>,
16}
17
18impl Watermark {
19 pub fn new(timestamp: DateTime<Utc>) -> Self {
21 Self { timestamp }
22 }
23
24 pub fn min() -> Self {
26 Self {
27 timestamp: DateTime::from_timestamp(0, 0).unwrap_or_else(Utc::now),
28 }
29 }
30
31 pub fn max() -> Self {
33 Self {
34 timestamp: DateTime::from_timestamp(i64::MAX / 1000, 0).unwrap_or_else(Utc::now),
35 }
36 }
37}
38
39#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
41pub enum WatermarkStrategy {
42 Ascending,
44
45 BoundedOutOfOrderness,
47
48 Periodic,
50
51 Punctuated,
53}
54
55#[derive(Debug, Clone, Serialize, Deserialize)]
57pub struct WatermarkConfig {
58 pub strategy: WatermarkStrategy,
60
61 pub max_out_of_orderness: Duration,
63
64 pub interval: Duration,
66
67 pub idle_timeout: Option<Duration>,
69}
70
71impl Default for WatermarkConfig {
72 fn default() -> Self {
73 Self {
74 strategy: WatermarkStrategy::BoundedOutOfOrderness,
75 max_out_of_orderness: Duration::seconds(5),
76 interval: Duration::seconds(1),
77 idle_timeout: Some(Duration::seconds(10)),
78 }
79 }
80}
81
82pub trait WatermarkGenerator: Send + Sync {
84 fn on_event(&mut self, element: &StreamElement) -> Option<Watermark>;
86
87 fn on_periodic_emit(&mut self) -> Option<Watermark>;
89
90 fn current_watermark(&self) -> Watermark;
92}
93
94pub struct PeriodicWatermarkGenerator {
96 config: WatermarkConfig,
97 max_timestamp: Option<DateTime<Utc>>,
98 current_watermark: Watermark,
99 last_emit: Option<DateTime<Utc>>,
100}
101
102impl PeriodicWatermarkGenerator {
103 pub fn new(config: WatermarkConfig) -> Self {
105 Self {
106 config,
107 max_timestamp: None,
108 current_watermark: Watermark::min(),
109 last_emit: None,
110 }
111 }
112}
113
114impl WatermarkGenerator for PeriodicWatermarkGenerator {
115 fn on_event(&mut self, element: &StreamElement) -> Option<Watermark> {
116 if let Some(max_ts) = self.max_timestamp {
117 if element.event_time > max_ts {
118 self.max_timestamp = Some(element.event_time);
119 }
120 } else {
121 self.max_timestamp = Some(element.event_time);
122 }
123
124 None
125 }
126
127 fn on_periodic_emit(&mut self) -> Option<Watermark> {
128 let now = Utc::now();
129 let should_emit = if let Some(last) = self.last_emit {
130 now - last >= self.config.interval
131 } else {
132 true
133 };
134
135 if should_emit && let Some(max_ts) = self.max_timestamp {
136 let new_watermark = match self.config.strategy {
137 WatermarkStrategy::Ascending => Watermark::new(max_ts),
138 WatermarkStrategy::BoundedOutOfOrderness => {
139 Watermark::new(max_ts - self.config.max_out_of_orderness)
140 }
141 _ => self.current_watermark,
142 };
143
144 if new_watermark > self.current_watermark {
145 self.current_watermark = new_watermark;
146 self.last_emit = Some(now);
147 return Some(new_watermark);
148 }
149 }
150
151 None
152 }
153
154 fn current_watermark(&self) -> Watermark {
155 self.current_watermark
156 }
157}
158
159pub struct PunctuatedWatermarkGenerator {
161 config: WatermarkConfig,
162 current_watermark: Watermark,
163 max_timestamp: Option<DateTime<Utc>>,
164}
165
166impl PunctuatedWatermarkGenerator {
167 pub fn new(config: WatermarkConfig) -> Self {
169 Self {
170 config,
171 current_watermark: Watermark::min(),
172 max_timestamp: None,
173 }
174 }
175
176 fn should_emit_watermark(&self, element: &StreamElement) -> bool {
178 if let Some(marker) = element.metadata.attributes.get("watermark_marker") {
179 marker == "true"
180 } else {
181 false
182 }
183 }
184}
185
186impl WatermarkGenerator for PunctuatedWatermarkGenerator {
187 fn on_event(&mut self, element: &StreamElement) -> Option<Watermark> {
188 if let Some(max_ts) = self.max_timestamp {
189 if element.event_time > max_ts {
190 self.max_timestamp = Some(element.event_time);
191 }
192 } else {
193 self.max_timestamp = Some(element.event_time);
194 }
195
196 if self.should_emit_watermark(element)
197 && let Some(max_ts) = self.max_timestamp
198 {
199 let new_watermark = Watermark::new(max_ts - self.config.max_out_of_orderness);
200
201 if new_watermark > self.current_watermark {
202 self.current_watermark = new_watermark;
203 return Some(new_watermark);
204 }
205 }
206
207 None
208 }
209
210 fn on_periodic_emit(&mut self) -> Option<Watermark> {
211 None
212 }
213
214 fn current_watermark(&self) -> Watermark {
215 self.current_watermark
216 }
217}
218
219pub struct MultiSourceWatermarkManager {
221 source_watermarks: Arc<RwLock<BTreeMap<String, Watermark>>>,
222 global_watermark: Arc<RwLock<Watermark>>,
223}
224
225impl MultiSourceWatermarkManager {
226 pub fn new() -> Self {
228 Self {
229 source_watermarks: Arc::new(RwLock::new(BTreeMap::new())),
230 global_watermark: Arc::new(RwLock::new(Watermark::min())),
231 }
232 }
233
234 pub async fn update_source_watermark(
236 &self,
237 source_id: String,
238 watermark: Watermark,
239 ) -> Result<()> {
240 let mut watermarks = self.source_watermarks.write().await;
241 watermarks.insert(source_id, watermark);
242
243 let min_watermark = watermarks
244 .values()
245 .min()
246 .copied()
247 .unwrap_or(Watermark::min());
248
249 let mut global = self.global_watermark.write().await;
250 if min_watermark > *global {
251 *global = min_watermark;
252 }
253
254 Ok(())
255 }
256
257 pub async fn global_watermark(&self) -> Watermark {
259 *self.global_watermark.read().await
260 }
261
262 pub async fn source_watermark(&self, source_id: &str) -> Option<Watermark> {
264 self.source_watermarks.read().await.get(source_id).copied()
265 }
266
267 pub async fn remove_source(&self, source_id: &str) -> Result<()> {
269 let mut watermarks = self.source_watermarks.write().await;
270 watermarks.remove(source_id);
271
272 let min_watermark = watermarks
273 .values()
274 .min()
275 .copied()
276 .unwrap_or(Watermark::max());
277
278 let mut global = self.global_watermark.write().await;
279 *global = min_watermark;
280
281 Ok(())
282 }
283}
284
285impl Default for MultiSourceWatermarkManager {
286 fn default() -> Self {
287 Self::new()
288 }
289}
290
291#[cfg(test)]
292mod tests {
293 use super::*;
294
295 #[test]
296 fn test_watermark_creation() {
297 let now = Utc::now();
298 let wm = Watermark::new(now);
299 assert_eq!(wm.timestamp, now);
300 }
301
302 #[test]
303 fn test_watermark_ordering() {
304 let wm1 = Watermark::new(Utc::now());
305 let wm2 = Watermark::new(Utc::now() + Duration::seconds(10));
306
307 assert!(wm1 < wm2);
308 assert!(wm2 > wm1);
309 }
310
311 #[tokio::test]
312 async fn test_periodic_watermark_generator() {
313 let config = WatermarkConfig::default();
314 let mut generator = PeriodicWatermarkGenerator::new(config);
315
316 let elem = StreamElement::new(vec![1, 2, 3], Utc::now());
317 generator.on_event(&elem);
318
319 let wm = generator.on_periodic_emit();
320 assert!(wm.is_some());
321 }
322
323 #[tokio::test]
324 async fn test_multi_source_watermark_manager() {
325 let manager = MultiSourceWatermarkManager::new();
326
327 let wm1 = Watermark::new(Utc::now());
328 let wm2 = Watermark::new(Utc::now() + Duration::seconds(10));
329
330 manager
331 .update_source_watermark("source1".to_string(), wm1)
332 .await
333 .expect("Test watermark update for source1 should succeed");
334 manager
335 .update_source_watermark("source2".to_string(), wm2)
336 .await
337 .expect("Test watermark update for source2 should succeed");
338
339 let global = manager.global_watermark().await;
340 assert_eq!(global, wm1);
341 }
342
343 #[tokio::test]
344 async fn test_remove_source_watermark() {
345 let manager = MultiSourceWatermarkManager::new();
346
347 let wm1 = Watermark::new(Utc::now());
348 let wm2 = Watermark::new(Utc::now() + Duration::seconds(10));
349
350 manager
351 .update_source_watermark("source1".to_string(), wm1)
352 .await
353 .expect("Test watermark update for source1 should succeed");
354 manager
355 .update_source_watermark("source2".to_string(), wm2)
356 .await
357 .expect("Test watermark update for source2 should succeed");
358
359 manager
360 .remove_source("source1")
361 .await
362 .expect("Test source removal should succeed");
363
364 let global = manager.global_watermark().await;
365 assert_eq!(global, wm2);
366 }
367}