diff --git a/aw-transform/src/merge.rs b/aw-transform/src/merge.rs index b46c4e3f..f04fbdb5 100644 --- a/aw-transform/src/merge.rs +++ b/aw-transform/src/merge.rs @@ -6,7 +6,13 @@ use aw_models::Event; /// /// Doesn't care about if events are neighbouring or not, this transform merges /// all events with the same key. -/// The timestamp will be the timestamp of the first event with a specific key value +/// Each merged event keeps the timestamp and the whole data of the first event +/// with those key values, and the durations are summed. Events missing any of +/// the keys are dropped, and an empty key list returns no events. The merged +/// events come out in the order their first event appeared in the input, so +/// the result is deterministic. That is not necessarily timestamp order: sort +/// by timestamp before passing the result to union_no_overlap, which expects +/// chronological input. /// /// # Example 1 /// A simple example only using one key @@ -22,7 +28,6 @@ use aw_models::Event; /// output: /// { duration: 3.0, data: { "a": 1 } } /// { duration: 1.0, data: { "a": 2 } } -/// { duration: 1.0, data: { "b": 1 } } /// ``` /// /// # Example 2 @@ -43,7 +48,10 @@ pub fn merge_events_by_keys(events: Vec, keys: Vec) -> Vec if keys.is_empty() { return vec![]; } - let mut merged_events_map: HashMap = HashMap::new(); + // Index into `merged` by key, so the output keeps first-seen order + // instead of HashMap iteration order, which differs between runs. + let mut index: HashMap = HashMap::new(); + let mut merged: Vec = Vec::new(); 'event: for mut event in events { let mut key_values = Vec::new(); for key in &keys { @@ -53,17 +61,18 @@ pub fn merge_events_by_keys(events: Vec, keys: Vec) -> Vec } } let summed_key = key_values.join("."); - match merged_events_map.entry(summed_key) { - std::collections::hash_map::Entry::Occupied(mut entry) => { - entry.get_mut().duration += event.duration; + match index.entry(summed_key) { + std::collections::hash_map::Entry::Occupied(entry) => { + merged[*entry.get()].duration += event.duration; } std::collections::hash_map::Entry::Vacant(entry) => { event.id = None; - entry.insert(event); + entry.insert(merged.len()); + merged.push(event); } } } - merged_events_map.into_values().collect() + merged } #[cfg(test)] @@ -150,4 +159,32 @@ mod tests { ]; assert_eq!(&res2, &expected); } + + #[test] + fn merge_keeps_first_seen_order() { + // HashMap iteration order would vary between runs; the output must + // follow the order in which each group first appears. + let mk = |secs: &str, app: &str| Event { + id: None, + timestamp: DateTime::from_str(secs).unwrap(), + duration: Duration::seconds(1), + data: json_map! {"app": json!(app)}, + }; + let apps = ["m", "c", "x", "a", "q", "b", "z", "d", "k", "e"]; + let mut events = Vec::new(); + for (i, app) in apps.iter().enumerate() { + events.push(mk(&format!("2000-01-01T00:00:{:02}Z", i), app)); + } + // Repeats merge into the first occurrence without reordering + events.push(mk("2000-01-01T00:00:30Z", "a")); + events.push(mk("2000-01-01T00:00:31Z", "m")); + let res = merge_events_by_keys(events, vec!["app".to_string()]); + let order: Vec<&str> = res + .iter() + .map(|e| e.data.get("app").unwrap().as_str().unwrap()) + .collect(); + assert_eq!(order, apps); + assert_eq!(res[0].duration, Duration::seconds(2)); + assert_eq!(res[3].duration, Duration::seconds(2)); + } }