From f524a0b5b4deeeed39c03dd0eea9a573341709d6 Mon Sep 17 00:00:00 2001 From: "Lei, HUANG" <6406592+v0y4g3r@users.noreply.github.com> Date: Wed, 29 Jul 2026 21:33:10 +0800 Subject: [PATCH] feat: support time range in manual compaction (#8669) * feat: support time range in manual compaction Signed-off-by: Lei, HUANG * fix: reject overflowing compaction range alignment Signed-off-by: Lei, HUANG * fix: preserve range across compaction continuations Signed-off-by: Lei, HUANG * perf: use graph traversal for compaction windows Signed-off-by: Lei, HUANG * chore: bump proto to commit on main Signed-off-by: Lei, HUANG * docs: explain compaction window dependency closure Signed-off-by: Lei, HUANG --------- Signed-off-by: Lei, HUANG --- Cargo.lock | 2 +- Cargo.toml | 2 +- .../function/src/admin/flush_compact_table.rs | 321 ++++++++++++------ src/mito2/src/compaction.rs | 267 +++++++++++++-- src/mito2/src/compaction/picker.rs | 5 +- src/mito2/src/compaction/twcs.rs | 144 +++++++- src/mito2/src/compaction/window.rs | 135 +++++++- src/mito2/src/worker/handle_compaction.rs | 24 +- src/operator/src/request.rs | 87 ++++- src/store-api/src/region_request.rs | 64 +++- src/table/src/requests.rs | 2 + .../function/admin/flush_compact_table.result | 24 ++ .../function/admin/flush_compact_table.sql | 12 + 13 files changed, 925 insertions(+), 164 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 1dc4e68bd7..fe566fd88a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -5991,7 +5991,7 @@ dependencies = [ [[package]] name = "greptime-proto" version = "0.1.0" -source = "git+https://github.com/GreptimeTeam/greptime-proto.git?rev=a937a645546b377a8ba88b43dd82a294e63d99ac#a937a645546b377a8ba88b43dd82a294e63d99ac" +source = "git+https://github.com/GreptimeTeam/greptime-proto.git?rev=5adb8a1637abe87bcbb455d32136eee5d538cc49#5adb8a1637abe87bcbb455d32136eee5d538cc49" dependencies = [ "prost 0.14.1", "prost-types 0.14.1", diff --git a/Cargo.toml b/Cargo.toml index ee45c64e6d..049d453e2c 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -158,7 +158,7 @@ fs2 = "0.4" fst = "0.4.7" futures = "0.3" futures-util = "0.3" -greptime-proto = { git = "https://github.com/GreptimeTeam/greptime-proto.git", rev = "a937a645546b377a8ba88b43dd82a294e63d99ac" } +greptime-proto = { git = "https://github.com/GreptimeTeam/greptime-proto.git", rev = "5adb8a1637abe87bcbb455d32136eee5d538cc49" } hex = "0.4" http = "1" humantime = "2.1" diff --git a/src/common/function/src/admin/flush_compact_table.rs b/src/common/function/src/admin/flush_compact_table.rs index 20b95e4379..4bec7b2e7d 100644 --- a/src/common/function/src/admin/flush_compact_table.rs +++ b/src/common/function/src/admin/flush_compact_table.rs @@ -23,6 +23,8 @@ use common_query::error::{ UnsupportedInputDataTypeSnafu, }; use common_telemetry::info; +use common_time::range::TimestampRange; +use common_time::{Timestamp, Timezone}; use datafusion_expr::{Signature, Volatility}; use datatypes::prelude::*; use session::context::QueryContextRef; @@ -121,6 +123,7 @@ fn compact_signature() -> Signature { /// - `[, , ]`: provides both type and type-specific options. /// - For `twcs`, it accepts `parallelism=[N]` where N is an unsigned 32 bits number /// - For `swcs`, it accepts two numeric parameter: `parallelism` and `window`. +/// - Both types accept `start_time` and `end_time` to constrain compaction windows. fn parse_compact_request( params: &[ValueRef<'_>], query_ctx: &QueryContextRef, @@ -129,26 +132,29 @@ fn parse_compact_request( !params.is_empty() && params.len() <= 3, InvalidFuncArgsSnafu { err_msg: format!( - "The length of the args is not correct, expect 1-4, have: {}", + "The length of the args is not correct, expect 1-3, have: {}", params.len() ), } ); - let (table_name, compact_type, parallelism) = match params { + let timezone = query_ctx.timezone(); + let (table_name, compact_type, parallelism, time_range) = match params { // 1. Only table name, strategy defaults to twcs and default parallelism. [ValueRef::String(table_name)] => ( table_name, compact_request::Options::Regular(Default::default()), DEFAULT_COMPACTION_PARALLELISM, + None, ), // 2. Both table name and strategy are provided. [ ValueRef::String(table_name), ValueRef::String(compact_ty_str), ] => { - let (compact_type, parallelism) = parse_compact_options(compact_ty_str, None)?; - (table_name, compact_type, parallelism) + let (compact_type, parallelism, time_range) = + parse_compact_options(compact_ty_str, None, &timezone)?; + (table_name, compact_type, parallelism, time_range) } // 3. Table name, strategy and strategy specific options [ @@ -156,9 +162,9 @@ fn parse_compact_request( ValueRef::String(compact_ty_str), ValueRef::String(options_str), ] => { - let (compact_type, parallelism) = - parse_compact_options(compact_ty_str, Some(options_str))?; - (table_name, compact_type, parallelism) + let (compact_type, parallelism, time_range) = + parse_compact_options(compact_ty_str, Some(options_str), &timezone)?; + (table_name, compact_type, parallelism, time_range) } _ => { return UnsupportedInputDataTypeSnafu { @@ -179,6 +185,7 @@ fn parse_compact_request( table_name, compact_options: compact_type, parallelism, + time_range, }) } @@ -187,118 +194,115 @@ fn parse_compact_request( fn parse_compact_options( type_str: &str, option: Option<&str>, -) -> Result<(compact_request::Options, u32)> { - if type_str.eq_ignore_ascii_case(COMPACT_TYPE_STRICT_WINDOW) - | type_str.eq_ignore_ascii_case(COMPACT_TYPE_STRICT_WINDOW_SHORT) - { - let Some(option_str) = option else { - return Ok(( - compact_request::Options::StrictWindow(StrictWindow { window_seconds: 0 }), - DEFAULT_COMPACTION_PARALLELISM, - )); + timezone: &Timezone, +) -> Result<(compact_request::Options, u32, Option)> { + let strict_window = type_str.eq_ignore_ascii_case(COMPACT_TYPE_STRICT_WINDOW) + || type_str.eq_ignore_ascii_case(COMPACT_TYPE_STRICT_WINDOW_SHORT); + let Some(option_str) = option else { + let options = if strict_window { + compact_request::Options::StrictWindow(StrictWindow { window_seconds: 0 }) + } else { + compact_request::Options::Regular(Default::default()) }; + return Ok((options, DEFAULT_COMPACTION_PARALLELISM, None)); + }; - // For compatibility, accepts single number as window size. - if let Ok(window_seconds) = i64::from_str(option_str) { - return Ok(( - compact_request::Options::StrictWindow(StrictWindow { window_seconds }), - DEFAULT_COMPACTION_PARALLELISM, - )); - }; - - // Parse keyword arguments in forms: `key1=value1,key2=value2` - let mut window_seconds = 0i64; - let mut parallelism = DEFAULT_COMPACTION_PARALLELISM; - - let pairs: Vec<&str> = option_str.split(',').collect(); - for pair in pairs { - let kv: Vec<&str> = pair.trim().split('=').collect(); - if kv.len() != 2 { - return InvalidFuncArgsSnafu { - err_msg: format!("Invalid key-value pair: {}", pair.trim()), - } - .fail(); - } - - let key = kv[0].trim(); - let value = kv[1].trim(); - - match key { - "window" | "window_seconds" => { - window_seconds = i64::from_str(value).map_err(|_| { - InvalidFuncArgsSnafu { - err_msg: format!("Invalid value for window: {}", value), - } - .build() - })?; - } - "parallelism" => { - parallelism = value.parse::().map_err(|_| { - InvalidFuncArgsSnafu { - err_msg: format!("Invalid value for parallelism: {}", value), - } - .build() - })?; - } - _ => { - return InvalidFuncArgsSnafu { - err_msg: format!("Unknown parameter: {}", key), - } - .fail(); - } - } - } - - Ok(( + // For compatibility, strict-window compaction accepts a single number as window size. + if strict_window && let Ok(window_seconds) = i64::from_str(option_str) { + return Ok(( compact_request::Options::StrictWindow(StrictWindow { window_seconds }), - parallelism, - )) - } else { - // TWCS strategy - let Some(option_str) = option else { - return Ok(( - compact_request::Options::Regular(Default::default()), - DEFAULT_COMPACTION_PARALLELISM, - )); - }; + DEFAULT_COMPACTION_PARALLELISM, + None, + )); + } - let mut parallelism = DEFAULT_COMPACTION_PARALLELISM; - let pairs: Vec<&str> = option_str.split(',').collect(); - for pair in pairs { - let kv: Vec<&str> = pair.trim().split('=').collect(); - if kv.len() != 2 { + let mut window_seconds = 0i64; + let mut parallelism = DEFAULT_COMPACTION_PARALLELISM; + let mut start_time = None; + let mut end_time = None; + + for pair in option_str.split(',') { + let Some((key, value)) = pair.trim().split_once('=') else { + return InvalidFuncArgsSnafu { + err_msg: format!("Invalid key-value pair: {}", pair.trim()), + } + .fail(); + }; + let key = key.trim(); + let value = value.trim(); + + match key { + "window" | "window_seconds" if strict_window => { + window_seconds = i64::from_str(value).map_err(|_| { + InvalidFuncArgsSnafu { + err_msg: format!("Invalid value for window: {}", value), + } + .build() + })?; + } + "parallelism" => { + parallelism = value.parse::().map_err(|_| { + InvalidFuncArgsSnafu { + err_msg: format!("Invalid value for parallelism: {}", value), + } + .build() + })?; + } + "start_time" => { + start_time = Some(Timestamp::from_str(value, Some(timezone)).map_err(|_| { + InvalidFuncArgsSnafu { + err_msg: format!("Invalid value for start_time: {}", value), + } + .build() + })?); + } + "end_time" => { + end_time = Some(Timestamp::from_str(value, Some(timezone)).map_err(|_| { + InvalidFuncArgsSnafu { + err_msg: format!("Invalid value for end_time: {}", value), + } + .build() + })?); + } + _ => { return InvalidFuncArgsSnafu { - err_msg: format!("Invalid key-value pair: {}", pair.trim()), + err_msg: format!("Unknown parameter: {}", key), } .fail(); } - - let key = kv[0].trim(); - let value = kv[1].trim(); - - match key { - "parallelism" => { - parallelism = value.parse::().map_err(|_| { - InvalidFuncArgsSnafu { - err_msg: format!("Invalid value for parallelism: {}", value), - } - .build() - })?; - } - _ => { - return InvalidFuncArgsSnafu { - err_msg: format!("Unknown parameter: {}", key), - } - .fail(); - } - } } - - Ok(( - compact_request::Options::Regular(Default::default()), - parallelism, - )) } + + let time_range = match (start_time, end_time) { + (None, None) => None, + (Some(start), Some(end)) if start < end => { + Some(TimestampRange::new(start, end).ok_or_else(|| { + InvalidFuncArgsSnafu { + err_msg: "invalid compaction time range".to_string(), + } + .build() + })?) + } + (Some(_), Some(_)) => { + return InvalidFuncArgsSnafu { + err_msg: "start_time must be earlier than end_time".to_string(), + } + .fail(); + } + _ => { + return InvalidFuncArgsSnafu { + err_msg: "start_time and end_time must be specified together".to_string(), + } + .fail(); + } + }; + + let options = if strict_window { + compact_request::Options::StrictWindow(StrictWindow { window_seconds }) + } else { + compact_request::Options::Regular(Default::default()) + }; + Ok((options, parallelism, time_range)) } #[cfg(test)] @@ -309,6 +313,8 @@ mod tests { use arrow::array::StringArray; use arrow::datatypes::{DataType, Field}; use common_catalog::consts::{DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME}; + use common_time::Timestamp; + use common_time::range::TimestampRange; use datafusion_expr::ColumnarValue; use session::context::QueryContext; @@ -409,6 +415,90 @@ mod tests { } } + #[test] + fn test_parse_compact_time_range() { + let params = [ + "table", + "regular", + "start_time=2026-01-01T00:00:00Z,end_time=2026-02-01T00:00:00Z", + ] + .into_iter() + .map(ValueRef::String) + .collect::>(); + + let request = parse_compact_request(¶ms, &QueryContext::arc()).unwrap(); + assert_eq!( + Some( + TimestampRange::new( + Timestamp::from_str_utc("2026-01-01T00:00:00Z").unwrap(), + Timestamp::from_str_utc("2026-02-01T00:00:00Z").unwrap(), + ) + .unwrap() + ), + request.time_range + ); + + let query_ctx = QueryContext::arc(); + query_ctx.set_timezone(Timezone::from_tz_string("Asia/Shanghai").unwrap()); + let params = [ + "table", + "regular", + "start_time=2026-01-01T00:00:00,end_time=2026-02-01T00:00:00", + ] + .into_iter() + .map(ValueRef::String) + .collect::>(); + let request = parse_compact_request(¶ms, &query_ctx).unwrap(); + assert_eq!( + Some( + TimestampRange::new( + Timestamp::from_str_utc("2025-12-31T16:00:00Z").unwrap(), + Timestamp::from_str_utc("2026-01-31T16:00:00Z").unwrap(), + ) + .unwrap() + ), + request.time_range + ); + } + + #[test] + fn test_parse_strict_window_compact_time_range() { + let params = [ + "table", + "strict_window", + "window=3600,parallelism=2,start_time=2026-01-01T00:00:00Z,end_time=2026-02-01T00:00:00Z", + ] + .into_iter() + .map(ValueRef::String) + .collect::>(); + + let request = parse_compact_request(¶ms, &QueryContext::arc()).unwrap(); + assert_eq!( + Options::StrictWindow(StrictWindow { + window_seconds: 3600, + }), + request.compact_options + ); + assert_eq!(2, request.parallelism); + assert!(request.time_range.is_some()); + } + + #[test] + fn test_parse_compact_time_range_requires_valid_bounds() { + for options in [ + "start_time=2026-01-01T00:00:00Z", + "end_time=2026-02-01T00:00:00Z", + "start_time=2026-02-01T00:00:00Z,end_time=2026-01-01T00:00:00Z", + "start_time=2026-01-01T00:00:00Z,end_time=2026-01-01T00:00:00Z", + ] { + let params = ["table", "regular", options] + .into_iter() + .map(ValueRef::String) + .collect::>(); + assert!(parse_compact_request(¶ms, &QueryContext::arc()).is_err()); + } + } + #[test] fn test_parse_compact_params() { check_parse_compact_params(&[ @@ -420,6 +510,7 @@ mod tests { table_name: "table".to_string(), compact_options: Options::Regular(Default::default()), parallelism: 1, + time_range: None, }, ), ( @@ -430,6 +521,7 @@ mod tests { table_name: "table".to_string(), compact_options: Options::Regular(Default::default()), parallelism: 1, + time_range: None, }, ), ( @@ -443,6 +535,7 @@ mod tests { table_name: "table".to_string(), compact_options: Options::Regular(Default::default()), parallelism: 1, + time_range: None, }, ), ( @@ -453,6 +546,7 @@ mod tests { table_name: "table".to_string(), compact_options: Options::Regular(Default::default()), parallelism: 1, + time_range: None, }, ), ( @@ -463,6 +557,7 @@ mod tests { table_name: "table".to_string(), compact_options: Options::StrictWindow(StrictWindow { window_seconds: 0 }), parallelism: 1, + time_range: None, }, ), ( @@ -475,6 +570,7 @@ mod tests { window_seconds: 3600, }), parallelism: 1, + time_range: None, }, ), ( @@ -487,6 +583,7 @@ mod tests { window_seconds: 120, }), parallelism: 1, + time_range: None, }, ), // Test with parallelism parameter @@ -498,6 +595,7 @@ mod tests { table_name: "table".to_string(), compact_options: Options::Regular(Default::default()), parallelism: 4, + time_range: None, }, ), ( @@ -510,6 +608,7 @@ mod tests { window_seconds: 3600, }), parallelism: 2, + time_range: None, }, ), ( @@ -522,6 +621,7 @@ mod tests { window_seconds: 3600, }), parallelism: 1, + time_range: None, }, ), ( @@ -534,6 +634,7 @@ mod tests { window_seconds: 7200, }), parallelism: 1, + time_range: None, }, ), ( @@ -546,6 +647,7 @@ mod tests { window_seconds: 1800, }), parallelism: 1, + time_range: None, }, ), ( @@ -556,6 +658,7 @@ mod tests { table_name: "table".to_string(), compact_options: Options::Regular(Default::default()), parallelism: 8, + time_range: None, }, ), ]); diff --git a/src/mito2/src/compaction.rs b/src/mito2/src/compaction.rs index 30a3512eca..41139a252e 100644 --- a/src/mito2/src/compaction.rs +++ b/src/mito2/src/compaction.rs @@ -30,7 +30,6 @@ use std::sync::{Arc, Mutex}; use std::time::Instant; use api::v1::region::compact_request; -use api::v1::region::compact_request::Options; use common_base::Plugins; use common_base::cancellation::CancellationHandle; use common_memory_manager::OnExhaustedPolicy; @@ -192,6 +191,13 @@ pub(crate) struct CompactionScheduler { next_plan_id: u64, } +fn requires_pending_compaction_slot( + options: &compact_request::Options, + time_range: Option, +) -> bool { + matches!(options, compact_request::Options::StrictWindow(_)) || time_range.is_some() +} + impl CompactionScheduler { #[allow(clippy::too_many_arguments)] pub(crate) fn new( @@ -241,6 +247,33 @@ impl CompactionScheduler { manifest_ctx: &ManifestContextRef, schema_metadata_manager: SchemaMetadataManagerRef, max_parallelism: usize, + ) -> Result { + self.schedule_compaction_with_time_range( + region_id, + compact_options, + version_control, + access_layer, + waiter, + manifest_ctx, + schema_metadata_manager, + max_parallelism, + None, + ) + } + + /// Schedules a compaction constrained by an optional time range. + #[allow(clippy::too_many_arguments)] + pub(crate) fn schedule_compaction_with_time_range( + &mut self, + region_id: RegionId, + compact_options: compact_request::Options, + version_control: &VersionControlRef, + access_layer: &AccessLayerRef, + waiter: OptionOutputTx, + manifest_ctx: &ManifestContextRef, + schema_metadata_manager: SchemaMetadataManagerRef, + max_parallelism: usize, + time_range: Option, ) -> Result { // skip compaction if region is in staging state let current_state = manifest_ctx.current_state(); @@ -267,22 +300,20 @@ impl CompactionScheduler { return Ok(false); } - match compact_options { - Options::Regular(_) => { - status.merge_regular_trigger(waiter); - } - options @ Options::StrictWindow(_) => { - // Incoming compaction request is manually triggered. - status.set_pending_request(PendingCompaction { - options, - waiter, - max_parallelism, - }); - info!( - "Region {} is compacting, manually compaction will be re-scheduled.", - region_id - ); - } + if requires_pending_compaction_slot(&compact_options, time_range) { + // Incoming compaction request is manually triggered. + status.set_pending_request(PendingCompaction { + options: compact_options, + waiter, + max_parallelism, + time_range, + }); + info!( + "Region {} is compacting, manually compaction will be re-scheduled.", + region_id + ); + } else { + status.merge_regular_trigger(waiter); } return Ok(false); } @@ -300,10 +331,10 @@ impl CompactionScheduler { max_parallelism, ); let plan_id = Self::next_plan_id(&mut self.next_plan_id); - status.start_picking(plan_id); + status.start_picking_with_time_range(plan_id, time_range); status.merge_waiter(waiter); self.region_status.insert(region_id, status); - self.dispatch_compaction_planning(plan_id, request, compact_options); + self.dispatch_compaction_planning(plan_id, request, compact_options, time_range); self.listener.on_compaction_scheduled(region_id); Ok(true) } @@ -331,6 +362,7 @@ impl CompactionScheduler { options, waiter, max_parallelism, + time_range, } = pending_request; let request = status.new_compaction_request( @@ -347,8 +379,8 @@ impl CompactionScheduler { // borrow stays alive; nothing could have removed the status since it // was fetched above. let plan_id = Self::next_plan_id(&mut self.next_plan_id); - status.start_picking(plan_id); - self.dispatch_compaction_planning(plan_id, request, options); + status.start_picking_with_time_range(plan_id, time_range); + self.dispatch_compaction_planning(plan_id, request, options, time_range); debug!( "Successfully scheduled manual compaction planning for region id: {}", region_id @@ -413,6 +445,7 @@ impl CompactionScheduler { manifest_ctx, schema_metadata_manager, Some(active), + None, ); return Vec::new(); } @@ -488,11 +521,13 @@ impl CompactionScheduler { return true; } + let time_range = status.time_range; self.schedule_next_compaction_with_active( region_id, manifest_ctx, schema_metadata_manager, None, + time_range, ) } @@ -502,6 +537,7 @@ impl CompactionScheduler { manifest_ctx: &ManifestContextRef, schema_metadata_manager: SchemaMetadataManagerRef, active: Option, + time_range: Option, ) -> bool { let Some(status) = self.region_status.get_mut(®ion_id) else { return false; @@ -520,11 +556,12 @@ impl CompactionScheduler { // borrow stays alive; nothing could have removed the status since it // was fetched above. let plan_id = Self::next_plan_id(&mut self.next_plan_id); - status.start_regular_picking(plan_id, active); + status.start_regular_picking(plan_id, active, time_range); self.dispatch_compaction_planning( plan_id, request, compact_request::Options::Regular(Default::default()), + time_range, ); debug!( "Successfully scheduled next compaction planning for region id: {}", @@ -664,14 +701,20 @@ impl CompactionScheduler { plan_id: u64, request: CompactionRequest, options: compact_request::Options, + time_range: Option, ) { let plugins = self.plugins.clone(); let max_background_compactions = self.engine_config.max_background_compactions; common_runtime::spawn_compact(async move { let region_id = request.region_id(); let request_sender = request.request_sender.clone(); - let planning = - Self::prepare_compaction(request, options, plugins, max_background_compactions); + let planning = Self::prepare_compaction( + request, + options, + plugins, + max_background_compactions, + time_range, + ); Self::notify_planning_result(region_id, plan_id, request_sender, planning).await; }); } @@ -729,6 +772,7 @@ impl CompactionScheduler { options: compact_request::Options, plugins: Plugins, max_background_compactions: usize, + time_range: Option, ) -> CompactionPlanningResult { let region_id = request.region_id(); let (dynamic_compaction_opts, ttl) = find_dynamic_options( @@ -750,6 +794,7 @@ impl CompactionScheduler { &dynamic_compaction_opts, request.current_version.options.append_mode, Some(max_background_compactions), + time_range, ); let region_id = request.region_id(); let CompactionRequest { @@ -954,6 +999,7 @@ impl CompactionScheduler { manifest_ctx, schema_metadata_manager, Some(active), + None, ); return Vec::new(); } @@ -1478,13 +1524,12 @@ struct CompactionStatus { // TODO: Remove idle statuses and make ActiveCompaction non-optional once chained // scheduling can recreate the status from region context. active: Option, + /// Optional range retained by automatic continuations of the current compaction. + time_range: Option, /// Pending compactions that are supposed to run as soon as current compaction task finished. /// - /// In production code, this can only hold a manual `Options::StrictWindow` compaction. - /// When a `Options::Regular` compaction request arrives while the region is already - /// compacting, it goes to `ActiveCompaction::regular_followup_waiters` or `waiters` - /// instead — it does not become a `pending_request`. See `schedule_compaction` for - /// the branching logic. + /// This holds strict-window requests and ranged regular requests. An unrestricted regular + /// request is instead merged into `ActiveCompaction::regular_followup_waiters` or `waiters`. pending_request: Option, /// Pending DDL requests that should run when compaction is done. /// @@ -1506,12 +1551,19 @@ impl CompactionStatus { version_control, access_layer, active: None, + time_range: None, pending_request: None, pending_ddl_requests: Vec::new(), } } + #[cfg(test)] fn start_picking(&mut self, plan_id: u64) { + self.start_picking_with_time_range(plan_id, None); + } + + fn start_picking_with_time_range(&mut self, plan_id: u64, time_range: Option) { + self.time_range = time_range; if let Some(active) = &mut self.active { active.start_picking(plan_id); } else { @@ -1519,7 +1571,13 @@ impl CompactionStatus { } } - fn start_regular_picking(&mut self, plan_id: u64, active: Option) { + fn start_regular_picking( + &mut self, + plan_id: u64, + active: Option, + time_range: Option, + ) { + self.time_range = time_range; self.active = Some(if let Some(mut active) = active { active.start_regular_picking(plan_id); active @@ -2031,13 +2089,14 @@ fn refresh_picker_output(output: PickerOutput, current: &SstVersion) -> Option

, } #[cfg(test)] @@ -2046,6 +2105,7 @@ mod tests { use std::time::Duration; use api::v1::region::StrictWindow; + use api::v1::region::compact_request::Options; use common_datasource::compression::CompressionType; use common_meta::key::schema_name::SchemaNameValue; use common_time::DatabaseTimeToLive; @@ -2064,6 +2124,28 @@ mod tests { use crate::test_util::scheduler_util::{SchedulerEnv, VecScheduler}; use crate::test_util::version_util::{VersionControlBuilder, apply_edit}; + #[test] + fn test_requires_pending_compaction_slot() { + let time_range = TimestampRange::new( + Timestamp::new_millisecond(1_000), + Timestamp::new_millisecond(2_000), + ) + .unwrap(); + + assert!(!requires_pending_compaction_slot( + &compact_request::Options::Regular(Default::default()), + None, + )); + assert!(requires_pending_compaction_slot( + &compact_request::Options::Regular(Default::default()), + Some(time_range), + )); + assert!(requires_pending_compaction_slot( + &compact_request::Options::StrictWindow(StrictWindow { window_seconds: 60 }), + None, + )); + } + struct FailingScheduler; struct FailingRemoteScheduler; @@ -3141,7 +3223,7 @@ mod tests { } #[tokio::test] - async fn test_manual_compaction_when_compaction_in_progress() { + async fn test_time_range_compaction_when_compaction_in_progress() { common_telemetry::init_default_ut_logging(); let job_scheduler = Arc::new(VecScheduler::default()); let env = SchedulerEnv::new().await.scheduler(job_scheduler.clone()); @@ -3225,25 +3307,37 @@ mod tests { .is_none() ); - // Schedule another manual compaction. + // Schedule another manual compaction with a time range. + let time_range = TimestampRange::new( + Timestamp::new_millisecond(0), + Timestamp::new_millisecond(end + 1), + ) + .unwrap(); let (tx, _rx) = oneshot::channel(); scheduler - .schedule_compaction( + .schedule_compaction_with_time_range( region_id, - compact_request::Options::StrictWindow(StrictWindow { window_seconds: 60 }), + compact_request::Options::Regular(Default::default()), &version_control, &env.access_layer, OptionOutputTx::new(Some(OutputTx::new(tx))), &manifest_ctx, schema_metadata_manager.clone(), 1, + Some(time_range), ) .unwrap(); assert_eq!(1, scheduler.region_status.len()); // Current job num should be 1 since compaction is in progress. assert_eq!(1, job_scheduler.num_jobs()); let status = scheduler.region_status.get(&builder.region_id()).unwrap(); - assert!(status.pending_request.is_some()); + assert_eq!( + Some(time_range), + status + .pending_request + .as_ref() + .and_then(|pending| pending.time_range) + ); // On compaction finished and schedule next compaction. scheduler @@ -3265,6 +3359,101 @@ mod tests { assert!(status.pending_request.is_none()); } + #[tokio::test] + async fn test_ranged_compaction_continuation_preserves_time_range() { + let job_scheduler = Arc::new(VecScheduler::default()); + let env = SchedulerEnv::new().await.scheduler(job_scheduler.clone()); + let (tx, mut rx) = mpsc::channel(4); + let mut scheduler = env.mock_compaction_scheduler(tx); + + let mut builder = VersionControlBuilder::new(); + for offset in [0, 10, 20, 30] { + builder.push_l0_file(offset, 1_000); + } + for offset in [0, 10, 20, 30] { + builder.push_l0_file(2 * 3_600_000 + offset, 2 * 3_600_000 + 1_000); + } + let region_id = builder.region_id(); + let version_control = Arc::new(builder.build()); + let manifest_ctx = env + .mock_manifest_context(version_control.current().version.metadata.clone()) + .await; + let (schema_metadata_manager, kv_backend) = mock_schema_metadata_manager(); + let mut schema_value = SchemaNameValue::default(); + schema_value + .extra_options + .insert("compaction.type".to_string(), "twcs".to_string()); + schema_value + .extra_options + .insert("compaction.twcs.time_window".to_string(), "1h".to_string()); + schema_metadata_manager + .register_region_table_info( + region_id.table_id(), + "t", + "c", + "s", + Some(schema_value), + kv_backend, + ) + .await; + let time_range = TimestampRange::new( + Timestamp::new_millisecond(0), + Timestamp::new_millisecond(3_600_000), + ) + .unwrap(); + + scheduler + .schedule_compaction_with_time_range( + region_id, + compact_request::Options::StrictWindow(StrictWindow { + window_seconds: 3_600, + }), + &version_control, + &env.access_layer, + OptionOutputTx::none(), + &manifest_ctx, + schema_metadata_manager.clone(), + 1, + Some(time_range), + ) + .unwrap(); + let first = recv_compaction_pick_finished(&mut rx).await; + assert!( + selected_files(&first) + .iter() + .all(|file| file.time_range().1 < Timestamp::new_millisecond(3_600_000)) + ); + scheduler + .handle_compaction_pick_finished(first, &manifest_ctx, schema_metadata_manager.clone()) + .await; + assert_eq!(1, job_scheduler.num_jobs()); + + scheduler + .on_compaction_finished(region_id, &manifest_ctx, schema_metadata_manager.clone()) + .await; + assert!(scheduler.schedule_next_compaction( + region_id, + &manifest_ctx, + schema_metadata_manager, + )); + + let continuation = recv_compaction_pick_finished(&mut rx).await; + match continuation.result { + CompactionPlanningResult::NoPlan => {} + CompactionPlanningResult::Prepared(prepared) => assert!( + prepared + .picker_output + .outputs + .iter() + .flat_map(|output| &output.inputs) + .all(|file| file.time_range().1 < Timestamp::new_millisecond(3_600_000)) + ), + CompactionPlanningResult::Error(err) => { + panic!("unexpected compaction planning error: {err}") + } + } + } + #[tokio::test] async fn test_compaction_bypass_in_staging_mode() { let env = SchedulerEnv::new().await; @@ -3378,6 +3567,7 @@ mod tests { options: compact_request::Options::StrictWindow(StrictWindow { window_seconds: 60 }), waiter: OptionOutputTx::from(first_manual_tx), max_parallelism: 1, + time_range: None, }); scheduler.region_status.insert(region_id, status); @@ -3629,6 +3819,7 @@ mod tests { options: compact_request::Options::StrictWindow(StrictWindow { window_seconds: 60 }), waiter: OptionOutputTx::from(manual_tx), max_parallelism: 1, + time_range: None, }); let (output_tx, _output_rx) = oneshot::channel(); @@ -3860,6 +4051,7 @@ mod tests { options: compact_request::Options::Regular(Default::default()), waiter: OptionOutputTx::from(manual_tx), max_parallelism: 1, + time_range: None, }); scheduler.region_status.insert(region_id, status); @@ -3987,6 +4179,7 @@ mod tests { options: compact_request::Options::Regular(Default::default()), waiter: OptionOutputTx::from(manual_tx), max_parallelism: 1, + time_range: None, }); scheduler.region_status.insert(region_id, status); diff --git a/src/mito2/src/compaction/picker.rs b/src/mito2/src/compaction/picker.rs index 7c5cccfb8c..403538606f 100644 --- a/src/mito2/src/compaction/picker.rs +++ b/src/mito2/src/compaction/picker.rs @@ -16,6 +16,7 @@ use std::fmt::Debug; use std::sync::Arc; use api::v1::region::compact_request; +use common_time::range::TimestampRange; use serde::{Deserialize, Serialize}; use crate::compaction::compactor::CompactionRegion; @@ -126,6 +127,7 @@ pub fn new_picker( compaction_options: &CompactionOptions, append_mode: bool, max_background_tasks: Option, + time_range: Option, ) -> Arc { if let compact_request::Options::StrictWindow(window) = compact_request_options { let window = if window.window_seconds == 0 { @@ -133,7 +135,7 @@ pub fn new_picker( } else { Some(window.window_seconds) }; - Arc::new(WindowedCompactionPicker::new(window)) as Arc<_> + Arc::new(WindowedCompactionPicker::new(window).with_time_range(time_range)) as Arc<_> } else { match compaction_options { CompactionOptions::Twcs(twcs_opts) => Arc::new(TwcsPicker { @@ -142,6 +144,7 @@ pub fn new_picker( max_output_file_size: twcs_opts.max_output_file_size.map(|r| r.as_bytes()), append_mode, max_background_tasks, + time_range, }) as Arc<_>, } } diff --git a/src/mito2/src/compaction/twcs.rs b/src/mito2/src/compaction/twcs.rs index fb95175655..9ec9fe647f 100644 --- a/src/mito2/src/compaction/twcs.rs +++ b/src/mito2/src/compaction/twcs.rs @@ -20,6 +20,7 @@ use std::num::NonZeroU64; use common_base::readable_size::ReadableSize; use common_telemetry::{debug, info}; use common_time::Timestamp; +use common_time::range::TimestampRange; use common_time::timestamp::TimeUnit; use common_time::timestamp_millis::BucketAligned; use rayon::prelude::*; @@ -55,15 +56,28 @@ pub struct TwcsPicker { pub append_mode: bool, /// Max background compaction tasks. pub max_background_tasks: Option, + /// Optional time range that constrains candidate compaction windows. + pub(crate) time_range: Option, } impl TwcsPicker { /// Builds compaction output from files. + #[cfg(test)] fn build_output( &self, region_id: RegionId, time_windows: &mut BTreeMap, active_window: Option, + ) -> Vec { + self.build_output_with_time_range(region_id, time_windows, active_window, None) + } + + fn build_output_with_time_range( + &self, + region_id: RegionId, + time_windows: &mut BTreeMap, + active_window: Option, + time_window_size: Option, ) -> Vec { let find_inputs = |files: &Window, windows: &BTreeMap| @@ -158,7 +172,18 @@ impl TwcsPicker { let mut output = vec![]; let windows = time_windows .values() - .filter(|w| !w.files.is_empty()) + .filter(|window| { + !window.files.is_empty() + && self.time_range.as_ref().is_none_or(|time_range| { + time_window_size.is_none_or(|time_window_size| { + time_window_intersects_range( + window.time_window, + time_window_size, + time_range, + ) + }) + }) + }) .collect::>(); let chunk_size = self.max_background_tasks.unwrap_or(windows.len()).max(1); 'chunks: for chunk in windows.chunks(chunk_size) { @@ -279,7 +304,12 @@ impl Picker for TwcsPicker { .filter(|file| !expired_file_ids.contains(&file.file_id())), time_window_size, ); - let outputs = self.build_output(region_id, &mut windows, active_window); + let outputs = self.build_output_with_time_range( + region_id, + &mut windows, + active_window, + Some(time_window_size), + ); if outputs.is_empty() && expired_ssts.is_empty() { return None; @@ -381,6 +411,39 @@ fn assign_to_windows<'a>( windows.into_iter().collect() } +fn time_window_intersects_range( + window_end: i64, + time_window_size: i64, + time_range: &TimestampRange, +) -> bool { + let first_window = match time_range.start() { + None => i64::MIN, + Some(start) => { + let Some(first_window) = start + .convert_to(TimeUnit::Second) + .and_then(|timestamp| timestamp.value().align_to_ceil_by_bucket(time_window_size)) + else { + return false; + }; + first_window + } + }; + let last_window = match time_range.end() { + None => i64::MAX, + Some(end) => { + let Some(last_window) = end + .convert_to_ceil(TimeUnit::Second) + .and_then(|timestamp| timestamp.value().checked_sub(1)) + .and_then(|timestamp| timestamp.align_to_ceil_by_bucket(time_window_size)) + else { + return false; + }; + last_window + } + }; + (first_window..=last_window).contains(&window_end) +} + fn window_has_overlap(this: &Window, windows: &BTreeMap) -> bool { windows .values() @@ -426,6 +489,7 @@ mod tests { use bytes::Bytes; use common_base::Plugins; + use common_time::range::TimestampRange; use store_api::storage::FileId; use super::*; @@ -490,6 +554,7 @@ mod tests { max_output_file_size: None, append_mode: false, max_background_tasks: None, + time_range: None, }; let compaction_region = compaction_region_with_expired_sst().await; @@ -887,6 +952,7 @@ mod tests { max_output_file_size: None, append_mode: false, max_background_tasks: None, + time_range: None, } .build_output(RegionId::from_u64(0), &mut windows, active_window); @@ -1028,6 +1094,7 @@ mod tests { max_output_file_size: None, append_mode: false, max_background_tasks: None, + time_range: None, } .build_output(RegionId::from_u64(0), &mut windows, active_window); @@ -1098,6 +1165,7 @@ mod tests { max_output_file_size: None, append_mode: false, max_background_tasks: None, + time_range: None, }; let active_window = find_latest_window_in_seconds(files.iter(), 3); @@ -1139,6 +1207,7 @@ mod tests { max_output_file_size: Some(1000), append_mode: true, max_background_tasks: None, + time_range: None, }; let active_window = find_latest_window_in_seconds(files.iter(), 3); @@ -1181,6 +1250,7 @@ mod tests { max_output_file_size: None, append_mode: false, max_background_tasks: Some(max_background_tasks), + time_range: None, }; let active_window = find_latest_window_in_seconds(files.iter(), 3); @@ -1201,6 +1271,7 @@ mod tests { max_output_file_size: None, append_mode: false, max_background_tasks: None, + time_range: None, }; let mut windows_no_limit = assign_to_windows(files.iter(), 3); @@ -1240,6 +1311,7 @@ mod tests { max_output_file_size: None, append_mode: false, max_background_tasks: Some(max_background_tasks), + time_range: None, }; let active_window = find_latest_window_in_seconds(files.iter(), 3); @@ -1285,6 +1357,7 @@ mod tests { max_output_file_size: None, append_mode: false, max_background_tasks: None, + time_range: None, }; let active_window = find_latest_window_in_seconds(files.iter(), 3); @@ -1325,6 +1398,7 @@ mod tests { max_output_file_size: None, append_mode: false, max_background_tasks: None, + time_range: None, }; let active_window = find_latest_window_in_seconds(files.iter(), 3); @@ -1334,5 +1408,71 @@ mod tests { assert_eq!(output[0].inputs.len(), 32); } + #[test] + fn test_filter_time_windows_by_time_range() { + let time_range = TimestampRange::new( + Timestamp::new_millisecond(1_200), + Timestamp::new_millisecond(1_800), + ) + .unwrap(); + + assert!(time_window_intersects_range(3, 3, &time_range)); + assert!(!time_window_intersects_range(9, 3, &time_range)); + + let boundary_range = + TimestampRange::new(Timestamp::new_second(0), Timestamp::new_second(3)).unwrap(); + assert!(time_window_intersects_range(0, 3, &boundary_range)); + assert!(time_window_intersects_range(3, 3, &boundary_range)); + assert!(!time_window_intersects_range(6, 3, &boundary_range)); + + let overflowing_range = TimestampRange::new( + Timestamp::new_second(i64::MAX - 1), + Timestamp::new_second(i64::MAX), + ) + .unwrap(); + assert!(!time_window_intersects_range(0, 4, &overflowing_range)); + } + + #[test] + fn test_time_range_filter_precedes_background_task_limit() { + let early_file_ids = [FileId::random(), FileId::random()]; + let selected_file_ids = [FileId::random(), FileId::random()]; + let files = [ + new_file_handle_with_sequence(early_file_ids[0], 1_000, 1_999, 0, 1), + new_file_handle_with_sequence(early_file_ids[1], 1_000, 1_999, 0, 2), + new_file_handle_with_sequence(selected_file_ids[0], 7_000, 7_999, 0, 3), + new_file_handle_with_sequence(selected_file_ids[1], 7_000, 7_999, 0, 4), + ]; + let mut windows = assign_to_windows(files.iter(), 3); + let picker = TwcsPicker { + trigger_file_num: 2, + time_window_seconds: Some(3), + max_output_file_size: None, + append_mode: false, + max_background_tasks: Some(1), + time_range: TimestampRange::new( + Timestamp::new_millisecond(7_200), + Timestamp::new_millisecond(7_800), + ), + }; + + let output = picker.build_output_with_time_range( + RegionId::from_u64(123), + &mut windows, + Some(9), + Some(3), + ); + + assert_eq!(1, output.len()); + assert_eq!( + selected_file_ids.into_iter().collect::>(), + output[0] + .inputs + .iter() + .map(|file| file.file_id().file_id()) + .collect::>() + ); + } + // TODO(hl): TTL tester that checks if get_expired_ssts function works as expected. } diff --git a/src/mito2/src/compaction/window.rs b/src/mito2/src/compaction/window.rs index f6d9cfae71..af7a35412c 100644 --- a/src/mito2/src/compaction/window.rs +++ b/src/mito2/src/compaction/window.rs @@ -12,7 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. -use std::collections::{BTreeMap, HashSet}; +use std::collections::{BTreeMap, HashMap, HashSet, VecDeque}; use std::fmt::Debug; use common_telemetry::info; @@ -34,15 +34,23 @@ use crate::sst::file::FileHandle; #[derive(Debug)] pub struct WindowedCompactionPicker { compaction_time_window_seconds: Option, + time_range: Option, } impl WindowedCompactionPicker { pub fn new(window_seconds: Option) -> Self { Self { compaction_time_window_seconds: window_seconds, + time_range: None, } } + /// Sets the time range used to select compaction windows. + pub(crate) fn with_time_range(mut self, time_range: Option) -> Self { + self.time_range = time_range; + self + } + // Computes compaction time window. First we respect user specified parameter, then // use persisted window. If persist window is not present, we check the time window // provided while creating table. If all of those are absent, we infer the window @@ -101,6 +109,7 @@ impl WindowedCompactionPicker { .flat_map(|level| level.files.values()) .filter(|file| !expired_file_ids.contains(&file.file_id())), ); + let windows = filter_time_windows(windows, self.time_range); (build_output(windows), expired_ssts, time_window) } @@ -123,6 +132,66 @@ impl Picker for WindowedCompactionPicker { } } +/// Keeps windows that overlap the requested range and their transitive dependencies. +/// +/// [`assign_files_to_time_windows`] adds an SST to every time window that the SST covers. If a +/// selected window contains such a cross-window SST, compaction will remove that input SST after +/// rewriting it. Keeping only the directly selected window would therefore omit the SST's rows in +/// the other windows. We must include every window covered by the SST, then repeat the process for +/// other cross-window SSTs in those windows, until the complete dependency closure is selected. +fn filter_time_windows( + mut windows: BTreeMap)>, + time_range: Option, +) -> BTreeMap)> { + let Some(time_range) = time_range else { + return windows; + }; + + let mut selected_windows = windows + .iter() + .filter_map(|(lower_bound, (upper_bound, _))| { + let window_start = Timestamp::new_second(*lower_bound); + let window_end = Timestamp::new_second(*upper_bound); + let starts_before_range_end = time_range + .end() + .is_none_or(|range_end| window_start < range_end); + let ends_after_range_start = time_range + .start() + .is_none_or(|range_start| range_start < window_end); + (starts_before_range_end && ends_after_range_start).then_some(*lower_bound) + }) + .collect::>(); + + let mut file_windows = HashMap::new(); + for (lower_bound, (_, files)) in &windows { + for file in files { + file_windows + .entry(file.file_id()) + .or_insert_with(Vec::new) + .push(*lower_bound); + } + } + + let mut pending_windows = selected_windows.iter().copied().collect::>(); + let mut visited_files = HashSet::new(); + while let Some(lower_bound) = pending_windows.pop_front() { + let (_, files) = &windows[&lower_bound]; + for file in files { + if !visited_files.insert(file.file_id()) { + continue; + } + for dependent_window in &file_windows[&file.file_id()] { + if selected_windows.insert(*dependent_window) { + pending_windows.push_back(*dependent_window); + } + } + } + } + + windows.retain(|lower_bound, _| selected_windows.contains(lower_bound)); + windows +} + fn build_output(windows: BTreeMap)>) -> Vec { let mut outputs = Vec::with_capacity(windows.len()); for (lower_bound, (upper_bound, files)) in windows { @@ -352,6 +421,70 @@ mod tests { ); } + #[test] + fn test_pick_time_range_expands_for_cross_window_files() { + let time_range = TimestampRange::new( + Timestamp::new_millisecond(HOUR / 2), + Timestamp::new_millisecond(HOUR * 3 / 4), + ) + .unwrap(); + let picker = + WindowedCompactionPicker::new(Some(HOUR / 1000)).with_time_range(Some(time_range)); + let files = vec![ + (FileId::random(), 0, 2 * HOUR - 1, 0), + (FileId::random(), HOUR, HOUR * 3 - 1, 0), + (FileId::random(), 4 * HOUR, 5 * HOUR - 1, 0), + ]; + let version = build_version(&files, None); + + let (outputs, _, _) = picker.pick_inner( + RegionId::new(0, 0), + &version, + Timestamp::new_millisecond(6 * HOUR), + ); + + assert_eq!(3, outputs.len()); + assert_eq!( + Some(TimestampRange::new( + Timestamp::new_millisecond(0), + Timestamp::new_millisecond(HOUR), + )), + outputs.first().map(|output| output.output_time_range) + ); + assert_eq!( + Some(TimestampRange::new( + Timestamp::new_millisecond(2 * HOUR), + Timestamp::new_millisecond(3 * HOUR), + )), + outputs.last().map(|output| output.output_time_range) + ); + } + + #[test] + fn test_pick_time_range_expands_long_dependency_chain() { + const CHAIN_LEN: i64 = 128; + + let time_range = TimestampRange::new( + Timestamp::new_millisecond(0), + Timestamp::new_millisecond(HOUR / 2), + ) + .unwrap(); + let picker = + WindowedCompactionPicker::new(Some(HOUR / 1000)).with_time_range(Some(time_range)); + let files = (0..CHAIN_LEN) + .map(|window| (FileId::random(), window * HOUR, (window + 2) * HOUR - 1, 0)) + .collect::>(); + let version = build_version(&files, None); + + let (outputs, _, _) = picker.pick_inner( + RegionId::new(0, 0), + &version, + Timestamp::new_millisecond((CHAIN_LEN + 2) * HOUR), + ); + + assert_eq!(CHAIN_LEN as usize + 1, outputs.len()); + } + #[test] fn test_assign_compacting_files_to_windows() { let picker = WindowedCompactionPicker::new(Some(HOUR / 1000)); diff --git a/src/mito2/src/worker/handle_compaction.rs b/src/mito2/src/worker/handle_compaction.rs index 9748afd3f9..7189503ff6 100644 --- a/src/mito2/src/worker/handle_compaction.rs +++ b/src/mito2/src/worker/handle_compaction.rs @@ -67,16 +67,20 @@ impl RegionWorkerLoop { }; COMPACTION_REQUEST_COUNT.inc(); let parallelism = req.parallelism.unwrap_or(1) as usize; - if let Err(e) = self.compaction_scheduler.schedule_compaction( - region.region_id, - req.options, - ®ion.version_control, - ®ion.access_layer, - sender, - ®ion.manifest_ctx, - self.schema_metadata_manager.clone(), - parallelism, - ) { + if let Err(e) = self + .compaction_scheduler + .schedule_compaction_with_time_range( + region.region_id, + req.options, + ®ion.version_control, + ®ion.access_layer, + sender, + ®ion.manifest_ctx, + self.schema_metadata_manager.clone(), + parallelism, + req.time_range, + ) + { error!(e; "Failed to schedule compaction task for region: {}", region_id); } else { info!( diff --git a/src/operator/src/request.rs b/src/operator/src/request.rs index ccfcfe8830..aea395d8ee 100644 --- a/src/operator/src/request.rs +++ b/src/operator/src/request.rs @@ -14,14 +14,19 @@ use std::sync::Arc; +use api::helper::to_pb_time_unit; use api::v1::region::region_request::Body as RegionRequestBody; -use api::v1::region::{BuildIndexRequest, CompactRequest, FlushRequest, RegionRequestHeader}; +use api::v1::region::{ + BuildIndexRequest, CompactRequest, CompactionTimeRange, FlushRequest, RegionRequestHeader, +}; use catalog::CatalogManagerRef; use common_catalog::build_db_string; use common_meta::node_manager::{AffectedRows, NodeManagerRef}; use common_meta::peer::Peer; use common_telemetry::tracing_context::TracingContext; use common_telemetry::{debug, error, info}; +use common_time::range::TimestampRange; +use common_time::timestamp::TimeUnit as TimestampUnit; use futures_util::future; use partition::cache::PhysicalPartitionInfo; use partition::manager::PartitionRuleManagerRef; @@ -145,6 +150,10 @@ impl Requester { .await? .partitions; + let time_range = request + .time_range + .map(to_pb_compaction_time_range) + .transpose()?; let requests = partitions .iter() .map(|partition| { @@ -152,6 +161,7 @@ impl Requester { region_id: partition.id.into(), parallelism: request.parallelism, options: Some(request.compact_options), + time_range, }) }) .collect(); @@ -190,6 +200,7 @@ impl Requester { region_id: region_id.into(), parallelism: 1, options: None, // todo(hl): maybe also support parameters in region compaction. + time_range: None, }); info!("Handle region manual compaction request: {region_id}"); @@ -197,6 +208,43 @@ impl Requester { } } +fn to_pb_compaction_time_range(range: TimestampRange) -> Result { + let (Some(start), Some(end)) = (*range.start(), *range.end()) else { + return crate::error::InvalidTimestampRangeSnafu { + start: format!("{:?}", range.start()), + end: format!("{:?}", range.end()), + } + .fail(); + }; + let original_start = start; + let original_end = end; + let start = start.convert_to(TimestampUnit::Second).with_context(|| { + crate::error::InvalidTimestampRangeSnafu { + start: format!("{original_start:?}"), + end: format!("{original_end:?}"), + } + })?; + let end = end + .convert_to_ceil(TimestampUnit::Second) + .with_context(|| crate::error::InvalidTimestampRangeSnafu { + start: format!("{original_start:?}"), + end: format!("{original_end:?}"), + })?; + ensure!( + start < end, + crate::error::InvalidTimestampRangeSnafu { + start: format!("{original_start:?}"), + end: format!("{original_end:?}"), + } + ); + + Ok(CompactionTimeRange { + start: start.value(), + end: end.value(), + time_unit: to_pb_time_unit(TimestampUnit::Second) as i32, + }) +} + impl Requester { async fn do_request( &self, @@ -280,3 +328,40 @@ impl Requester { }) } } + +#[cfg(test)] +mod tests { + use api::v1::TimeUnit; + use common_time::Timestamp; + use common_time::range::TimestampRange; + + use super::*; + + #[test] + fn test_to_pb_compaction_time_range_normalizes_mixed_units() { + let range = TimestampRange::new( + Timestamp::new_millisecond(1_500), + Timestamp::new_microsecond(2_500_000), + ) + .unwrap(); + + let pb_range = to_pb_compaction_time_range(range).unwrap(); + assert_eq!(1, pb_range.start); + assert_eq!(3, pb_range.end); + assert_eq!(TimeUnit::Second as i32, pb_range.time_unit); + } + + #[test] + fn test_to_pb_compaction_time_range() { + let range = TimestampRange::new( + Timestamp::new_microsecond(1_000), + Timestamp::new_microsecond(2_000), + ) + .unwrap(); + + let pb_range = to_pb_compaction_time_range(range).unwrap(); + assert_eq!(0, pb_range.start); + assert_eq!(1, pb_range.end); + assert_eq!(TimeUnit::Second as i32, pb_range.time_unit); + } +} diff --git a/src/store-api/src/region_request.rs b/src/store-api/src/region_request.rs index d420ed66a3..1f799200bb 100644 --- a/src/store-api/src/region_request.rs +++ b/src/store-api/src/region_request.rs @@ -16,7 +16,7 @@ use std::collections::HashMap; use std::fmt::{self, Display}; use std::time::Duration; -use api::helper::{ColumnDataTypeWrapper, from_pb_time_ranges}; +use api::helper::{ColumnDataTypeWrapper, from_pb_time_ranges, from_pb_time_unit}; use api::v1::add_column_location::LocationType; use api::v1::column_def::{ as_fulltext_option_analyzer, as_fulltext_option_backend, as_skipping_index_type, @@ -36,6 +36,7 @@ pub use common_base::AffectedRows; use common_base::readable_size::ReadableSize; use common_grpc::flight::FlightDecoder; use common_recordbatch::DfRecordBatch; +use common_time::range::TimestampRange; use common_time::{TimeToLive, Timestamp}; use datatypes::prelude::ConcreteDataType; use datatypes::schema::{FulltextOptions, SkippingIndexOptions}; @@ -391,11 +392,41 @@ fn make_region_compact(compact: CompactRequest) -> Result= {}", + range.start, range.end + ), + } + ); + TimestampRange::new(start, end).with_context(|| InvalidRegionRequestSnafu { + region_id, + err: "invalid compaction time range".to_string(), + }) + }) + .transpose()?; Ok(vec![( region_id, RegionRequest::Compact(RegionCompactRequest { options, parallelism, + time_range, }), )]) } @@ -1569,6 +1600,7 @@ pub struct RegionFlushRequest { pub struct RegionCompactRequest { pub options: compact_request::Options, pub parallelism: Option, + pub time_range: Option, } impl Default for RegionCompactRequest { @@ -1577,6 +1609,7 @@ impl Default for RegionCompactRequest { // Default to regular compaction. options: compact_request::Options::Regular(Default::default()), parallelism: None, + time_range: None, } } } @@ -1726,12 +1759,41 @@ mod tests { use api::v1::region::RegionColumnDef; use api::v1::{ColumnDataType, ColumnDef}; + use common_time::range::TimestampRange; use datatypes::prelude::ConcreteDataType; use datatypes::schema::{ColumnSchema, FulltextAnalyzer, FulltextBackend}; use super::*; use crate::metadata::RegionMetadataBuilder; + #[test] + fn test_make_region_compact_with_time_range() { + let requests = make_region_compact(CompactRequest { + region_id: 42, + time_range: Some(api::v1::region::CompactionTimeRange { + start: 1_000, + end: 2_000, + time_unit: api::v1::TimeUnit::Microsecond as i32, + }), + ..Default::default() + }) + .unwrap(); + + let RegionRequest::Compact(request) = &requests[0].1 else { + unreachable!(); + }; + assert_eq!( + Some( + TimestampRange::new( + Timestamp::new_microsecond(1_000), + Timestamp::new_microsecond(2_000), + ) + .unwrap() + ), + request.time_range + ); + } + #[test] fn test_from_proto_location() { let proto_location = v1::AddColumnLocation { diff --git a/src/table/src/requests.rs b/src/table/src/requests.rs index 12331bf8bb..afaa5da7b5 100644 --- a/src/table/src/requests.rs +++ b/src/table/src/requests.rs @@ -439,6 +439,7 @@ pub struct CompactTableRequest { pub table_name: String, pub compact_options: compact_request::Options, pub parallelism: u32, + pub time_range: Option, } impl Default for CompactTableRequest { @@ -449,6 +450,7 @@ impl Default for CompactTableRequest { table_name: Default::default(), compact_options: compact_request::Options::Regular(Default::default()), parallelism: 1, + time_range: None, } } } diff --git a/tests/cases/standalone/common/function/admin/flush_compact_table.result b/tests/cases/standalone/common/function/admin/flush_compact_table.result index 0c3fdbc0de..8ea0377f5a 100644 --- a/tests/cases/standalone/common/function/admin/flush_compact_table.result +++ b/tests/cases/standalone/common/function/admin/flush_compact_table.result @@ -35,6 +35,30 @@ ADMIN COMPACT_TABLE('test'); | 0 | +-----------------------------+ +ADMIN COMPACT_TABLE( + 'test', + 'regular', + 'start_time=1970-01-01T00:00:00Z,end_time=1970-01-01T01:00:00Z' +); + ++---------------------------------------------------------------------------------------------------------+ +| ADMIN COMPACT_TABLE('test', 'regular', 'start_time=1970-01-01T00:00:00Z,end_time=1970-01-01T01:00:00Z') | ++---------------------------------------------------------------------------------------------------------+ +| 0 | ++---------------------------------------------------------------------------------------------------------+ + +ADMIN COMPACT_TABLE( + 'test', + 'strict_window', + 'window=3600,start_time=1970-01-01T00:00:00Z,end_time=1970-01-01T01:00:00Z' +); + ++---------------------------------------------------------------------------------------------------------------------------+ +| ADMIN COMPACT_TABLE('test', 'strict_window', 'window=3600,start_time=1970-01-01T00:00:00Z,end_time=1970-01-01T01:00:00Z') | ++---------------------------------------------------------------------------------------------------------------------------+ +| 0 | ++---------------------------------------------------------------------------------------------------------------------------+ + SELECT FLUSH_TABLE('test'); +---------------------------+ diff --git a/tests/cases/standalone/common/function/admin/flush_compact_table.sql b/tests/cases/standalone/common/function/admin/flush_compact_table.sql index a1a316b35c..d5e6c9072a 100644 --- a/tests/cases/standalone/common/function/admin/flush_compact_table.sql +++ b/tests/cases/standalone/common/function/admin/flush_compact_table.sql @@ -10,6 +10,18 @@ ADMIN FLUSH_TABLE('test'); ADMIN COMPACT_TABLE('test'); +ADMIN COMPACT_TABLE( + 'test', + 'regular', + 'start_time=1970-01-01T00:00:00Z,end_time=1970-01-01T01:00:00Z' +); + +ADMIN COMPACT_TABLE( + 'test', + 'strict_window', + 'window=3600,start_time=1970-01-01T00:00:00Z,end_time=1970-01-01T01:00:00Z' +); + SELECT FLUSH_TABLE('test'); SELECT COMPACT_TABLE('test');