Skip to main content

oxigdal_streaming/windowing/
watermark.rs

1//! Watermark generation for event-time processing.
2
3use 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/// A watermark representing event-time progress.
12#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
13pub struct Watermark {
14    /// Watermark timestamp
15    pub timestamp: DateTime<Utc>,
16}
17
18impl Watermark {
19    /// Create a new watermark.
20    pub fn new(timestamp: DateTime<Utc>) -> Self {
21        Self { timestamp }
22    }
23
24    /// Get the minimum watermark (beginning of time).
25    pub fn min() -> Self {
26        Self {
27            timestamp: DateTime::from_timestamp(0, 0).unwrap_or_else(Utc::now),
28        }
29    }
30
31    /// Get the maximum watermark (end of time).
32    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/// Strategy for generating watermarks.
40#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
41pub enum WatermarkStrategy {
42    /// Ascending timestamps (watermark = max observed timestamp)
43    Ascending,
44
45    /// Bounded out-of-orderness (watermark = max timestamp - max delay)
46    BoundedOutOfOrderness,
47
48    /// Periodic watermarks
49    Periodic,
50
51    /// Punctuated watermarks (based on special markers)
52    Punctuated,
53}
54
55/// Configuration for watermark generation.
56#[derive(Debug, Clone, Serialize, Deserialize)]
57pub struct WatermarkConfig {
58    /// Watermark strategy
59    pub strategy: WatermarkStrategy,
60
61    /// Maximum out-of-orderness
62    pub max_out_of_orderness: Duration,
63
64    /// Watermark interval (for periodic strategy)
65    pub interval: Duration,
66
67    /// Idle timeout (emit watermark if no data for this duration)
68    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
82/// Generates watermarks for a stream.
83pub trait WatermarkGenerator: Send + Sync {
84    /// Process an element and potentially generate a watermark.
85    fn on_event(&mut self, element: &StreamElement) -> Option<Watermark>;
86
87    /// Generate a periodic watermark.
88    fn on_periodic_emit(&mut self) -> Option<Watermark>;
89
90    /// Get the current watermark.
91    fn current_watermark(&self) -> Watermark;
92}
93
94/// Periodic watermark generator.
95pub 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    /// Create a new periodic watermark generator.
104    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
159/// Punctuated watermark generator.
160pub struct PunctuatedWatermarkGenerator {
161    config: WatermarkConfig,
162    current_watermark: Watermark,
163    max_timestamp: Option<DateTime<Utc>>,
164}
165
166impl PunctuatedWatermarkGenerator {
167    /// Create a new punctuated watermark generator.
168    pub fn new(config: WatermarkConfig) -> Self {
169        Self {
170            config,
171            current_watermark: Watermark::min(),
172            max_timestamp: None,
173        }
174    }
175
176    /// Check if an element should trigger a watermark.
177    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
219/// Multi-source watermark manager.
220pub struct MultiSourceWatermarkManager {
221    source_watermarks: Arc<RwLock<BTreeMap<String, Watermark>>>,
222    global_watermark: Arc<RwLock<Watermark>>,
223}
224
225impl MultiSourceWatermarkManager {
226    /// Create a new multi-source watermark manager.
227    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    /// Update watermark for a source.
235    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    /// Get the global watermark (minimum of all source watermarks).
258    pub async fn global_watermark(&self) -> Watermark {
259        *self.global_watermark.read().await
260    }
261
262    /// Get watermark for a specific source.
263    pub async fn source_watermark(&self, source_id: &str) -> Option<Watermark> {
264        self.source_watermarks.read().await.get(source_id).copied()
265    }
266
267    /// Remove a source.
268    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}