Skip to content

Commit eb99d92

Browse files
zzylolclaude
andcommitted
Make MutableEpoch insert O(1) amortized: append-only raw buffer
MutableEpoch is now a plain append-only Vec (no BTreeMap, no sorted Vec maintenance during writes). All sorting is deferred to seal() — paid once at epoch rotation, not on every insert. Mirrors VictoriaMetrics' rawRows → in-memory part pipeline. Insert path: Vec::push + HashSet::insert + 2 scalar min/max updates = O(1). seal() path: sort_unstable_by_key = O(M log M), called once per epoch. Query on active epoch: linear scan O(M), bounded by epoch_capacity × L. Acceptable since sealed epochs hold most historical data and use binary search. Replace min_tr/max_tr (TimestampRange) with min_start/max_end (u64) in both MutableEpoch and SealedEpoch for accurate epoch-skip bounds. Callers updated to use time_bounds() -> Option<(u64, u64)>. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
1 parent 9597ccb commit eb99d92

3 files changed

Lines changed: 92 additions & 112 deletions

File tree

‎asap-query-engine/src/stores/simple_map_store/common.rs‎

Lines changed: 78 additions & 91 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
use crate::data_model::{AggregateCore, KeyByLabelValues};
2-
use std::collections::{BTreeMap, HashMap};
2+
use std::collections::{HashMap, HashSet};
33
use std::sync::Arc;
44

55
pub type MetricID = u32;
@@ -46,122 +46,105 @@ impl InternTable {
4646
}
4747
}
4848

49-
/// Mutable (active) epoch: accepts inserts.
49+
/// Mutable (active) epoch: pure append-only insert, O(1) amortized.
5050
///
51-
/// Optimization 1: flat `HashMap<(MetricID, TimestampRange), _>` replaces the nested
52-
/// `HashMap<MetricID, BTreeMap<TimestampRange, _>>`, giving O(1) inserts and lookups.
51+
/// Raw entries are stored in insertion order — no sorting, no deduplication, no index
52+
/// maintenance during writes. All ordering work is deferred to `seal()`, which is
53+
/// called at most once per epoch (at rotation time). This matches VictoriaMetrics'
54+
/// rawRows → in-memory part pipeline.
5355
///
54-
/// Optimization 3 (time-primary index): `window_to_ids` is a `BTreeMap` keyed by
55-
/// TimestampRange so range queries scan windows in order rather than scanning every
56-
/// label. Each value is a sorted `Vec<MetricID>` (binary-search dedup on insert).
56+
/// Queries on the active epoch do a bounded linear scan (epoch size ≤ epoch_capacity ×
57+
/// labels), which is acceptable because the vast majority of historical data lives in
58+
/// sealed (already-sorted) epochs.
5759
pub struct MutableEpoch {
58-
/// Primary storage: (MetricID, TimestampRange) → aggregates.
59-
pub data: HashMap<(MetricID, TimestampRange), Vec<Arc<dyn AggregateCore>>>,
60-
/// Time-primary index: window → sorted Vec<MetricID>. BTreeMap enables O(log N)
61-
/// range scan without touching labels that have no data in the query window.
62-
pub window_to_ids: BTreeMap<TimestampRange, Vec<MetricID>>,
60+
/// Append-only raw inserts. Sorted only at seal() time.
61+
pub raw: Vec<(TimestampRange, MetricID, Arc<dyn AggregateCore>)>,
62+
/// Distinct windows for rotation threshold — O(1) insert, O(1) len.
63+
windows: HashSet<TimestampRange>,
64+
/// Epoch time bounds for O(1) skip check, updated incrementally on insert.
65+
min_start: Option<u64>,
66+
max_end: Option<u64>,
6367
}
6468

