Skip to content
Open
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
53 changes: 45 additions & 8 deletions aw-transform/src/merge.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand All @@ -43,7 +48,10 @@ pub fn merge_events_by_keys(events: Vec<Event>, keys: Vec<String>) -> Vec<Event>
if keys.is_empty() {
return vec![];
}
let mut merged_events_map: HashMap<String, Event> = 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<String, usize> = HashMap::new();
let mut merged: Vec<Event> = Vec::new();
'event: for mut event in events {
let mut key_values = Vec::new();
for key in &keys {
Expand All @@ -53,17 +61,18 @@ pub fn merge_events_by_keys(events: Vec<Event>, keys: Vec<String>) -> Vec<Event>
}
}
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)]
Expand Down Expand Up @@ -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));
}
}
Loading