Merge per-shard partial aggregates into one result per group-by key. Groups by `group_key`; uses `ArrayAggPartial::merge` for each group.
(shard_resps: &[ArrayShardAggResp])
| 132 | /// |
| 133 | /// Groups by `group_key`; uses `ArrayAggPartial::merge` for each group. |
| 134 | pub fn reduce_agg_partials(shard_resps: &[ArrayShardAggResp]) -> Vec<ArrayAggPartial> { |
| 135 | use std::collections::BTreeMap; |
| 136 | let mut buckets: BTreeMap<i64, ArrayAggPartial> = BTreeMap::new(); |
| 137 | for resp in shard_resps { |
| 138 | for partial in &resp.partials { |
| 139 | buckets |
| 140 | .entry(partial.group_key) |
| 141 | .and_modify(|existing| existing.merge(partial)) |
| 142 | .or_insert_with(|| partial.clone()); |
| 143 | } |
| 144 | } |
| 145 | buckets.into_values().collect() |
| 146 | } |
| 147 | |
| 148 | #[cfg(test)] |
| 149 | mod tests { |