feat(mito2): introduce TWCS active window compaction (#9011)

* feat(mito2): support independent TWCS trigger_file_num for active and inactive windows

Split the single TWCS trigger_file_num into per-window-state thresholds:
the active window keeps the existing trigger (default 4, legacy
compaction.twcs.trigger_file_num stays a compatible alias), while
inactive windows use a new trigger (default 2). Inactive windows
additionally fall back from balanced L0-only/L1-only candidates to a
progress-making unbalanced mixed candidate so historical windows can
converge; the active window retains the row/byte balance guards.

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>

* fix(mito2): bound inactive TWCS window convergence by rewrite budget

Inactive windows that cannot compact within one level previously either
stayed stuck (a threshold-qualified but unbalanced level returned no
candidate without trying any fallback) or fell back to a mixed merge
with no balance checks at all, which could rewrite a huge compacted file
to absorb tiny fresh files.

Inactive windows now converge progressively: threshold-qualified
balanced picks, sub-threshold balanced single-level picks, a mixed merge
whose total rewrite must fit in the output file budget, and finally an
L0-only merge without balance checks. Windows that qualify for none of
these are left uncompacted, bounding write amplification.

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>

* test(mito2): cover TWCS window trigger options in alter_table_options sqlness case

Exercise SET/UNSET of compaction.twcs.active_window.trigger_file_num and
compaction.twcs.inactive_window.trigger_file_num end to end, including
that setting the canonical active key removes the legacy
compaction.twcs.trigger_file_num alias from the table options.

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>

* fix(mito2): derive TWCS active window from the max-sequence file

The active window was determined by the max event-time window among
level-0 files. Between an L0 compaction removing its inputs and the next
flush landing, level 0 is empty, so the active window transiently became
None and every window fell back to the inactive rules - triggering
full-window convergence merges during ongoing ingestion whose outputs
are then superseded by new data.

Flush and compaction outputs both inherit the max input sequence, so the
file with the highest sequence across all levels always tracks the most
recent write. Use its window as the active window, falling back to the
previous L0-based rule when no file carries a sequence (legacy files).

The new helper deliberately computes window keys with the
assign_to_windows convention (truncate to seconds, then align up),
because the result is compared against window keys produced there; the
older ceil-based helper is kept unchanged for the legacy fallback path.

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>

* feat(mito2): add active-window L1 compaction safety trigger

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>

* test(compat): cover TWCS active window options

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>

* fix(mito2): align TWCS window trigger validation

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>

* fix(mito2): resolve database TWCS trigger aliases

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>

* fix(mito2): prioritize newer compaction windows

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>

* fix(mito2): repick serial compaction outputs

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>

* feat(mito2): configure inactive-window L1 trigger

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>

* feat(mito2): prioritize TWCS compaction candidates

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>

* fix(meta): preserve TWCS trigger downgrade compatibility

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>

* fix(sql): distinguish invalid database option values

Separate database option key and value validation so recognized keys report the invalid value and its constraint. Add parser coverage for invalid, valid, and unknown options.

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>

* docs(meta): explain TWCS legacy key compatibility

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>

* refactor(mito2): clarify active window trigger field

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>

* fix(options): normalize TWCS trigger aliases

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>

* fix(mito2): ignore ineligible files for active window

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>

* fix(mito2): mark explicit TWCS options as overrides

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>

* fix(mito2): normalize zero compaction output size

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>

* fix(options): validate database TWCS trigger values

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>

* style(store-api): collapse TWCS alias conflict condition

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>

---------

Signed-off-by: Lei, HUANG <ratuthomm@gmail.com>
This commit is contained in:
Lei, HUANG
2026-09-09 13:02:26 +00:00
committed by GitHub
parent 7cf84892d2
commit 6fa1023b7f
33 changed files with 2763 additions and 212 deletions
+19 -1
View File
@@ -135,6 +135,22 @@ pub enum Error {
location: Location,
},
#[snafu(display(
"Conflicting table options: {}={} and {}={}",
first_key,
first_value,
second_key,
second_value
))]
ConflictingTableOptions {
first_key: String,
first_value: String,
second_key: String,
second_value: String,
#[snafu(implicit)]
location: Location,
},
#[snafu(display("Invalid alter table({}) request: {}", table, err))]
InvalidAlterRequest {
table: String,
@@ -223,7 +239,9 @@ impl ErrorExt for Error {
Error::InvalidColumnOption { .. } => StatusCode::InvalidArguments,
Error::ColumnNotExists { .. } => StatusCode::TableColumnNotFound,
Error::Unsupported { .. } => StatusCode::Unsupported,
Error::ParseTableOption { .. } => StatusCode::InvalidArguments,
Error::ParseTableOption { .. } | Error::ConflictingTableOptions { .. } => {
StatusCode::InvalidArguments
}
Error::MissingTimeIndexColumn { .. } => StatusCode::IllegalState,
Error::SetSkippingOptions { .. }
| Error::UnsetSkippingOptions { .. }
+60 -2
View File
@@ -32,6 +32,7 @@ use store_api::metric_engine_consts::PHYSICAL_TABLE_METADATA_KEY;
use store_api::mito_engine_options::{
APPEND_MODE_KEY, AUTO_FLUSH_INTERVAL_KEY, COMPACTION_TYPE, COMPACTION_TYPE_TWCS,
MAX_ROW_GROUP_ROW_COUNT, MERGE_MODE_KEY, PRESERVE_ROW_SEQUENCE, SKIP_WAL_KEY, SST_FORMAT_KEY,
TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM, TWCS_TRIGGER_FILE_NUM,
};
use store_api::region_request::{SetRegionOption, UnsetRegionOption};
use store_api::storage::{ColumnDescriptor, ColumnDescriptorBuilder, ColumnId};
@@ -362,8 +363,22 @@ impl TableMeta {
new_options.ttl = *new_ttl;
}
SetRegionOption::Twsc(key, value) => {
let persisted_key = if matches!(
key.as_str(),
TWCS_TRIGGER_FILE_NUM | TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM
) {
new_options.extra_options.remove(TWCS_TRIGGER_FILE_NUM);
new_options
.extra_options
.remove(TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM);
TWCS_TRIGGER_FILE_NUM
} else {
key
};
if !value.is_empty() {
new_options.extra_options.insert(key.clone(), value.clone());
new_options
.extra_options
.insert(persisted_key.to_string(), value.clone());
// Ensure node restart correctly.
new_options.extra_options.insert(
COMPACTION_TYPE.to_string(),
@@ -371,7 +386,7 @@ impl TableMeta {
);
} else {
// Invalidate the previous change option if an empty value has been set.
new_options.extra_options.remove(key.as_str());
new_options.extra_options.remove(persisted_key);
}
}
SetRegionOption::Format(value) => {
@@ -2224,6 +2239,49 @@ mod tests {
);
}
#[test]
fn test_set_twcs_trigger_persists_legacy_key() {
for key in [TWCS_TRIGGER_FILE_NUM, TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM] {
let mut table_options = TableOptions::default();
table_options.extra_options.insert(
TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM.to_string(),
"4".to_string(),
);
let meta = TableMetaBuilder::empty()
.schema(Arc::new(new_test_schema()))
.primary_key_indices(vec![0])
.engine("engine")
.next_column_id(3)
.options(table_options)
.build()
.unwrap();
let alter_kind = AlterKind::SetTableOptions {
options: vec![SetRegionOption::Twsc(key.to_string(), "8".to_string())],
};
let new_meta = meta
.builder_with_alter_kind("my_table", &alter_kind)
.unwrap()
.build()
.unwrap();
assert_eq!(
Some("8"),
new_meta
.options
.extra_options
.get(TWCS_TRIGGER_FILE_NUM)
.map(String::as_str)
);
assert!(
!new_meta
.options
.extra_options
.contains_key(TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM)
);
}
}
#[test]
fn test_set_and_unset_max_row_group_row_count() {
let meta = TableMetaBuilder::empty()
+131 -4
View File
@@ -39,12 +39,14 @@ use store_api::mito_engine_options::{
APPEND_MODE_KEY, COMPACTION_TYPE, MEMTABLE_BULK_ENCODE_BYTES_THRESHOLD,
MEMTABLE_BULK_ENCODE_ROW_THRESHOLD, MEMTABLE_BULK_MAX_MERGE_GROUPS,
MEMTABLE_BULK_MERGE_THRESHOLD, MEMTABLE_TYPE, MERGE_MODE_KEY, SST_FORMAT_KEY,
TWCS_FALLBACK_TO_LOCAL, TWCS_MAX_OUTPUT_FILE_SIZE, TWCS_TIME_WINDOW, TWCS_TRIGGER_FILE_NUM,
is_mito_engine_option_key,
TWCS_ACTIVE_WINDOW_L1_MERGE_TRIGGER, TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM,
TWCS_FALLBACK_TO_LOCAL, TWCS_INACTIVE_WINDOW_L1_MERGE_TRIGGER,
TWCS_INACTIVE_WINDOW_TRIGGER_FILE_NUM, TWCS_MAX_OUTPUT_FILE_SIZE, TWCS_TIME_WINDOW,
TWCS_TRIGGER_FILE_NUM, is_mito_engine_option_key, normalize_twcs_trigger_options,
};
use store_api::region_request::{SetRegionOption, UnsetRegionOption};
use crate::error::{ParseTableOptionSnafu, Result};
use crate::error::{ConflictingTableOptionsSnafu, ParseTableOptionSnafu, Result};
use crate::metadata::{TableId, TableVersion};
use crate::table_reference::TableReference;
@@ -118,6 +120,10 @@ static VALID_DB_OPT_KEYS: Lazy<HashSet<&str>> = Lazy::new(|| {
set.insert(TWCS_FALLBACK_TO_LOCAL);
set.insert(TWCS_TIME_WINDOW);
set.insert(TWCS_TRIGGER_FILE_NUM);
set.insert(TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM);
set.insert(TWCS_ACTIVE_WINDOW_L1_MERGE_TRIGGER);
set.insert(TWCS_INACTIVE_WINDOW_TRIGGER_FILE_NUM);
set.insert(TWCS_INACTIVE_WINDOW_L1_MERGE_TRIGGER);
set.insert(TWCS_MAX_OUTPUT_FILE_SIZE);
set.insert(SST_FORMAT_KEY);
set
@@ -128,6 +134,32 @@ pub fn validate_database_option(key: &str) -> bool {
VALID_DB_OPT_KEYS.contains(&key)
}
/// Validates a database option value, returning the violated constraint on error.
pub fn validate_database_option_value(
key: &str,
value: Option<&str>,
) -> std::result::Result<(), &'static str> {
let (minimum, constraint) = match key {
TWCS_TRIGGER_FILE_NUM
| TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM
| TWCS_INACTIVE_WINDOW_TRIGGER_FILE_NUM => {
(0, "expected a non-negative integer fitting in usize")
}
TWCS_ACTIVE_WINDOW_L1_MERGE_TRIGGER | TWCS_INACTIVE_WINDOW_L1_MERGE_TRIGGER => {
(2, "expected an integer greater than or equal to 2")
}
_ => return Ok(()),
};
if value
.and_then(|value| value.parse::<usize>().ok())
.is_some_and(|files| files >= minimum)
{
Ok(())
} else {
Err(constraint)
}
}
/// Returns true if the `key` is a valid key for any engine or storage.
pub fn validate_table_option(key: &str) -> bool {
if is_supported_in_s3(key) {
@@ -183,11 +215,21 @@ impl TableOptions {
) -> Result<TableOptions> {
let mut options = TableOptions::default();
let kvs: HashMap<String, String> = iter
let mut kvs: HashMap<String, String> = iter
.into_iter()
.map(|(k, v)| (k.to_string(), v.to_string()))
.collect();
normalize_twcs_trigger_options(&mut kvs).map_err(|conflict| {
ConflictingTableOptionsSnafu {
first_key: TWCS_TRIGGER_FILE_NUM,
first_value: conflict.legacy_value,
second_key: TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM,
second_value: conflict.canonical_value,
}
.build()
})?;
if let Some(write_buffer_size) = kvs.get(WRITE_BUFFER_SIZE_KEY) {
let size = ReadableSize::from_str(write_buffer_size).map_err(|_| {
ParseTableOptionSnafu {
@@ -716,6 +758,9 @@ pub struct CopyQueryToRequest {
mod tests {
use std::time::Duration;
use common_error::ext::ErrorExt;
use common_error::status_code::StatusCode;
use super::*;
#[test]
@@ -748,9 +793,60 @@ mod tests {
MEMTABLE_BULK_ENCODE_BYTES_THRESHOLD
));
assert!(validate_database_option(MEMTABLE_BULK_MAX_MERGE_GROUPS));
assert!(validate_database_option(
"compaction.twcs.active_window.trigger_file_num"
));
assert!(validate_database_option(
"compaction.twcs.active_window.l1_merge_trigger"
));
assert!(validate_database_option(
"compaction.twcs.inactive_window.trigger_file_num"
));
assert!(validate_database_option(
"compaction.twcs.inactive_window.l1_merge_trigger"
));
assert!(!validate_database_option("foo"));
}
#[test]
fn test_database_trigger_value_boundaries() {
let maximum = usize::MAX.to_string();
let overflow = format!("{maximum}0");
for (key, minimum) in [
(TWCS_TRIGGER_FILE_NUM, 0),
(TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM, 0),
(TWCS_INACTIVE_WINDOW_TRIGGER_FILE_NUM, 0),
(TWCS_ACTIVE_WINDOW_L1_MERGE_TRIGGER, 2),
(TWCS_INACTIVE_WINDOW_L1_MERGE_TRIGGER, 2),
] {
for invalid in [
None,
Some(""),
Some("invalid"),
Some("-1"),
Some(overflow.as_str()),
] {
assert!(
validate_database_option_value(key, invalid).is_err(),
"{key}: {invalid:?}"
);
}
for valid in ["2", maximum.as_str()] {
assert!(
validate_database_option_value(key, Some(valid)).is_ok(),
"{key}: {valid}"
);
}
for boundary in ["0", "1"] {
assert_eq!(
validate_database_option_value(key, Some(boundary)).is_ok(),
minimum == 0,
"{key}: {boundary}"
);
}
}
}
#[test]
fn test_serialize_table_options() {
let options = TableOptions {
@@ -820,6 +916,37 @@ mod tests {
assert_eq!(options, serialized);
}
#[test]
fn test_table_options_normalizes_twcs_trigger_aliases() {
for options in [
vec![(TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM, "4")],
vec![
(TWCS_TRIGGER_FILE_NUM, "4"),
(TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM, "4"),
],
] {
let table_options = TableOptions::try_from_iter(options).unwrap();
assert_eq!(
HashMap::from([(TWCS_TRIGGER_FILE_NUM.to_string(), "4".to_string())]),
table_options.extra_options
);
}
}
#[test]
fn test_table_options_rejects_conflicting_twcs_trigger_aliases() {
let error = TableOptions::try_from_iter([
(TWCS_TRIGGER_FILE_NUM, "4"),
(TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM, "8"),
])
.unwrap_err();
assert_eq!(StatusCode::InvalidArguments, error.status_code());
assert_eq!(
"Conflicting table options: compaction.twcs.trigger_file_num=4 and compaction.twcs.active_window.trigger_file_num=8",
error.to_string()
);
}
#[test]
fn test_table_options_to_string() {
let options = TableOptions {