Skip to main content

fluxus_api/operators/
window_skipper.rs

1use async_trait::async_trait;
2use fluxus_transformers::Operator;
3use fluxus_utils::{
4    models::{Record, StreamResult},
5    window::WindowConfig,
6};
7use std::{collections::HashMap, marker::PhantomData};
8
9pub struct WindowSkipper<T> {
10    window_config: WindowConfig,
11    n: usize,
12    buffer: HashMap<u64, Vec<T>>,
13    _phantom: PhantomData<T>,
14}
15
16impl<T> WindowSkipper<T>
17where
18    T: Clone,
19{
20    pub fn new(window_config: WindowConfig, n: usize) -> Self {
21        Self {
22            window_config,
23            n,
24            buffer: HashMap::new(),
25            _phantom: PhantomData,
26        }
27    }
28
29    fn get_window_keys(&self, timestamp: i64) -> Vec<u64> {
30        self.window_config.window_type.get_window_keys(timestamp)
31    }
32}
33
34#[async_trait]
35impl<T> Operator<T, Vec<T>> for WindowSkipper<T>
36where
37    T: Clone + Send + Sync + 'static,
38{
39    async fn process(&mut self, record: Record<T>) -> StreamResult<Vec<Record<Vec<T>>>> {
40        let mut results = Vec::new();
41
42        for window_key in self.get_window_keys(record.timestamp) {
43            let records = self.buffer.entry(window_key).or_default();
44            records.push(record.data.clone());
45            let new_records = records.iter().skip(self.n).cloned().collect::<Vec<_>>();
46            results.push(Record {
47                data: new_records,
48                timestamp: record.timestamp,
49            });
50        }
51
52        Ok(results)
53    }
54}