diff --git a/aw-transform/benches/bench.rs b/aw-transform/benches/bench.rs index 0631e67c..712fc416 100644 --- a/aw-transform/benches/bench.rs +++ b/aw-transform/benches/bench.rs @@ -1,5 +1,5 @@ use chrono::Duration; -use criterion::{criterion_group, criterion_main, Criterion}; +use criterion::{criterion_group, criterion_main, BatchSize, BenchmarkId, Criterion}; use serde_json::json; use serde_json::Map; use serde_json::Value; @@ -50,5 +50,35 @@ fn bench_filter_period_intersect(c: &mut Criterion) { }); } -criterion_group!(benches, bench_filter_period_intersect); +fn bench_merge_events_by_keys(c: &mut Criterion) { + let mut group = c.benchmark_group("merge_events_by_keys"); + for distinct in [100, 50_000] { + let mut events = create_events(50_000); + for (i, event) in events.iter_mut().enumerate() { + event.data = json_map! { + "app": "browser", + "title": format!("Page {}", i % distinct), + "url": "https://example.com/a/long/path" + }; + } + group.bench_with_input( + BenchmarkId::new("distinct", distinct), + &events, + |b, events| { + b.iter_batched( + || (events.clone(), vec!["app".into(), "title".into()]), + |(events, keys)| merge_events_by_keys(events, keys), + BatchSize::LargeInput, + ); + }, + ); + } + group.finish(); +} + +criterion_group!( + benches, + bench_filter_period_intersect, + bench_merge_events_by_keys +); criterion_main!(benches); diff --git a/aw-transform/src/merge.rs b/aw-transform/src/merge.rs index f2034dda..b46c4e3f 100644 --- a/aw-transform/src/merge.rs +++ b/aw-transform/src/merge.rs @@ -39,13 +39,12 @@ use aw_models::Event; /// { duration: 1.0, data: { "a": 2, "b": 2 } } /// { duration: 1.0, data: { "a": 1, "b": 2 } } /// ``` -#[allow(clippy::map_entry)] 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(); - 'event: for event in events { + 'event: for mut event in events { let mut key_values = Vec::new(); for key in &keys { match event.data.get(key) { @@ -54,28 +53,17 @@ pub fn merge_events_by_keys(events: Vec, keys: Vec) -> Vec } } let summed_key = key_values.join("."); - if merged_events_map.contains_key(&summed_key) { - let merged_event = merged_events_map.get_mut(&summed_key).unwrap(); - merged_event.duration += event.duration; - } else { - let mut data = HashMap::new(); - for key in &keys { - data.insert(key.clone(), event.data.get(key).unwrap()); + match merged_events_map.entry(summed_key) { + std::collections::hash_map::Entry::Occupied(mut entry) => { + entry.get_mut().duration += event.duration; + } + std::collections::hash_map::Entry::Vacant(entry) => { + event.id = None; + entry.insert(event); } - let merged_event = Event { - id: None, - timestamp: event.timestamp, - duration: event.duration, - data: event.data.clone(), - }; - merged_events_map.insert(summed_key, merged_event); } } - let mut merged_events_list = Vec::new(); - for (_key, event) in merged_events_map.drain() { - merged_events_list.push(event); - } - merged_events_list + merged_events_map.into_values().collect() } #[cfg(test)] @@ -92,6 +80,30 @@ mod tests { use super::merge_events_by_keys; + #[test] + fn merge_preserves_first_payload_and_clears_id() { + let first = Event { + id: Some(42), + timestamp: DateTime::from_str("2000-01-01T00:00:01Z").unwrap(), + duration: Duration::seconds(2), + data: json_map! {"app": json!("browser"), "title": json!("page"), "extra": json!({"nested": [1, 2]})}, + }; + let mut second = first.clone(); + second.timestamp += Duration::seconds(10); + second.data.insert("extra".into(), json!("different")); + let mut missing = first.clone(); + missing.data.remove("title"); + let result = merge_events_by_keys( + vec![first.clone(), second, missing], + vec!["app".into(), "title".into()], + ); + assert_eq!(result.len(), 1); + assert_eq!(result[0].id, None); + assert_eq!(result[0].timestamp, first.timestamp); + assert_eq!(result[0].data, first.data); + assert_eq!(result[0].duration, Duration::seconds(4)); + } + #[test] fn test_merge_events_by_key() { let e1 = Event {