reifydb_core/value/column/view/
group_by.rs1use std::{
5 iter::{Enumerate, FilterMap},
6 mem,
7 vec::IntoIter as VecIntoIter,
8};
9
10use indexmap::IndexMap;
11use reifydb_codec::key::{encoded::EncodedKey, serializer::KeySerializer};
12use reifydb_value::{Result, error::Error, value::Value};
13
14use crate::{
15 error::CoreError,
16 metrics::heap::HeapSize,
17 value::column::{ColumnBuffer, columns::Columns},
18};
19
20pub type GroupKey = Vec<Value>;
21
22#[repr(transparent)]
23#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
24pub struct GroupId(pub u32);
25
26impl GroupId {
27 pub fn index(self) -> usize {
28 self.0 as usize
29 }
30}
31
32pub type GroupRows = Vec<(GroupId, Vec<usize>)>;
33
34#[derive(Debug, Clone)]
35pub struct GroupSlots<T> {
36 slots: Vec<Option<T>>,
37 occupied: usize,
38}
39
40impl<T> Default for GroupSlots<T> {
41 fn default() -> Self {
42 Self::new()
43 }
44}
45
46impl<T> GroupSlots<T> {
47 pub fn new() -> Self {
48 Self {
49 slots: Vec::new(),
50 occupied: 0,
51 }
52 }
53
54 pub fn len(&self) -> usize {
55 self.occupied
56 }
57
58 pub fn is_empty(&self) -> bool {
59 self.occupied == 0
60 }
61
62 pub fn get(&self, group: GroupId) -> Option<&T> {
63 self.slots.get(group.index()).and_then(Option::as_ref)
64 }
65
66 pub fn insert(&mut self, group: GroupId, value: T) {
67 self.reserve_for(group);
68 if self.slots[group.index()].replace(value).is_none() {
69 self.occupied += 1;
70 }
71 }
72
73 pub fn get_or_insert_with(&mut self, group: GroupId, default: impl FnOnce() -> T) -> &mut T {
74 self.reserve_for(group);
75 if self.slots[group.index()].is_none() {
76 self.occupied += 1;
77 }
78 self.slots[group.index()].get_or_insert_with(default)
79 }
80
81 pub fn or_insert(&mut self, group: GroupId, default: T) -> &mut T {
82 self.get_or_insert_with(group, || default)
83 }
84
85 pub fn remove(&mut self, group: GroupId) -> Option<T> {
86 let removed = self.slots.get_mut(group.index()).and_then(Option::take);
87 if removed.is_some() {
88 self.occupied -= 1;
89 }
90 removed
91 }
92
93 pub fn iter(&self) -> impl Iterator<Item = (GroupId, &T)> {
94 self.slots
95 .iter()
96 .enumerate()
97 .filter_map(|(index, slot)| slot.as_ref().map(|value| (GroupId(index as u32), value)))
98 }
99
100 pub fn drain(&mut self) -> impl Iterator<Item = (GroupId, T)> + '_ {
101 self.occupied = 0;
102 self.slots
103 .drain(..)
104 .enumerate()
105 .filter_map(|(index, slot)| slot.map(|value| (GroupId(index as u32), value)))
106 }
107
108 fn reserve_for(&mut self, group: GroupId) {
109 if self.slots.len() <= group.index() {
110 self.slots.resize_with(group.index() + 1, || None);
111 }
112 }
113}
114
115impl<T: HeapSize> HeapSize for GroupSlots<T> {
116 fn heap_size(&self) -> usize {
117 self.slots.heap_size()
118 }
119}
120
121fn occupied_slot<T>((index, slot): (usize, Option<T>)) -> Option<(GroupId, T)> {
122 slot.map(|value| (GroupId(index as u32), value))
123}
124
125impl<T> IntoIterator for GroupSlots<T> {
126 type Item = (GroupId, T);
127 type IntoIter = FilterMap<Enumerate<VecIntoIter<Option<T>>>, fn((usize, Option<T>)) -> Option<(GroupId, T)>>;
128
129 fn into_iter(self) -> Self::IntoIter {
130 self.slots
131 .into_iter()
132 .enumerate()
133 .filter_map(occupied_slot as fn((usize, Option<T>)) -> Option<(GroupId, T)>)
134 }
135}
136
137#[derive(Debug, Default, Clone)]
138pub struct GroupKeyDict {
139 entries: IndexMap<EncodedKey, GroupKey>,
140}
141
142impl GroupKeyDict {
143 pub fn new() -> Self {
144 Self {
145 entries: IndexMap::new(),
146 }
147 }
148
149 pub fn len(&self) -> usize {
150 self.entries.len()
151 }
152
153 pub fn is_empty(&self) -> bool {
154 self.entries.is_empty()
155 }
156
157 pub fn values(&self, group: GroupId) -> Option<&GroupKey> {
158 self.entries.get_index(group.index()).map(|(_, values)| values)
159 }
160
161 pub fn iter(&self) -> impl Iterator<Item = (GroupId, &GroupKey)> {
162 self.entries.values().enumerate().map(|(index, values)| (GroupId(index as u32), values))
163 }
164
165 fn intern(&mut self, encoded: &EncodedKey, materialize: impl FnOnce() -> GroupKey) -> GroupId {
166 if let Some(index) = self.entries.get_index_of(encoded) {
167 return GroupId(index as u32);
168 }
169 let (index, _) = self.entries.insert_full(encoded.clone(), materialize());
170 GroupId(index as u32)
171 }
172}
173
174impl HeapSize for GroupKeyDict {
175 fn heap_size(&self) -> usize {
176 self.entries.capacity()
177 * (mem::size_of::<EncodedKey>() + mem::size_of::<GroupKey>() + mem::size_of::<usize>())
178 + self.entries.iter().map(|(key, values)| key.heap_size() + values.heap_size()).sum::<usize>()
179 }
180}
181
182impl Columns {
183 pub fn group_by_ids(&self, keys: &[&str], dict: &mut GroupKeyDict) -> Result<GroupRows> {
184 let row_count = self.columns.first().map_or(0, |c| c.len());
185 let key_columns = self.key_columns(keys)?;
186
187 let mut rows_by_group: IndexMap<GroupId, Vec<usize>> = IndexMap::new();
188
189 for row in 0..row_count {
190 let mut serializer = KeySerializer::new();
191 for column in &key_columns {
192 column.extend_key(row, &mut serializer);
193 }
194 let encoded = serializer.to_encoded_key();
195
196 let group = dict
197 .intern(&encoded, || key_columns.iter().map(|column| column.get_value(row)).collect());
198 rows_by_group.entry(group).or_default().push(row);
199 }
200
201 Ok(rows_by_group.into_iter().collect())
202 }
203
204 fn key_columns(&self, keys: &[&str]) -> Result<Vec<&ColumnBuffer>> {
205 let mut key_columns: Vec<&ColumnBuffer> = Vec::with_capacity(keys.len());
206 for &key in keys {
207 let pos = self.names.iter().position(|n| n.text() == key).ok_or_else(|| {
208 Error::from(CoreError::FrameError {
209 message: format!("Column '{}' not found", key),
210 })
211 })?;
212 key_columns.push(&self.columns[pos]);
213 }
214 Ok(key_columns)
215 }
216}