Skip to content

Commit b9c8037

Browse files
zzylolclaude
andcommitted
Apply three VictoriaMetrics-inspired optimizations to SimpleMapStore
1. Label interning (MetricID: u32) Introduce InternTable in common.rs that assigns a compact u32 to each unique Option<KeyByLabelValues> on first insert. All internal maps (label_map, window_to_ids) use MetricID as key — O(1) hash/compare instead of O(label_bytes). Label strings stored once; resolved back to KeyByLabelValues only when building the returned TimestampedBucketsMap. 2. Time-epoch partitioning + O(1) rotation cleanup Replace the single BTreeMap-per-label with epoch-partitioned storage: BTreeMap<EpochID, EpochData>. epoch_capacity = num_aggregates_to_retain (set on first insert). When the current epoch reaches capacity, a new epoch is opened and the oldest is dropped — O(1) drop vs the previous O(k·m) BTreeSet walk + targeted BTreeMap removals. ReadBased cleanup scans read_counts then calls EpochData::remove_windows on each epoch. 3. Sorted Vec posting lists Replace HashSet<Option<KeyByLabelValues>> in window_to_labels with Vec<MetricID> maintained in sorted order via partition_point + insert. Cache-friendly iteration for exact queries; enables merge-intersection for future label-predicate pushdown. Both per_key.rs and global.rs updated. Shared types extracted to common.rs. Public Store trait interface and all 329 existing tests are unchanged. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
1 parent 75d6d8b commit b9c8037

4 files changed

Lines changed: 588 additions & 533 deletions

File tree

Lines changed: 160 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,160 @@
1+
use crate::data_model::{AggregateCore, KeyByLabelValues};
2+
use std::collections::{BTreeMap, BTreeSet, HashMap};
3+
use std::sync::Arc;
4+
5+
pub type MetricID = u32;
6+
pub type EpochID = u64;
7+
pub type TimestampRange = (u64, u64);
8+
9+
/// Assigns a compact MetricID (u32) to each unique label combination.
10+
/// Label strings stored once; all internal maps use MetricID (O(1) key ops).
11+
pub struct InternTable {
12+
label_to_id: HashMap<Option<KeyByLabelValues>, MetricID>,
13+
id_to_label: Vec<Option<KeyByLabelValues>>,
14+
}
15+
16+
impl InternTable {
17+
pub fn new() -> Self {
18+
Self {
19+
label_to_id: HashMap::new(),
20+
id_to_label: Vec::new(),
21+
}
22+
}
23+
24+
/// Intern a label, assigning a new MetricID if first seen.
25+
/// Uses HashMap::entry to avoid double-hashing.
26+
pub fn intern(&mut self, label: Option<KeyByLabelValues>) -> MetricID {
27+
let next_id = self.id_to_label.len() as MetricID;
28+
match self.label_to_id.entry(label) {
29+
std::collections::hash_map::Entry::Occupied(e) => *e.get(),
30+
std::collections::hash_map::Entry::Vacant(e) => {
31+
self.id_to_label.push(e.key().clone());
32+
*e.insert(next_id)
33+
}
34+
}
35+
}
36+
37+
/// O(1) resolution by MetricID.
38+
pub fn resolve(&self, id: MetricID) -> &Option<KeyByLabelValues> {
39+
&self.id_to_label[id as usize]
40+
}
41+
}
42+
43+
/// One epoch slot: holds up to `epoch_capacity` distinct time windows.
44+
pub struct EpochData {
45+
/// Primary inverted index: MetricID → time-sorted aggregates.
46+
pub label_map: HashMap<MetricID, BTreeMap<TimestampRange, Vec<Arc<dyn AggregateCore>>>>,
47+
/// Reverse index: window → sorted Vec<MetricID> (Optimization 3).
48+
pub window_to_ids: HashMap<TimestampRange, Vec<MetricID>>,
49+
/// All distinct time windows in this epoch, sorted.
50+
pub time_ranges: BTreeSet<TimestampRange>,
51+
}
52+
53+
impl EpochData {
54+
pub fn new() -> Self {
55+
Self {
56+
label_map: HashMap::new(),
57+
window_to_ids: HashMap::new(),
58+
time_ranges: BTreeSet::new(),
59+
}
60+
}
61+
62+
pub fn window_count(&self) -> usize {
63+
self.time_ranges.len()
64+
}
65+
66+
pub fn is_empty(&self) -> bool {
67+
self.time_ranges.is_empty()
68+
}
69+
70+
/// Insert (metric_id, range, aggregate) into this epoch.
71+
pub fn insert(
72+
&mut self,
73+
metric_id: MetricID,
74+
range: TimestampRange,
75+
agg: Arc<dyn AggregateCore>,
76+
) {
77+
self.time_ranges.insert(range);
78+
self.label_map
79+
.entry(metric_id)
80+
.or_default()
81+
.entry(range)
82+
.or_default()
83+
.push(agg);
84+
// Maintain sorted Vec<MetricID> in reverse index
85+
let ids = self.window_to_ids.entry(range).or_default();
86+
let pos = ids.partition_point(|&id| id < metric_id);
87+
if ids.get(pos) != Some(&metric_id) {
88+
ids.insert(pos, metric_id);
89+
}
90+
}
91+
92+
/// Remove windows from this epoch (ReadBased cleanup).
93+
pub fn remove_windows(&mut self, windows: &[TimestampRange]) {
94+
for &window in windows {
95+
self.time_ranges.remove(&window);
96+
let Some(ids) = self.window_to_ids.remove(&window) else {
97+
continue;
98+
};
99+
for metric_id in ids {
100+
let remove_label = if let Some(btree) = self.label_map.get_mut(&metric_id) {
101+
btree.remove(&window);
102+
btree.is_empty()
103+
} else {
104+
false
105+
};
106+
if remove_label {
107+
self.label_map.remove(&metric_id);
108+
}
109+
}
110+
}
111+
}
112+
113+
/// Collect all results matching [start, end].
114+
pub fn range_query(
115+
&self,
116+
start: u64,
117+
end: u64,
118+
) -> Vec<(MetricID, TimestampRange, Arc<dyn AggregateCore>)> {
119+
let mut out = Vec::new();
120+
for (&metric_id, btree) in &self.label_map {
121+
for (&tr, aggs) in btree.range((start, 0)..=(end, u64::MAX)) {
122+
if tr.1 > end {
123+
continue;
124+
}
125+
for agg in aggs {
126+
out.push((metric_id, tr, Arc::clone(agg)));
127+
}
128+
}
129+
}
130+
out
131+
}
132+
133+
/// Collect results for an exact window match using the reverse index.
134+
pub fn exact_query(
135+
&self,
136+
range: TimestampRange,
137+
) -> Option<Vec<(MetricID, Arc<dyn AggregateCore>)>> {
138+
let ids = self.window_to_ids.get(&range)?;
139+
if ids.is_empty() {
140+
return None;
141+
}
142+
let mut out = Vec::new();
143+
for &metric_id in ids {
144+
if let Some(aggs) = self
145+
.label_map
146+
.get(&metric_id)
147+
.and_then(|b| b.get(&range))
148+
{
149+
for agg in aggs {
150+
out.push((metric_id, Arc::clone(agg)));
151+
}
152+
}
153+
}
154+
if out.is_empty() {
155+
None
156+
} else {
157+
Some(out)
158+
}
159+
}
160+
}

0 commit comments

Comments
 (0)