6569
impl MutableEpoch {
6670
pub fn new() -> Self {
6771
Self {
68-
data: HashMap::new(),
69-
window_to_ids: BTreeMap::new(),
72+
raw: Vec::new(),
73+
windows: HashSet::new(),
74+
min_start: None,
75+
max_end: None,
7076
}
7177
}
7278

7379
pub fn window_count(&self) -> usize {
74-
self.window_to_ids.len()
80+
self.windows.len()
7581
}
7682

77-
#[allow(dead_code)]
78-
pub fn is_empty(&self) -> bool {
79-
self.window_to_ids.is_empty()
80-
}
81-
82-
pub fn min_tr(&self) -> Option<TimestampRange> {
83-
self.window_to_ids.keys().next().copied()
84-
}
85-
86-
pub fn max_tr(&self) -> Option<TimestampRange> {
87-
self.window_to_ids.keys().next_back().copied()
83+
/// Returns `(min_start, max_end)` across all windows, or `None` if empty.
84+
/// Used by callers for the epoch-skip check: `min_start > end || max_end < start`.
85+
pub fn time_bounds(&self) -> Option<(u64, u64)> {
86+
match (self.min_start, self.max_end) {
87+
(Some(s), Some(e)) => Some((s, e)),
88+
_ => None,
89+
}
8890
}
8991

92+
/// O(1) amortized: Vec push + HashSet insert + two scalar comparisons.
9093
pub fn insert(
9194
&mut self,
9295
metric_id: MetricID,
9396
range: TimestampRange,
9497
agg: Arc<dyn AggregateCore>,
9598
) {
96-
self.data.entry((metric_id, range)).or_default().push(agg);
97-
// Maintain sorted Vec<MetricID> in time-primary index.
98-
let ids = self.window_to_ids.entry(range).or_default();
99-
let pos = ids.partition_point(|&id| id < metric_id);
100-
if ids.get(pos) != Some(&metric_id) {
101-
ids.insert(pos, metric_id);
102-
}
99+
self.raw.push((range, metric_id, agg));
100+
self.windows.insert(range);
101+
self.min_start = Some(self.min_start.map_or(range.0, |m| m.min(range.0)));
102+
self.max_end = Some(self.max_end.map_or(range.1, |m| m.max(range.1)));
103103
}
104104

105-
/// Seal this epoch into a cache-friendly flat sorted array.
105+
/// Consume this epoch and produce an immutable SealedEpoch by sorting in-place.
106+
/// O(M log M) where M = number of raw entries — paid once at rotation, not at query time.
106107
pub fn seal(self) -> SealedEpoch {
107-
let mut entries: Vec<(TimestampRange, MetricID, Arc<dyn AggregateCore>)> = self
108-
.data
109-
.into_iter()
110-
.flat_map(|((metric_id, tr), aggs)| {
111-
aggs.into_iter().map(move |agg| (tr, metric_id, agg))
112-
})
113-
.collect();
114-
// Sort by (TimestampRange, MetricID): windows contiguous, labels ordered within window.
108+
let min_start = self.min_start;
109+
let max_end = self.max_end;
110+
let mut entries = self.raw;
115111
entries.sort_unstable_by_key(|(tr, metric_id, _)| (*tr, *metric_id));
116-
let min_tr = entries.first().map(|(tr, _, _)| *tr);
117-
let max_tr = entries.last().map(|(tr, _, _)| *tr);
118112
SealedEpoch {
119113
entries,
120-
min_tr,
121-
max_tr,
114+
min_start,
115+
max_end,
122116
}
123117
}
124118

125-
/// Stream results for [start, end] into `out` using the time-primary BTreeMap index.
126-
/// Only visits labels that actually have data in matching windows — O(log N + actual_matches).
119+
/// Linear scan over raw entries for [start, end] — O(M) where M ≤ epoch_capacity × L.
120+
/// Acceptable because: (a) the epoch is bounded, (b) most data is in sealed epochs.
127121
pub fn range_query_into(
128122
&self,
129123
start: u64,
130124
end: u64,
131125
out: &mut MetricBucketMap,
132126
matched_windows: &mut Vec<TimestampRange>,
133127
) {
134-
for (&tr, metric_ids) in self.window_to_ids.range((start, 0)..=(end, u64::MAX)) {
135-
if tr.1 > end {
128+
for (tr, metric_id, agg) in &self.raw {
129+
if tr.0 < start || tr.0 > end || tr.1 > end {
136130
continue;
137131
}
138-
for &metric_id in metric_ids {
139-
if let Some(aggs) = self.data.get(&(metric_id, tr)) {
140-
let slot = out.entry(metric_id).or_default();
141-
for agg in aggs {
142-
slot.push((tr, Arc::clone(agg)));
143-
matched_windows.push(tr);
144-
}
145-
}
146-
}
132+
out.entry(*metric_id)
133+
.or_default()
134+
.push((*tr, Arc::clone(agg)));
135+
matched_windows.push(*tr);
147136
}
148137
}
149138

150-
/// Exact match for a single window using the time-primary index.
139+
/// Linear scan for exact window match — O(M), bounded.
151140
pub fn exact_query(
152141
&self,
153142
range: TimestampRange,
154143
) -> Option<Vec<(MetricID, Arc<dyn AggregateCore>)>> {
155-
let ids = self.window_to_ids.get(&range)?;
156-
if ids.is_empty() {
157-
return None;
158-
}
159144
let mut out = Vec::new();
160-
for &metric_id in ids {
161-
if let Some(aggs) = self.data.get(&(metric_id, range)) {
162-
for agg in aggs {
163-
out.push((metric_id, Arc::clone(agg)));
164-
}
145+
for (tr, metric_id, agg) in &self.raw {
146+
if *tr == range {
147+
out.push((*metric_id, Arc::clone(agg)));
165148
}
166149
}
167150
if out.is_empty() {
@@ -173,36 +156,41 @@ impl MutableEpoch {
173156

174157
/// Remove specific windows (ReadBased cleanup).
175158
pub fn remove_windows(&mut self, windows: &[TimestampRange]) {
176-
for &window in windows {
177-
if let Some(ids) = self.window_to_ids.remove(&window) {
178-
for metric_id in ids {
179-
self.data.remove(&(metric_id, window));
180-
}
181-
}
182-
}
159+
let window_set: HashSet<TimestampRange> = windows.iter().copied().collect();
160+
self.raw.retain(|(tr, _, _)| !window_set.contains(tr));
161+
self.windows.retain(|tr| !window_set.contains(tr));
162+
// Recompute bounds (cleanup is rare, linear scan is fine).
163+
self.min_start = self.raw.iter().map(|(tr, _, _)| tr.0).min();
164+
self.max_end = self.raw.iter().map(|(tr, _, _)| tr.1).max();
183165
}
184166
}
185167

186168
/// Sealed (immutable) epoch: flat sorted `Vec` for cache-friendly range scans.
187169
///
188-
/// Optimization 2: once an epoch is full and rotated, it is converted to a contiguous
189-
/// array sorted by `(TimestampRange, MetricID)`. Range queries use binary search to
190-
/// find the start position and then do a linear scan — no pointer chasing through
191-
/// nested HashMap/BTreeMap nodes.
170+
/// Produced by `MutableEpoch::seal()`. Entries are sorted by `(TimestampRange, MetricID)`:
171+
/// all entries for the same window are contiguous, which is cache-friendly for both
172+
/// range queries (binary-search start + linear scan) and exact queries.
192173
pub struct SealedEpoch {
193-
/// Sorted by (TimestampRange, MetricID). All entries for the same window are
194-
/// contiguous; within a window entries are ordered by MetricID.
174+
/// Sorted by (TimestampRange, MetricID).
195175
pub entries: Vec<(TimestampRange, MetricID, Arc<dyn AggregateCore>)>,
196-
/// Precomputed min/max for O(1) epoch-skip check.
197-
pub min_tr: Option<TimestampRange>,
198-
pub max_tr: Option<TimestampRange>,
176+
/// Precomputed for O(1) epoch-skip check.
177+
pub min_start: Option<u64>,
178+
pub max_end: Option<u64>,
199179
}
200180

201181
impl SealedEpoch {
202182
pub fn is_empty(&self) -> bool {
203183
self.entries.is_empty()
204184
}
205185

186+
/// Returns `(min_start, max_end)`, or `None` if empty.
187+
pub fn time_bounds(&self) -> Option<(u64, u64)> {
188+
match (self.min_start, self.max_end) {
189+
(Some(s), Some(e)) => Some((s, e)),
190+
_ => None,
191+
}
192+
}
193+
206194
/// Binary-search start + linear scan — O(log N + actual_matches), cache-friendly.
207195
pub fn range_query_into(
208196
&self,
@@ -246,16 +234,15 @@ impl SealedEpoch {
246234
}
247235
}
248236

249-
/// Remove specific windows (ReadBased cleanup). Rebuilds the Vec in one pass.
237+
/// Remove specific windows (ReadBased cleanup). Rebuilds Vec in one pass.
250238
pub fn remove_windows(&mut self, windows: &[TimestampRange]) {
251-
let window_set: std::collections::HashSet<TimestampRange> =
252-
windows.iter().copied().collect();
239+
let window_set: HashSet<TimestampRange> = windows.iter().copied().collect();
253240
self.entries.retain(|(tr, _, _)| !window_set.contains(tr));
254-
self.min_tr = self.entries.first().map(|(tr, _, _)| *tr);
255-
self.max_tr = self.entries.last().map(|(tr, _, _)| *tr);
241+
self.min_start = self.entries.iter().map(|(tr, _, _)| tr.0).min();
242+
self.max_end = self.entries.iter().map(|(tr, _, _)| tr.1).max();
256243
}
257244

258-
/// Deduplicated windows (entries are sorted so consecutive dupes are adjacent).
245+
/// Deduplicated windows (entries sorted, so consecutive dupes are adjacent).
259246
/// Used to purge `read_counts` when this epoch is dropped.
260247
pub fn unique_windows(&self) -> Vec<TimestampRange> {
261248
let mut windows: Vec<TimestampRange> = self.entries.iter().map(|(tr, _, _)| *tr).collect();

‎asap-query-engine/src/stores/simple_map_store/global.rs‎

Lines changed: 7 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -362,11 +362,8 @@ impl Store for SimpleMapStoreGlobal {
362362
let mut mid: MetricBucketMap = HashMap::with_capacity(per_key.intern.len());
363363

364364
// Query current (mutable) epoch.
365-
if let (Some(min), Some(max)) = (
366-
per_key.current_epoch.min_tr(),
367-
per_key.current_epoch.max_tr(),
368-
) {
369-
if !(min.0 > end || max.1 < start) {
365+
if let Some((min_start, max_end)) = per_key.current_epoch.time_bounds() {
366+
if !(min_start > end || max_end < start) {
370367
per_key.current_epoch.range_query_into(
371368
start,
372369
end,
@@ -378,13 +375,11 @@ impl Store for SimpleMapStoreGlobal {
378375

379376
// Query sealed epochs; skip those with no overlap.
380377
for epoch in per_key.sealed_epochs.values() {
381-
match (epoch.min_tr, epoch.max_tr) {
382-
(Some(min), Some(max)) => {
383-
if min.0 > end || max.1 < start {
384-
continue;
385-
}
386-
}
387-
_ => continue, // empty epoch
378+
let Some((min_start, max_end)) = epoch.time_bounds() else {
379+
continue;
380+
};
381+
if min_start > end || max_end < start {
382+
continue;
388383
}
389384
epoch.range_query_into(start, end, &mut mid, &mut matched_windows);
390385
}

‎asap-query-engine/src/stores/simple_map_store/per_key.rs‎

Lines changed: 7 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -450,22 +450,20 @@ impl Store for SimpleMapStorePerKey {
450450
let mut mid: MetricBucketMap = HashMap::with_capacity(data.intern.len());
451451

452452
// Query current (mutable) epoch.
453-
if let (Some(min), Some(max)) = (data.current_epoch.min_tr(), data.current_epoch.max_tr()) {
454-
if !(min.0 > end || max.1 < start) {
453+
if let Some((min_start, max_end)) = data.current_epoch.time_bounds() {
454+
if !(min_start > end || max_end < start) {
455455
data.current_epoch
456456
.range_query_into(start, end, &mut mid, &mut matched_windows);
457457
}
458458
}
459459

460460
// Query sealed epochs; skip those with no overlap.
461461
for epoch in data.sealed_epochs.values() {
462-
match (epoch.min_tr, epoch.max_tr) {
463-
(Some(min), Some(max)) => {
464-
if min.0 > end || max.1 < start {
465-
continue;
466-
}
467-
}
468-
_ => continue, // empty epoch
462+
let Some((min_start, max_end)) = epoch.time_bounds() else {
463+
continue;
464+
};
465+
if min_start > end || max_end < start {
466+
continue;
469467
}
470468
epoch.range_query_into(start, end, &mut mid, &mut matched_windows);
471469
}

0 commit comments

Comments
 (0)