Skip to main content

mito2/
compaction.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
15mod buckets;
16pub mod compactor;
17mod json2;
18mod last_non_null;
19pub mod memory_manager;
20mod overlap;
21pub mod picker;
22mod reader;
23pub mod run;
24mod scheduler;
25mod task;
26#[cfg(test)]
27mod test_util;
28mod twcs;
29mod window;
30
31use std::collections::HashMap;
32
33use common_meta::key::SchemaMetadataManagerRef;
34use common_telemetry::{debug, error};
35use common_time::TimeToLive;
36use common_time::range::TimestampRange;
37pub(crate) use json2::{
38    Json2RewritePlans, collect_json2_rewrite_plans, rewrite_json2_batch, rewrite_json2_schema,
39};
40pub use scheduler::CompactionRequest;
41pub(crate) use scheduler::{
42    CompactionExecution, CompactionPickFinished, CompactionScheduler, CompactionTransition,
43};
44use serde::{Deserialize, Serialize};
45use snafu::ResultExt;
46use store_api::mito_engine_options::{TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM, TWCS_TRIGGER_FILE_NUM};
47use store_api::storage::RegionId;
48
49use crate::error::{GetSchemaMetadataSnafu, Result, TimeoutSnafu};
50use crate::sst::file::{FileHandle, FileMeta, Level};
51
52/// Finds compaction options and TTL together with a single metadata fetch to reduce RTT.
53async fn find_dynamic_options(
54    region_id: RegionId,
55    region_options: &crate::region::options::RegionOptions,
56    schema_metadata_manager: &SchemaMetadataManagerRef,
57) -> Result<(crate::region::options::CompactionOptions, TimeToLive)> {
58    let table_id = region_id.table_id();
59    if let (true, Some(ttl)) = (region_options.compaction_override, region_options.ttl) {
60        debug!(
61            "Use region options directly for table {}: compaction={:?}, ttl={:?}",
62            table_id, region_options.compaction, region_options.ttl
63        );
64        return Ok((region_options.compaction.clone(), ttl));
65    }
66
67    let db_options = tokio::time::timeout(
68        crate::config::FETCH_OPTION_TIMEOUT,
69        schema_metadata_manager.get_schema_options_by_table_id(table_id),
70    )
71    .await
72    .context(TimeoutSnafu)?
73    .context(GetSchemaMetadataSnafu)?;
74
75    let ttl = if let Some(ttl) = region_options.ttl {
76        debug!(
77            "Use region TTL directly for table {}: ttl={:?}",
78            table_id, region_options.ttl
79        );
80        ttl
81    } else {
82        db_options
83            .as_ref()
84            .and_then(|options| options.ttl)
85            .unwrap_or_default()
86            .into()
87    };
88
89    let compaction = if !region_options.compaction_override {
90        if let Some(schema_opts) = db_options {
91            let mut map: HashMap<String, String> = schema_opts
92                .extra_options
93                .iter()
94                .filter_map(|(k, v)| {
95                    if k.starts_with("compaction.") {
96                        Some((k.clone(), v.clone()))
97                    } else {
98                        None
99                    }
100                })
101                .collect();
102            // Historical metadata may contain both aliases; prefer the canonical key.
103            if map.contains_key(TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM) {
104                map.remove(TWCS_TRIGGER_FILE_NUM);
105            }
106            if map.is_empty() {
107                region_options.compaction.clone()
108            } else {
109                crate::region::options::RegionOptions::try_from_options(region_id, &map)
110                    .map(|o| o.compaction)
111                    .unwrap_or_else(|e| {
112                        error!(e; "Failed to create RegionOptions from map");
113                        region_options.compaction.clone()
114                    })
115            }
116        } else {
117            debug!(
118                "DB options is None for table {}, use region compaction: compaction={:?}",
119                table_id, region_options.compaction
120            );
121            region_options.compaction.clone()
122        }
123    } else {
124        debug!(
125            "No schema options for table {}, use region compaction: compaction={:?}",
126            table_id, region_options.compaction
127        );
128        region_options.compaction.clone()
129    };
130
131    debug!(
132        "Resolved dynamic options for table {}: compaction={:?}, ttl={:?}",
133        table_id, compaction, ttl
134    );
135    Ok((compaction, ttl))
136}
137
138#[derive(Debug, Clone)]
139pub struct CompactionOutput {
140    /// Compaction output file level.
141    pub output_level: Level,
142    /// Compaction input files.
143    pub inputs: Vec<FileHandle>,
144    /// Whether to remove deletion markers.
145    pub filter_deleted: bool,
146    /// Compaction output time range. Only windowed compaction specifies output time range.
147    pub output_time_range: Option<TimestampRange>,
148}
149
150/// SerializedCompactionOutput is a serialized version of [CompactionOutput] by replacing [FileHandle] with [FileMeta].
151#[derive(Debug, Clone, Serialize, Deserialize)]
152pub struct SerializedCompactionOutput {
153    output_level: Level,
154    inputs: Vec<FileMeta>,
155    filter_deleted: bool,
156    output_time_range: Option<TimestampRange>,
157}