Skip to main content

mito2/compaction/
picker.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::fmt::Debug;
16use std::sync::Arc;
17
18use api::v1::region::compact_request;
19use common_time::range::TimestampRange;
20use common_time::{TimeToLive, Timestamp};
21use serde::{Deserialize, Serialize};
22
23use crate::compaction::compactor::CompactionRegion;
24use crate::compaction::last_non_null::LastNonNullPicker;
25use crate::compaction::twcs::TwcsPicker;
26use crate::compaction::window::WindowedCompactionPicker;
27use crate::compaction::{CompactionOutput, SerializedCompactionOutput};
28use crate::error::Result;
29use crate::region::options::{CompactionOptions, MergeMode, RegionOptions};
30use crate::sst::file::{FileHandle, FileMeta};
31use crate::sst::file_purger::FilePurger;
32use crate::sst::version::LevelMeta;
33
34#[async_trait::async_trait]
35pub(crate) trait CompactionTask: Debug + Send + Sync + 'static {
36    async fn run(&mut self);
37}
38
39/// Picker picks input SST files for compaction.
40/// Different compaction strategy may implement different pickers.
41#[async_trait::async_trait]
42pub trait Picker: Debug + Send + Sync + 'static {
43    /// Picks input SST files for compaction.
44    async fn pick(&self, compaction_region: &CompactionRegion) -> Result<Option<PickerOutput>>;
45}
46
47/// PickerOutput is the output of a [`Picker`].
48/// It contains the outputs of the compaction and the expired SST files.
49#[derive(Default, Clone, Debug)]
50pub struct PickerOutput {
51    pub outputs: Vec<CompactionOutput>,
52    pub expired_ssts: Vec<FileHandle>,
53    pub time_window_size: i64,
54    /// Max single output file size in bytes.
55    pub max_file_size: Option<usize>,
56}
57
58/// SerializedPickerOutput is a serialized version of PickerOutput by replacing [CompactionOutput] and [FileHandle] with [SerializedCompactionOutput] and [FileMeta].
59#[derive(Default, Clone, Debug, Serialize, Deserialize)]
60pub struct SerializedPickerOutput {
61    pub outputs: Vec<SerializedCompactionOutput>,
62    pub expired_ssts: Vec<FileMeta>,
63    pub time_window_size: i64,
64    pub max_file_size: Option<usize>,
65}
66
67impl From<&PickerOutput> for SerializedPickerOutput {
68    fn from(input: &PickerOutput) -> Self {
69        let outputs = input
70            .outputs
71            .iter()
72            .map(|output| SerializedCompactionOutput {
73                output_level: output.output_level,
74                inputs: output.inputs.iter().map(|s| s.meta_ref().clone()).collect(),
75                filter_deleted: output.filter_deleted,
76                output_time_range: output.output_time_range,
77            })
78            .collect();
79        let expired_ssts = input
80            .expired_ssts
81            .iter()
82            .map(|s| s.meta_ref().clone())
83            .collect();
84        Self {
85            outputs,
86            expired_ssts,
87            time_window_size: input.time_window_size,
88            max_file_size: input.max_file_size,
89        }
90    }
91}
92
93impl PickerOutput {
94    /// Converts a [SerializedPickerOutput] to a [PickerOutput]. File statistics
95    /// retain their original encoding; comparisons supply their own schema context.
96    pub fn from_serialized(
97        input: SerializedPickerOutput,
98        file_purger: Arc<dyn FilePurger>,
99    ) -> Self {
100        let outputs = input
101            .outputs
102            .into_iter()
103            .map(|output| CompactionOutput {
104                output_level: output.output_level,
105                inputs: output
106                    .inputs
107                    .into_iter()
108                    .map(|file_meta| FileHandle::new(file_meta, file_purger.clone()))
109                    .collect(),
110                filter_deleted: output.filter_deleted,
111                output_time_range: output.output_time_range,
112            })
113            .collect();
114
115        let expired_ssts = input
116            .expired_ssts
117            .into_iter()
118            .map(|file_meta| FileHandle::new(file_meta, file_purger.clone()))
119            .collect();
120
121        Self {
122            outputs,
123            expired_ssts,
124            time_window_size: input.time_window_size,
125            max_file_size: input.max_file_size,
126        }
127    }
128}
129
130/// Creates a picker for the request and the region's compaction and merge modes.
131pub fn new_picker(
132    compact_request_options: &compact_request::Options,
133    region_options: &RegionOptions,
134    max_background_tasks: Option<usize>,
135    time_range: Option<TimestampRange>,
136) -> Arc<dyn Picker> {
137    if let compact_request::Options::StrictWindow(window) = compact_request_options {
138        let window = if window.window_seconds == 0 {
139            None
140        } else {
141            Some(window.window_seconds)
142        };
143        // Strict-window outputs already include every overlapping file segment
144        // within each disjoint, half-open output range, including LastNonNull.
145        Arc::new(WindowedCompactionPicker::new(window).with_time_range(time_range)) as Arc<_>
146    } else {
147        let picker = match &region_options.compaction {
148            CompactionOptions::Twcs(twcs_opts) => TwcsPicker {
149                trigger_file_num: twcs_opts.active_window_trigger_file_num,
150                active_window_l1_merge_trigger: twcs_opts.active_window_l1_merge_trigger,
151                inactive_window_trigger_file_num: twcs_opts.inactive_window_trigger_file_num,
152                inactive_window_l1_merge_trigger: twcs_opts.inactive_window_l1_merge_trigger,
153                time_window_seconds: twcs_opts.time_window_seconds(),
154                max_output_file_size: twcs_opts.max_output_file_size.map(|r| r.as_bytes()),
155                append_mode: region_options.append_mode,
156                max_background_tasks,
157                time_range,
158            },
159        };
160        if region_options.merge_mode() == MergeMode::LastNonNull {
161            // LastNonNull correctness spans windows and levels, so wrap the seed
162            // picker instead of trusting TWCS's independent windows.
163            Arc::new(LastNonNullPicker::new(picker))
164        } else {
165            Arc::new(picker)
166        }
167    }
168}
169
170/// Finds all expired SSTs across levels.
171pub(super) fn get_expired_ssts(
172    levels: &[LevelMeta],
173    ttl: Option<TimeToLive>,
174    now: Timestamp,
175) -> Vec<FileHandle> {
176    let Some(ttl) = ttl else {
177        return vec![];
178    };
179
180    levels
181        .iter()
182        .flat_map(|l| l.get_expired_files(&now, &ttl).into_iter())
183        .collect()
184}
185
186#[cfg(test)]
187mod tests {
188    use store_api::storage::FileId;
189
190    use super::*;
191    use crate::compaction::test_util::new_file_handle;
192    use crate::test_util::new_noop_file_purger;
193
194    #[test]
195    fn test_picker_output_serialization() {
196        let inputs_file_handle = vec![
197            new_file_handle(FileId::random(), 0, 999, 0),
198            new_file_handle(FileId::random(), 0, 999, 0),
199            new_file_handle(FileId::random(), 0, 999, 0),
200        ];
201        let expired_ssts_file_handle = vec![
202            new_file_handle(FileId::random(), 0, 999, 0),
203            new_file_handle(FileId::random(), 0, 999, 0),
204        ];
205
206        let picker_output = PickerOutput {
207            outputs: vec![
208                CompactionOutput {
209                    output_level: 0,
210                    inputs: inputs_file_handle.clone(),
211                    filter_deleted: false,
212                    output_time_range: None,
213                },
214                CompactionOutput {
215                    output_level: 0,
216                    inputs: inputs_file_handle.clone(),
217                    filter_deleted: false,
218                    output_time_range: None,
219                },
220            ],
221            expired_ssts: expired_ssts_file_handle.clone(),
222            time_window_size: 1000,
223            max_file_size: None,
224        };
225
226        let picker_output_str =
227            serde_json::to_string(&SerializedPickerOutput::from(&picker_output)).unwrap();
228        let serialized_picker_output: SerializedPickerOutput =
229            serde_json::from_str(&picker_output_str).unwrap();
230        let picker_output_from_serialized =
231            PickerOutput::from_serialized(serialized_picker_output, new_noop_file_purger());
232
233        picker_output
234            .expired_ssts
235            .iter()
236            .zip(picker_output_from_serialized.expired_ssts.iter())
237            .for_each(|(expected, actual)| {
238                assert_eq!(expected.meta_ref(), actual.meta_ref());
239            });
240
241        picker_output
242            .outputs
243            .iter()
244            .zip(picker_output_from_serialized.outputs.iter())
245            .for_each(|(expected, actual)| {
246                assert_eq!(expected.output_level, actual.output_level);
247                expected
248                    .inputs
249                    .iter()
250                    .zip(actual.inputs.iter())
251                    .for_each(|(expected, actual)| {
252                        assert_eq!(expected.meta_ref(), actual.meta_ref());
253                    });
254                assert_eq!(expected.filter_deleted, actual.filter_deleted);
255                assert_eq!(expected.output_time_range, actual.output_time_range);
256            });
257    }
258}