fluxus_api/operators/
window_skipper.rs1use 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}