|
2 | 2 |
|
3 | 3 | ## Overview |
4 | 4 |
|
5 | | -The `SimpleMapStore` uses an **inverted index** (label-primary) layout to store precomputed aggregates. This design aligns the storage structure with the query return type (`HashMap<Option<KeyByLabelValues>, Vec<TimestampedBucket>>`), eliminating the need for regrouping at query time. |
| 5 | +`SimpleMapStore` uses an **epoch-partitioned inverted index** to store precomputed aggregates. Three VictoriaMetrics-inspired optimizations are applied on top of the basic label-primary layout: |
6 | 6 |
|
7 | | -## Data Structure |
| 7 | +1. **Label Interning** — label combinations are mapped to compact `MetricID` (u32), reducing key size and hash cost. |
| 8 | +2. **Epoch Partitioning** — data is split into fixed-capacity epoch slots; the oldest epoch is dropped O(1) when the cap is exceeded (CircularBuffer policy). |
| 9 | +3. **Sorted Vec Posting Lists** — the reverse index (`window_to_ids`) stores `Vec<MetricID>` maintained in sorted order, enabling binary-search deduplication on insert and cache-friendly iteration on lookup. |
8 | 10 |
|
9 | | -### Per-Key Store (`per_key.rs`) |
| 11 | +--- |
10 | 12 |
|
11 | | -Each `aggregation_id` maps to a `StoreKeyData` protected by an `RwLock`: |
| 13 | +## Data Structures |
| 14 | + |
| 15 | +### Types (common.rs) |
| 16 | + |
| 17 | +```rust |
| 18 | +pub type MetricID = u32; // compact interned label ID |
| 19 | +pub type EpochID = u64; // monotonically increasing epoch counter |
| 20 | +pub type TimestampRange = (u64, u64); // (start_timestamp, end_timestamp) |
| 21 | +pub type MetricBucketMap = HashMap<MetricID, Vec<(TimestampRange, Arc<dyn AggregateCore>)>>; |
| 22 | +``` |
| 23 | + |
| 24 | +### InternTable (common.rs) |
| 25 | + |
| 26 | +``` |
| 27 | +InternTable { |
| 28 | + label_to_id: HashMap<Option<KeyByLabelValues>, MetricID> |
| 29 | + id_to_label: Vec<Option<KeyByLabelValues>> |
| 30 | +} |
| 31 | +``` |
| 32 | + |
| 33 | +- `intern(label)` → O(1) amortized, no double-hashing (uses `HashMap::entry`) |
| 34 | +- `resolve(id)` → O(1) indexed Vec lookup |
| 35 | +- All internal index maps use `MetricID` (u32) as keys, not full label strings |
| 36 | + |
| 37 | +### EpochData (common.rs) |
| 38 | + |
| 39 | +One epoch holds up to `epoch_capacity` distinct time windows. |
| 40 | + |
| 41 | +``` |
| 42 | +EpochData { |
| 43 | + label_map: HashMap<MetricID, BTreeMap<TimestampRange, Vec<Arc<dyn AggregateCore>>>> |
| 44 | + window_to_ids: HashMap<TimestampRange, Vec<MetricID>> // sorted (Optimization 3) |
| 45 | + time_ranges: BTreeSet<TimestampRange> |
| 46 | +} |
| 47 | +``` |
| 48 | + |
| 49 | +- **`label_map`** (primary index): inverted index MetricID → time-sorted BTreeMap of aggregates. Enables O(log N + k) range queries per label. |
| 50 | +- **`window_to_ids`** (reverse index): for each time window, sorted `Vec<MetricID>` of labels that contain data. Used for exact queries and targeted cleanup without full label scans. |
| 51 | +- **`time_ranges`** (secondary index): all distinct windows in this epoch, sorted. Used for epoch range filtering (skip epochs that don't overlap the query interval) and cleanup ordering. |
| 52 | + |
| 53 | +### Per-Key Store (per_key.rs) |
| 54 | + |
| 55 | +Each aggregation_id gets its own `StoreKeyData` behind a per-key `RwLock`: |
12 | 56 |
|
13 | 57 | ``` |
14 | 58 | DashMap<aggregation_id, Arc<RwLock<StoreKeyData>>> |
15 | 59 |
|
16 | 60 | StoreKeyData { |
17 | | - label_map: HashMap<Option<KeyByLabelValues>, BTreeMap<(start, end), Vec<Arc<dyn AggregateCore>>>> |
18 | | - window_to_labels: HashMap<(start, end), HashSet<Option<KeyByLabelValues>>> |
19 | | - time_ranges: BTreeSet<(start, end)> |
20 | | - read_counts: Mutex<HashMap<(start, end), u64>> |
| 61 | + intern: InternTable |
| 62 | + epochs: BTreeMap<EpochID, EpochData> |
| 63 | + current_epoch_id: EpochID |
| 64 | + epoch_capacity: Option<usize> // None = unlimited |
| 65 | + max_epochs: usize // default 4 |
| 66 | + read_counts: Mutex<HashMap<TimestampRange, u64>> |
21 | 67 | } |
22 | 68 | ``` |
23 | 69 |
|
24 | | -- **`label_map`** (primary index): Inverted index from label key to a time-sorted BTreeMap of aggregates. Enables O(log n + k) range queries per label. |
25 | | -- **`window_to_labels`** (reverse index): For each time window, tracks exactly which labels contain data. Enables exact queries and cleanup to avoid full label scans. |
26 | | -- **`time_ranges`** (secondary index): All known timestamp ranges across all labels. Used for cleanup counting and read-count tracking. |
27 | | -- **`read_counts`**: Wrapped in `Mutex` so queries can use a read lock on the outer `RwLock` (only needs brief exclusive access to increment counts). |
| 70 | +`read_counts` is behind an inner `Mutex` so queries can hold a read lock on the outer `RwLock` and still update counts (brief inner lock, no write-lock upgrade needed). |
28 | 71 |
|
29 | | -### Global Store (`global.rs`) |
| 72 | +### Global Store (global.rs) |
30 | 73 |
|
31 | | -Same inverted index structure, but nested under a single `Mutex<StoreData>`: |
| 74 | +Same per-key epoch structure, but all aggregation_ids share a single `Mutex<StoreData>`: |
32 | 75 |
|
33 | 76 | ``` |
34 | 77 | Mutex<StoreData> |
35 | 78 |
|
36 | 79 | StoreData { |
37 | | - store: HashMap<aggregation_id, HashMap<Option<KeyByLabelValues>, BTreeMap<(start, end), Vec<Arc<dyn AggregateCore>>>>> |
38 | | - window_to_labels: HashMap<aggregation_id, HashMap<(start, end), HashSet<Option<KeyByLabelValues>>>> |
39 | | - time_ranges: HashMap<aggregation_id, BTreeSet<(start, end)>> |
40 | | - read_counts: HashMap<aggregation_id, HashMap<(start, end), u64>> |
| 80 | + stores: HashMap<aggregation_id, PerKeyState> |
| 81 | + read_counts: HashMap<aggregation_id, HashMap<TimestampRange, u64>> |
| 82 | +} |
| 83 | +
|
| 84 | +PerKeyState { |
| 85 | + intern: InternTable |
| 86 | + epochs: BTreeMap<EpochID, EpochData> |
| 87 | + current_epoch_id: EpochID |
| 88 | + epoch_capacity: Option<usize> |
| 89 | + max_epochs: usize |
41 | 90 | } |
42 | 91 | ``` |
43 | 92 |
|
44 | | -No inner Mutex for `read_counts` since the outer Mutex already serializes all access. |
| 93 | +No inner `Mutex` for `read_counts` — the outer `Mutex` already serializes all access. |
| 94 | + |
| 95 | +--- |
45 | 96 |
|
46 | 97 | ## Theoretical Complexity |
47 | 98 |
|
48 | 99 | ### Variables |
49 | 100 |
|
50 | 101 | | Symbol | Meaning | |
51 | | -|---|---| |
| 102 | +|--------|---------| |
52 | 103 | | A | Number of distinct aggregation IDs | |
53 | 104 | | L | Number of distinct label combinations (cardinality) | |
54 | 105 | | N | Number of distinct time windows stored per (agg_id, label) | |
55 | | -| k | Number of results matched or entries removed in a given operation | |
| 106 | +| E | Number of epochs (bounded by `max_epochs`, default 4) | |
| 107 | +| k | Number of results matched or entries removed | |
56 | 108 | | m | Number of labels present in a specific time window | |
57 | | -| V | Number of aggregate objects stored per (label, window) slot (typically 1) | |
| 109 | +| V | Aggregate objects per (label, window) slot (typically 1) | |
58 | 110 |
|
59 | 111 | ### Time Complexity |
60 | 112 |
|
61 | 113 | | Operation | Time | Notes | |
62 | | -|---|---|---| |
63 | | -| **Insert** (single entry) | **O(log N)** | DashMap O(1) + RwLock O(1) + HashMap O(1) + BTreeMap O(log N) + BTreeSet O(log N) | |
64 | | -| **Insert** (batch of B entries, same agg_id) | **O(B · log N)** | One write-lock acquisition amortized over B items | |
65 | | -| **Range query** | **O(L · (log N + k))** | BTreeMap::range per label in O(log N + k_L); results already grouped by label | |
66 | | -| **Exact query** | **O(m · log N)** | window_to_labels lookup O(1) + BTreeMap point get O(log N) per matching label | |
67 | | -| **CircularBuffer cleanup** | **O(k · m)** amortized | BTreeSet iteration O(k) + targeted label-map removals via window_to_labels | |
68 | | -| **ReadBased cleanup** | **O(N + k · m)** | Full read_counts scan O(N) + targeted removals O(k · m) | |
69 | | -| **get_earliest_timestamp** | **O(A)** | DashMap iteration over A entries with atomic loads | |
| 114 | +|-----------|------|-------| |
| 115 | +| **Insert** (single entry) | O(log N) | DashMap O(1) + RwLock O(1) + InternTable O(1) + BTreeMap O(log N) + BTreeSet O(log N) + sorted-Vec insert O(L) worst | |
| 116 | +| **Insert** (batch B, same agg_id) | O(B · log N) | One write-lock acquisition amortized over B items | |
| 117 | +| **Epoch rotation** (CircularBuffer) | O(1) amortized | BTreeMap insert new epoch + BTreeMap pop oldest | |
| 118 | +| **Range query** | O(E · L · (log N + k)) | Per epoch: skip check O(1) + range scan per label O(log N + k_L); MetricID→label resolution O(L) | |
| 119 | +| **Exact query** | O(E · m · log N) | Per epoch: reverse-index lookup O(1) + point get O(log N) per matching label; stops at first match | |
| 120 | +| **CircularBuffer cleanup** | O(1) amortized | Epoch rotation drops entire oldest epoch | |
| 121 | +| **ReadBased cleanup** | O(N + k · m) | Scan read_counts O(N) + targeted removals via window_to_ids O(k · m) | |
| 122 | +| **get_earliest_timestamp** | O(A) | DashMap iteration with AtomicU64 loads | |
70 | 123 |
|
71 | 124 | ### Space Complexity |
72 | 125 |
|
73 | 126 | | Structure | Space | Notes | |
74 | | -|---|---|---| |
75 | | -| `label_map` | O(A · L · N · V) | Primary index: agg_id → label → BTreeMap(window → Vec<Arc<Agg>>) | |
76 | | -| `window_to_labels` | O(A · N · L) | Reverse index: agg_id → window → HashSet\<label\> | |
77 | | -| `time_ranges` | O(A · N) | Secondary index: agg_id → BTreeSet of all windows | |
78 | | -| `read_counts` | O(A · N) | agg_id → HashMap\<window, u64\> | |
79 | | -| **Total** | **O(A · L · N · V)** | Dominated by the primary label_map | |
80 | | - |
81 | | -Arc-sharing means query results hold references into the store; no deep copies are made for read paths. |
82 | | - |
83 | | -### Operation Complexity Summary |
84 | | - |
85 | | -| Operation | Complexity | |
86 | | -|---|---| |
87 | | -| Range query | O(L × (log N + k)) via `BTreeMap::range()`, already grouped by label | |
88 | | -| Exact query | O(m × log N) where m = labels present in target window (via reverse index) | |
89 | | -| Insert | O(log N) BTreeMap insert per label | |
90 | | -| CircularBuffer cleanup | O(k × m) iterate first k from `BTreeSet` + targeted removals via `window_to_labels` | |
91 | | -| ReadBased cleanup | O(N + k × m) scan `read_counts` + targeted removals via `window_to_labels` | |
92 | | -| Space | O(A × L × N × V) — proportional to stored aggregates, not index overhead | |
| 127 | +|-----------|-------|-------| |
| 128 | +| `InternTable` | O(L) per agg_id | Stores each label string once | |
| 129 | +| `label_map` (per epoch) | O(L · N · V) | Primary index across all epochs | |
| 130 | +| `window_to_ids` | O(N · m) | Reverse index, bounded by epoch | |
| 131 | +| `time_ranges` | O(N) per epoch | BTreeSet of distinct windows | |
| 132 | +| `read_counts` | O(N) total | Counts keyed by TimestampRange | |
| 133 | +| **Total** | **O(A · E · L · N · V)** | E bounded by `max_epochs` (default 4); dominated by label_map | |
| 134 | + |
| 135 | +Arc-sharing means query results reference aggregate objects already in the store — no deep copies on read paths. |
| 136 | + |
| 137 | +--- |
93 | 138 |
|
94 | 139 | ## Query Mechanics |
95 | 140 |
|
96 | 141 | ### Range Query |
97 | 142 |
|
98 | | -For a query with `[start, end]`: |
| 143 | +For a query `[start, end]`: |
99 | 144 |
|
100 | | -1. For each label in `label_map`, use `btree.range((start, 0)..=(end, u64::MAX))` to find candidate entries in O(log n) |
101 | | -2. Filter by `range_end <= end` (BTreeMap range only bounds `range_start`) |
102 | | -3. Results are already in chronological order (BTreeMap iteration order) and grouped by label |
103 | | -4. Update `read_counts` via the `time_ranges` secondary index |
| 145 | +1. Acquire **read lock** on `StoreKeyData` (concurrent queries run in parallel) |
| 146 | +2. For each epoch in `epochs.values()`: |
| 147 | + - Skip if `min_tr.0 > end || max_tr.1 < start` (epoch range check, O(1) via BTreeSet first/last) |
| 148 | + - For each label in `label_map`, call `btree.range((start, 0)..=(end, u64::MAX))`, filter `tr.1 <= end` |
| 149 | + - Stream results directly into a `MetricBucketMap` (grouped by MetricID, no intermediate flat vec) |
| 150 | +3. Resolve MetricIDs → label strings in one pass via `InternTable` |
| 151 | +4. Lock inner `Mutex` briefly to update `read_counts` |
104 | 152 |
|
105 | 153 | ### Exact Query |
106 | 154 |
|
107 | 155 | For exact match `(exact_start, exact_end)`: |
108 | 156 |
|
109 | | -1. Use `window_to_labels` to get labels that actually have that window |
110 | | -2. For those labels only, use `btree.get(&(exact_start, exact_end))` for O(log n) lookup |
111 | | -2. Results are already grouped by label |
| 157 | +1. Acquire **read lock** |
| 158 | +2. Iterate epochs newest-first (`epochs.values().rev()`): |
| 159 | + - Use `window_to_ids.get(&range)` to get the sorted `Vec<MetricID>` of labels with that window |
| 160 | + - For each MetricID, use `label_map[id].get(&range)` — O(log N) point lookup |
| 161 | + - Stop at the first epoch that has the window (break after first match) |
| 162 | +3. Resolve MetricIDs → labels, update `read_counts` |
| 163 | + |
| 164 | +--- |
112 | 165 |
|
113 | 166 | ## Cleanup Policies |
114 | 167 |
|
115 | 168 | ### CircularBuffer |
116 | 169 |
|
117 | | -Retains the newest `configured_limit * 4` time ranges: |
| 170 | +Epoch-based eviction — O(1) amortized: |
118 | 171 |
|
119 | | -1. Check `time_ranges.len()` against the retention limit |
120 | | -2. Iterate `time_ranges` from the start (oldest first, already sorted by BTreeSet) |
121 | | -3. Remove excess entries from `time_ranges`, `read_counts`, and reverse index |
122 | | -4. Remove from only affected label BTrees using `window_to_labels` membership |
| 172 | +1. On first insert, set `epoch_capacity` from `num_aggregates_to_retain` |
| 173 | +2. After each item insert, call `maybe_rotate_epoch()`: |
| 174 | + - If current epoch's `window_count() >= epoch_capacity`, open a new epoch (`current_epoch_id + 1`) |
| 175 | + - If `epochs.len() > max_epochs`, pop the oldest epoch (BTreeMap first entry) — O(1) drop of entire epoch |
| 176 | + - Purge dropped epoch's windows from `read_counts` |
123 | 177 |
|
124 | 178 | ### ReadBased |
125 | 179 |
|
126 | | -Removes entries that have been read `>= threshold` times: |
| 180 | +Read-count triggered eviction: |
| 181 | + |
| 182 | +1. Scan `read_counts` for windows with `count >= threshold` |
| 183 | +2. For each such window, call `EpochData::remove_windows()`: |
| 184 | + - Remove from `time_ranges`, `window_to_ids`, and only the affected label BTrees (via sorted `Vec<MetricID>` from reverse index) |
| 185 | +3. Drop any epochs that are now empty; re-create `current_epoch_id` entry if it was dropped |
| 186 | + |
| 187 | +### NoCleanup |
| 188 | + |
| 189 | +No eviction — data accumulates indefinitely. |
127 | 190 |
|
128 | | -1. Scan `read_counts` for entries meeting the threshold |
129 | | -2. Remove from `read_counts`, `time_ranges`, and reverse index |
130 | | -3. Remove from only affected label BTrees using `window_to_labels` membership |
| 191 | +--- |
131 | 192 |
|
132 | 193 | ## Concurrency (Per-Key Store) |
133 | 194 |
|
134 | | -The per-key store uses a read-lock optimization: |
| 195 | +| Operation | Lock acquired | |
| 196 | +|-----------|--------------| |
| 197 | +| **Insert** | `DashMap` shard lock (briefly) → `RwLock::write` for the duration of the batch | |
| 198 | +| **Range/Exact query** | `DashMap` shard lock (briefly) → `RwLock::read` (concurrent queries run in parallel) → `Mutex::lock` on `read_counts` (briefly, while holding read lock) | |
| 199 | +| **Cleanup** | Runs under the existing write lock; accesses `read_counts` via `Mutex::get_mut()` (no lock overhead — `&mut self` guarantees exclusivity) | |
135 | 200 |
|
136 | | -- **Insert**: Acquires a write lock on the `RwLock` (exclusive access needed for `label_map` and `time_ranges`) |
137 | | -- **Query**: Acquires a read lock on the `RwLock` (multiple queries can run concurrently). Updates `read_counts` by briefly locking the inner `Mutex` |
138 | | -- **Cleanup**: Runs during insert (under write lock), accesses `read_counts` via `Mutex::get_mut()` (no lock needed since `&mut self` guarantees exclusive access) |
| 201 | +Multiple readers per aggregation_id can proceed concurrently. Writers only block readers of the same aggregation_id, not other aggregation_ids. |
0 commit comments