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