Skip to main content

mito2/sst/
version.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
15//! SST version.
16use std::collections::HashMap;
17use std::fmt;
18use std::sync::Arc;
19
20use common_time::{TimeToLive, Timestamp};
21use store_api::metadata::RegionMetadataRef;
22use store_api::storage::{FileId, RegionId};
23
24use crate::sst::file::{FileHandle, FileMeta, FileTimeRange, Level, MAX_LEVEL};
25use crate::sst::file_purger::FilePurgerRef;
26use crate::sst::primary_key::PrimaryKeyRangeMapper;
27
28/// A version of all SSTs in a region.
29#[derive(Debug, Clone)]
30pub(crate) struct SstVersion {
31    /// SST metadata organized by levels.
32    levels: Arc<LevelMetaArray>,
33    primary_key_mapper: Arc<PrimaryKeyRangeMapper>,
34}
35
36pub(crate) type SstVersionRef = Arc<SstVersion>;
37
38impl SstVersion {
39    /// Returns a new [SstVersion].
40    pub(crate) fn new(metadata: RegionMetadataRef) -> SstVersion {
41        SstVersion {
42            levels: Arc::new(new_level_meta_vec()),
43            primary_key_mapper: Arc::new(PrimaryKeyRangeMapper::new(metadata)),
44        }
45    }
46
47    /// Changes the target schema without copying the SST list or rebinding handles.
48    pub(crate) fn set_metadata(&mut self, metadata: RegionMetadataRef) {
49        self.primary_key_mapper = Arc::new(self.primary_key_mapper.with_metadata(metadata));
50    }
51
52    /// Shares the target schema and encoded defaults with comparisons.
53    pub(crate) fn primary_key_mapper(&self) -> Arc<PrimaryKeyRangeMapper> {
54        self.primary_key_mapper.clone()
55    }
56
57    /// Returns a slice to metadatas of all levels.
58    pub(crate) fn levels(&self) -> &[LevelMeta] {
59        self.levels.as_ref()
60    }
61
62    /// Returns the current handle matching the selected file's identity in its immutable level.
63    pub(crate) fn file_for_compaction(&self, selected: &FileHandle) -> Option<&FileHandle> {
64        let current = self
65            .levels
66            .get(selected.level() as usize)?
67            .files
68            .get(&selected.file_id().file_id())?;
69        (current.file_id() == selected.file_id()).then_some(current)
70    }
71
72    /// Add files to the version. If a file with the same `file_id` already exists,
73    /// it will be overwritten with the new file.
74    ///
75    /// # Panics
76    /// Panics if level of [FileMeta] is greater than [MAX_LEVEL].
77    pub(crate) fn add_files(
78        &mut self,
79        file_purger: FilePurgerRef,
80        files_to_add: impl Iterator<Item = FileMeta>,
81    ) {
82        let levels = Arc::make_mut(&mut self.levels);
83        for file in files_to_add {
84            let level = file.level;
85            let new_index_version = file.index_version;
86            // If the file already exists, then we should only replace the handle when the index is outdated.
87            levels[level as usize]
88                .files
89                .entry(file.file_id)
90                .and_modify(|f| {
91                    if *f.meta_ref() == file || f.meta_ref().is_index_up_to_date(&file) {
92                        // same file meta or current file handle's index is up-to-date, skip adding
93                        if f.index_id().version > new_index_version {
94                            // what does it mean for us to see older index version?
95                            common_telemetry::warn!(
96                                "Adding file with older index version, existing: {:?}, new: {:?}, ignoring new file",
97                                f.meta_ref(),
98                                file
99                            );
100                        }
101                    } else {
102                        // include case like old file have no index or index is outdated
103                        *f = FileHandle::new(file.clone(), file_purger.clone());
104                    }
105                })
106                .or_insert_with(|| {
107                    FileHandle::new(file.clone(), file_purger.clone())
108                });
109        }
110    }
111
112    /// Remove files from the version.
113    ///
114    /// # Panics
115    /// Panics if level of [FileMeta] is greater than [MAX_LEVEL].
116    pub(crate) fn remove_files(&mut self, files_to_remove: impl Iterator<Item = FileMeta>) {
117        let levels = Arc::make_mut(&mut self.levels);
118        for file in files_to_remove {
119            let level = file.level;
120            if let Some(handle) = levels[level as usize].files.remove(&file.file_id) {
121                handle.mark_deleted();
122            }
123        }
124    }
125
126    /// Marks all SSTs in this version as deleted.
127    pub(crate) fn mark_all_deleted(&self) {
128        for level_meta in self.levels.iter() {
129            for file_handle in level_meta.files.values() {
130                file_handle.mark_deleted();
131            }
132        }
133    }
134
135    /// Returns the number of rows in SST files owned by `region_id`.
136    ///
137    /// Rows from SST files referenced from other regions, for example after
138    /// repartition, are not counted.
139    /// For historical reasons, the result is not precise for old SST files.
140    pub(crate) fn owned_num_rows(&self, region_id: RegionId) -> u64 {
141        self.levels
142            .iter()
143            .map(|level_meta| {
144                level_meta
145                    .files
146                    .values()
147                    .filter(|file_handle| file_handle.region_id() == region_id)
148                    .map(|file_handle| {
149                        let meta = file_handle.meta_ref();
150                        meta.num_rows
151                    })
152                    .sum::<u64>()
153            })
154            .sum()
155    }
156
157    /// Returns the number of SST files owned by `region_id`.
158    pub(crate) fn owned_num_files(&self, region_id: RegionId) -> u64 {
159        self.levels
160            .iter()
161            .map(|level_meta| {
162                level_meta
163                    .files
164                    .values()
165                    .filter(|file_handle| file_handle.region_id() == region_id)
166                    .count() as u64
167            })
168            .sum()
169    }
170
171    /// Returns the time range covered by every file in this version, including
172    /// files referenced from other regions after a repartition.
173    pub(crate) fn time_range(&self) -> Option<FileTimeRange> {
174        self.levels
175            .iter()
176            .flat_map(|level_meta| level_meta.files.values())
177            .map(|file_handle| file_handle.time_range())
178            .reduce(|(min_a, max_a), (min_b, max_b)| (min_a.min(min_b), max_a.max(max_b)))
179    }
180
181    /// Returns the space occupied by SST data files owned by `region_id`.
182    pub(crate) fn owned_sst_usage(&self, region_id: RegionId) -> u64 {
183        self.levels
184            .iter()
185            .map(|level_meta| {
186                level_meta
187                    .files
188                    .values()
189                    .filter(|file_handle| file_handle.region_id() == region_id)
190                    .map(|file_handle| {
191                        let meta = file_handle.meta_ref();
192                        meta.file_size
193                    })
194                    .sum::<u64>()
195            })
196            .sum()
197    }
198
199    /// Returns the space occupied by SST index files owned by `region_id`.
200    pub(crate) fn owned_index_usage(&self, region_id: RegionId) -> u64 {
201        self.levels
202            .iter()
203            .map(|level_meta| {
204                level_meta
205                    .files
206                    .values()
207                    .filter(|file_handle| file_handle.region_id() == region_id)
208                    .map(|file_handle| {
209                        let meta = file_handle.meta_ref();
210                        meta.index_file_size
211                    })
212                    .sum::<u64>()
213            })
214            .sum()
215    }
216}
217
218// We only has fixed number of level, so we use array to hold elements. This implementation
219// detail of LevelMetaArray should not be exposed to users of [LevelMetas].
220type LevelMetaArray = [LevelMeta; MAX_LEVEL as usize];
221
222/// Metadata of files in the same SST level.
223#[derive(Clone)]
224pub struct LevelMeta {
225    /// Level number.
226    pub level: Level,
227    /// Handles of SSTs in this level.
228    pub files: HashMap<FileId, FileHandle>,
229}
230
231impl LevelMeta {
232    /// Returns an empty meta of specific `level`.
233    pub(crate) fn new(level: Level) -> LevelMeta {
234        LevelMeta {
235            level,
236            files: HashMap::new(),
237        }
238    }
239
240    /// Returns expired SSTs from current level.
241    pub fn get_expired_files(&self, now: &Timestamp, ttl: &TimeToLive) -> Vec<FileHandle> {
242        self.files
243            .values()
244            .filter(|v| {
245                let (_, end) = v.time_range();
246
247                match ttl.is_expired(&end, now) {
248                    Ok(expired) => expired,
249                    Err(e) => {
250                        common_telemetry::error!(e; "Failed to calculate region TTL expire time");
251                        false
252                    }
253                }
254            })
255            .cloned()
256            .collect()
257    }
258
259    pub fn files(&self) -> impl Iterator<Item = &FileHandle> {
260        self.files.values()
261    }
262}
263
264impl fmt::Debug for LevelMeta {
265    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
266        f.debug_struct("LevelMeta")
267            .field("level", &self.level)
268            .field("files", &self.files.keys())
269            .finish()
270    }
271}
272
273fn new_level_meta_vec() -> LevelMetaArray {
274    (0u8..MAX_LEVEL)
275        .map(LevelMeta::new)
276        .collect::<Vec<_>>()
277        .try_into()
278        .unwrap() // safety: LevelMetaArray is a fixed length array with length MAX_LEVEL
279}
280
281#[cfg(test)]
282mod tests {
283    use super::*;
284    use crate::test_util::new_noop_file_purger;
285
286    #[test]
287    fn time_range_spans_files_referenced_from_other_regions() {
288        let purger = new_noop_file_purger();
289        let owned = FileMeta {
290            file_id: FileId::random(),
291            region_id: RegionId::new(1, 1),
292            time_range: (
293                Timestamp::new_millisecond(200),
294                Timestamp::new_millisecond(300),
295            ),
296            ..Default::default()
297        };
298        let referenced = FileMeta {
299            file_id: FileId::random(),
300            region_id: RegionId::new(2, 1),
301            time_range: (
302                Timestamp::new_millisecond(50),
303                Timestamp::new_millisecond(100),
304            ),
305            ..Default::default()
306        };
307
308        let mut version = SstVersion::new(crate::test_util::memtable_util::metadata_for_test());
309        version.add_files(purger, [owned, referenced].into_iter());
310
311        assert_eq!(
312            version.time_range(),
313            Some((
314                Timestamp::new_millisecond(50),
315                Timestamp::new_millisecond(300)
316            ))
317        );
318    }
319
320    #[test]
321    fn time_range_compares_across_units() {
322        let purger = new_noop_file_purger();
323        let seconds = FileMeta {
324            file_id: FileId::random(),
325            time_range: (Timestamp::new_second(1), Timestamp::new_second(2)),
326            ..Default::default()
327        };
328        let millis = FileMeta {
329            file_id: FileId::random(),
330            time_range: (
331                Timestamp::new_millisecond(500),
332                Timestamp::new_millisecond(2500),
333            ),
334            ..Default::default()
335        };
336
337        let mut version = SstVersion::new(crate::test_util::memtable_util::metadata_for_test());
338        version.add_files(purger, [seconds, millis].into_iter());
339
340        assert_eq!(
341            version.time_range(),
342            Some((
343                Timestamp::new_millisecond(500),
344                Timestamp::new_millisecond(2500)
345            ))
346        );
347    }
348
349    #[test]
350    fn time_range_is_none_without_files() {
351        let version = SstVersion::new(crate::test_util::memtable_util::metadata_for_test());
352        assert_eq!(version.time_range(), None);
353    }
354
355    #[test]
356    fn test_add_files() {
357        let purger = new_noop_file_purger();
358
359        let files = (1..=3)
360            .map(|_| FileMeta {
361                file_id: FileId::random(),
362                ..Default::default()
363            })
364            .collect::<Vec<_>>();
365
366        let mut version = SstVersion::new(crate::test_util::memtable_util::metadata_for_test());
367        // files[1] is added multiple times, and that's ok.
368        version.add_files(purger.clone(), files[..=1].iter().cloned());
369        version.add_files(purger, files[1..].iter().cloned());
370
371        let added_files = &version.levels()[0].files;
372        assert_eq!(added_files.len(), 3);
373        files.iter().for_each(|f| {
374            assert!(added_files.contains_key(&f.file_id));
375        });
376    }
377
378    #[test]
379    fn test_file_for_compaction_uses_selected_level() {
380        let purger = new_noop_file_purger();
381        let file_id = FileId::random();
382        let selected = FileHandle::new(
383            FileMeta {
384                file_id,
385                level: 1,
386                ..Default::default()
387            },
388            purger.clone(),
389        );
390        let mut version = SstVersion::new(crate::test_util::memtable_util::metadata_for_test());
391        version.add_files(
392            purger,
393            [
394                FileMeta {
395                    file_id,
396                    level: 0,
397                    ..Default::default()
398                },
399                selected.meta_ref().clone(),
400            ]
401            .into_iter(),
402        );
403
404        let current = version.file_for_compaction(&selected).unwrap();
405        assert_eq!(selected.file_id(), current.file_id());
406        assert_eq!(selected.level(), current.level());
407    }
408
409    #[test]
410    fn test_usage_only_counts_owned_files() {
411        let purger = new_noop_file_purger();
412        let region_id = RegionId::new(1, 1);
413        let other_region_id = RegionId::new(1, 2);
414
415        let files = [
416            FileMeta {
417                region_id,
418                file_id: FileId::random(),
419                file_size: 100,
420                index_file_size: 10,
421                num_rows: 1,
422                ..Default::default()
423            },
424            FileMeta {
425                region_id,
426                file_id: FileId::random(),
427                file_size: 200,
428                index_file_size: 20,
429                num_rows: 2,
430                ..Default::default()
431            },
432            FileMeta {
433                region_id: other_region_id,
434                file_id: FileId::random(),
435                file_size: 300,
436                index_file_size: 30,
437                num_rows: 3,
438                ..Default::default()
439            },
440        ];
441
442        let mut version = SstVersion::new(crate::test_util::memtable_util::metadata_for_test());
443        version.add_files(purger, files.iter().cloned());
444
445        assert_eq!(3, version.owned_num_rows(region_id));
446        assert_eq!(2, version.owned_num_files(region_id));
447        assert_eq!(300, version.owned_sst_usage(region_id));
448        assert_eq!(30, version.owned_index_usage(region_id));
449        assert_eq!(3, version.owned_num_rows(other_region_id));
450        assert_eq!(1, version.owned_num_files(other_region_id));
451        assert_eq!(300, version.owned_sst_usage(other_region_id));
452        assert_eq!(30, version.owned_index_usage(other_region_id));
453    }
454}