Skip to main content

mito2/compaction/
overlap.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::collections::HashMap;
16use std::ops::Range;
17
18use common_time::Timestamp;
19use store_api::metadata::{RegionMetadata, RegionMetadataRef};
20
21use crate::compaction::run::primary_key_ranges_overlap;
22use crate::sst::file::{FileHandle, RegionFileId};
23use crate::sst::primary_key::PrimaryKeyRangeMapper;
24
25/// Snapshot-local overlap candidates, ordered by start time. Each subtree stores
26/// its maximum active end. Removing visited files prunes dense internal overlaps.
27/// PK-disjoint files with overlapping times can still cost O(N) per query, so a
28/// complete closure has O(N²) worst-case time and O(N) auxiliary space.
29pub(super) struct FileOverlapIndex<'a> {
30    /// Current region metadata used to decide whether PK bounds can safely exclude overlaps.
31    metadata: &'a RegionMetadata,
32    primary_key_mapper: PrimaryKeyRangeMapper,
33    /// Candidates in original snapshot order, retained after removal so their
34    /// indices remain stable and drained matches can preserve merge input order.
35    files: Vec<FileHandle>,
36    /// Indices into `files`, sorted by start time with the snapshot index as a
37    /// tie-breaker. Sorted offset `i` corresponds to leaf `leaf_base + i`.
38    by_start: Vec<usize>,
39    /// Active region/file identities mapped to leaf indices in `max_ends` for
40    /// removal without searching `by_start`. Entries are deleted when visited.
41    positions: HashMap<RegionFileId, usize>,
42    /// Array-backed segment tree rooted at index 1 (index 0 is unused).
43    /// Leaves store inclusive file end times; internal nodes store the maximum
44    /// active end in their subtree. Removed/padding leaves and empty subtrees
45    /// hold `None`, allowing overlap queries to prune them.
46    max_ends: Vec<Option<Timestamp>>,
47    /// First leaf index in `max_ends`, also the padded leaf count: a power of two
48    /// with at least one leaf, even for an empty snapshot.
49    leaf_base: usize,
50}
51
52struct OverlapQuery<'a> {
53    input: &'a FileHandle,
54    min_end: Timestamp,
55    candidates_end: usize,
56}
57
58impl<'a> FileOverlapIndex<'a> {
59    pub(super) fn new(files: Vec<FileHandle>, metadata: &'a RegionMetadataRef) -> Self {
60        let mut by_start: Vec<_> = (0..files.len()).collect();
61        by_start.sort_unstable_by_key(|&i| (files[i].time_range().0, i));
62        // A single empty leaf also handles an empty snapshot.
63        let leaf_base = files.len().next_power_of_two();
64        let mut max_ends = vec![None; 2 * leaf_base];
65        let mut positions = HashMap::with_capacity(files.len());
66        for (offset, &i) in by_start.iter().enumerate() {
67            max_ends[leaf_base + offset] = Some(files[i].time_range().1);
68            positions.insert(files[i].file_id(), leaf_base + offset);
69        }
70        for node in (1..leaf_base).rev() {
71            max_ends[node] = max_ends[2 * node].max(max_ends[2 * node + 1]);
72        }
73        Self {
74            metadata,
75            primary_key_mapper: PrimaryKeyRangeMapper::new(metadata.clone()),
76            files,
77            by_start,
78            positions,
79            max_ends,
80            leaf_base,
81        }
82    }
83
84    pub(super) fn remove(&mut self, file_id: RegionFileId) {
85        let Some(mut node) = self.positions.remove(&file_id) else {
86            return;
87        };
88        self.max_ends[node] = None;
89        while node > 1 {
90            node /= 2;
91            self.max_ends[node] = self.max_ends[2 * node].max(self.max_ends[2 * node + 1]);
92        }
93    }
94
95    /// Removes dependencies so expansion never enumerates visited files.
96    /// Preserve snapshot order rather than exposing the index's time ordering
97    /// to the merge reader (including its handling of equal-sequence rows).
98    pub(super) fn drain_overlaps(&mut self, input: &FileHandle) -> Vec<FileHandle> {
99        let mut matches = Vec::new();
100        while let Some(i) = self.find_overlap(input) {
101            self.remove(self.files[i].file_id());
102            matches.push(i);
103        }
104        matches.sort_unstable();
105        matches.into_iter().map(|i| self.files[i].clone()).collect()
106    }
107
108    fn find_overlap(&self, input: &FileHandle) -> Option<usize> {
109        let (start, end) = input.time_range();
110        let query = OverlapQuery {
111            input,
112            min_end: start,
113            candidates_end: self
114                .by_start
115                .partition_point(|&i| self.files[i].time_range().0 <= end),
116        };
117        self.find_in_subtree(1, 0..self.leaf_base, &query)
118    }
119
120    fn find_in_subtree(
121        &self,
122        node: usize,
123        leaves: Range<usize>,
124        query: &OverlapQuery<'_>,
125    ) -> Option<usize> {
126        if leaves.start >= query.candidates_end
127            || self.max_ends[node].is_none_or(|end| end < query.min_end)
128        {
129            return None;
130        }
131        if leaves.len() == 1 {
132            let i = self.by_start[leaves.start];
133            return files_may_overlap(
134                query.input,
135                &self.files[i],
136                self.metadata,
137                &self.primary_key_mapper,
138            )
139            .then_some(i);
140        }
141        let mid = leaves.start + leaves.len() / 2;
142        self.find_in_subtree(2 * node, leaves.start..mid, query)
143            .or_else(|| self.find_in_subtree(2 * node + 1, mid..leaves.end, query))
144    }
145}
146
147/// SST bounds are inclusive. Missing statistics, foreign encodings and schema
148/// evolution must not exclude a possible logical-key dependency.
149fn files_may_overlap(
150    lhs: &FileHandle,
151    rhs: &FileHandle,
152    metadata: &RegionMetadata,
153    mapper: &PrimaryKeyRangeMapper,
154) -> bool {
155    let (lhs_start, lhs_end) = lhs.time_range();
156    let (rhs_start, rhs_end) = rhs.time_range();
157    if lhs_start.max(rhs_start) > lhs_end.min(rhs_end) {
158        return false;
159    }
160
161    // FileMeta does not record the PK schema/encoding. Adding a tag can change
162    // encoded keys without changing the logical keys of old rows.
163    // TODO: Use schema-aware PK bounds before relaxing this fallback. Appending
164    // defaults to dense keys can make raw ranges miss logical overlaps, breaking
165    // closure picking and risking premature tombstone removal in TWCS.
166    if metadata.schema_version != 0
167        || lhs.region_id() != metadata.region_id
168        || rhs.region_id() != metadata.region_id
169    {
170        return true;
171    }
172    match (lhs.primary_key_range(mapper), rhs.primary_key_range(mapper)) {
173        (Some(lhs), Some(rhs)) if lhs.0 <= lhs.1 && rhs.0 <= rhs.1 => {
174            primary_key_ranges_overlap(&lhs, &rhs)
175        }
176        _ => true,
177    }
178}
179
180#[cfg(test)]
181mod tests {
182    use std::collections::HashSet;
183    use std::sync::Arc;
184
185    use rand::{Rng, SeedableRng};
186    use store_api::storage::FileId;
187
188    use super::*;
189    use crate::compaction::test_util::{pk_range, primary_key_metadata_for_test};
190    use crate::sst::file::FileMeta;
191    use crate::test_util::memtable_util::metadata_for_test;
192    use crate::test_util::new_noop_file_purger;
193
194    fn file(start: i64, end: i64, pk: Option<(&str, &str)>) -> FileHandle {
195        FileHandle::new_with_primary_key_range(
196            FileMeta {
197                file_id: FileId::random(),
198                time_range: (
199                    Timestamp::new_millisecond(start),
200                    Timestamp::new_millisecond(end),
201                ),
202                ..Default::default()
203            },
204            new_noop_file_purger(),
205            pk.and_then(|(start, end)| pk_range(start.as_bytes(), end.as_bytes())),
206        )
207    }
208
209    #[rstest::rstest]
210    #[case(10, Some(("b", "c")), 0, false, true)]
211    #[case(11, Some(("b", "c")), 0, false, false)]
212    #[case(10, Some(("c", "d")), 0, false, false)]
213    #[case(10, None, 0, false, true)]
214    #[case(10, Some(("z", "a")), 0, false, true)]
215    #[case(10, Some(("c", "d")), 1, false, true)]
216    #[case(10, Some(("c", "d")), 0, true, true)]
217    #[case(11, None, 1, true, false)]
218    fn test_conservative_overlap(
219        #[case] start: i64,
220        #[case] pk: Option<(&str, &str)>,
221        #[case] schema_version: u64,
222        #[case] foreign: bool,
223        #[case] expected: bool,
224    ) {
225        let mut metadata = (*primary_key_metadata_for_test()).clone();
226        metadata.region_id = 0.into();
227        metadata.schema_version = schema_version;
228        let metadata = Arc::new(metadata);
229        let ranges = PrimaryKeyRangeMapper::new(metadata.clone());
230        let lhs = file(0, 10, Some(("a", "b")));
231        let mut rhs = file(start, start + 10, pk);
232        if foreign {
233            let mut meta = rhs.meta_ref().clone();
234            meta.region_id = 1.into();
235            rhs = FileHandle::new_with_primary_key_range(
236                meta,
237                new_noop_file_purger(),
238                rhs.raw_primary_key_range(),
239            );
240        }
241        assert_eq!(expected, files_may_overlap(&lhs, &rhs, &metadata, &ranges));
242        assert_eq!(expected, files_may_overlap(&rhs, &lhs, &metadata, &ranges));
243    }
244
245    #[test]
246    fn test_index_removal_distinguishes_region_owners() {
247        let metadata = metadata_for_test();
248        let lhs = file(0, 10, None);
249        let mut meta = lhs.meta_ref().clone();
250        meta.region_id = 1.into();
251        let rhs = FileHandle::new(meta, new_noop_file_purger());
252        let mut index = FileOverlapIndex::new(vec![lhs.clone(), rhs.clone()], &metadata);
253        index.remove(lhs.file_id());
254        let overlaps = index.drain_overlaps(&lhs);
255        assert_eq!(1, overlaps.len());
256        assert_eq!(rhs.file_id(), overlaps[0].file_id());
257        assert!(index.drain_overlaps(&lhs).is_empty());
258    }
259
260    #[test]
261    fn test_index_matches_linear_scan_after_removals() {
262        let metadata = primary_key_metadata_for_test();
263        let ranges = PrimaryKeyRangeMapper::new(metadata.clone());
264        let mut rng = rand::rngs::StdRng::seed_from_u64(9146);
265        for count in [0, 1, 7, 32, 127] {
266            let files: Vec<_> = (0..count)
267                .map(|_| {
268                    let start = rng.random_range(-50..50);
269                    let end = start + rng.random_range(0..40);
270                    let pk = match rng.random_range(0..4) {
271                        0 => None,
272                        1 => Some(("a", "b")),
273                        2 => Some(("b", "c")),
274                        _ => Some(("z", "z")),
275                    };
276                    file(start, end, pk)
277                })
278                .collect();
279            let mut index = FileOverlapIndex::new(files.clone(), &metadata);
280            let mut removed = HashSet::new();
281            for f in files.iter().step_by(3) {
282                index.remove(f.file_id());
283                removed.insert(f.file_id());
284            }
285            for _ in 0..32 {
286                let start = rng.random_range(-60..60);
287                let query = file(start, start + rng.random_range(0..25), Some(("b", "c")));
288                let expected: Vec<_> = files
289                    .iter()
290                    .filter(|f| {
291                        !removed.contains(&f.file_id())
292                            && files_may_overlap(&query, f, &metadata, &ranges)
293                    })
294                    .map(FileHandle::file_id)
295                    .collect();
296                let actual: Vec<_> = index
297                    .drain_overlaps(&query)
298                    .iter()
299                    .map(FileHandle::file_id)
300                    .collect();
301                assert_eq!(expected, actual, "snapshot size {count}");
302                removed.extend(expected);
303            }
304        }
305    }
306}