From 06631be8595348ca4ce1a9becee75dbb292249b2 Mon Sep 17 00:00:00 2001 From: Ruihang Xia Date: Fri, 10 Jul 2026 20:06:41 +0800 Subject: [PATCH] fix: preserve metric value and scalar semantics Signed-off-by: Ruihang Xia --- src/metric-engine/src/engine/create.rs | 53 +++++++- .../src/sst/parquet/metric_value_split.rs | 2 +- src/query/src/promql/planner.rs | 31 +++-- src/store-api/src/metric_engine_consts.rs | 6 +- .../common/alter/alter_metric_table.result | 46 +++---- .../common/create/create_metric_table.result | 44 +++---- .../common/insert/logical_metric_table.result | 124 ++++++++++++++---- .../common/insert/logical_metric_table.sql | 30 +++++ 8 files changed, 246 insertions(+), 90 deletions(-) diff --git a/src/metric-engine/src/engine/create.rs b/src/metric-engine/src/engine/create.rs index b412ac8d86..5c3158a5c7 100644 --- a/src/metric-engine/src/engine/create.rs +++ b/src/metric-engine/src/engine/create.rs @@ -28,14 +28,14 @@ use snafu::{OptionExt, ResultExt, ensure}; use store_api::metadata::ColumnMetadata; use store_api::metric_engine_consts::{ ALTER_PHYSICAL_EXTENSION_KEY, DATA_REGION_SUBDIR, DATA_SCHEMA_TABLE_ID_COLUMN_NAME, - DATA_SCHEMA_TSID_COLUMN_NAME, DATA_SCHEMA_VALUE_INT_COLUMN_SUFFIX, LOGICAL_TABLE_METADATA_KEY, + DATA_SCHEMA_TSID_COLUMN_NAME, DATA_SCHEMA_VALUE_INT_COLUMN_PREFIX, LOGICAL_TABLE_METADATA_KEY, METADATA_REGION_SUBDIR, METADATA_SCHEMA_KEY_COLUMN_INDEX, METADATA_SCHEMA_KEY_COLUMN_NAME, METADATA_SCHEMA_TIMESTAMP_COLUMN_INDEX, METADATA_SCHEMA_TIMESTAMP_COLUMN_NAME, METADATA_SCHEMA_VALUE_COLUMN_INDEX, METADATA_SCHEMA_VALUE_COLUMN_NAME, is_metric_engine_internal_column, is_metric_engine_value_int_column, metric_engine_value_int_column_name, }; -use store_api::mito_engine_options::{TTL_KEY, WAL_OPTIONS_KEY}; +use store_api::mito_engine_options::{MERGE_MODE_KEY, TTL_KEY, WAL_OPTIONS_KEY}; use store_api::region_engine::RegionEngine; use store_api::region_request::{AffectedRows, PathType, RegionCreateRequest, RegionRequest}; use store_api::storage::RegionId; @@ -553,7 +553,14 @@ impl MetricEngineInner { let table_id_col_def = request.column_metadatas.iter().any(is_metric_name_col); let tsid_col_def = request.column_metadatas.iter().any(is_tsid_col); - append_metric_value_int_columns(&mut data_region_request.column_metadatas); + let value_split_enabled = request + .options + .get(MERGE_MODE_KEY) + .is_none_or(|mode| mode != "last_non_null"); + configure_metric_value_int_columns( + &mut data_region_request.column_metadatas, + value_split_enabled, + ); // change nullability for tag columns data_region_request @@ -593,7 +600,7 @@ fn is_valid_physical_metric_value_int_column(request: &RegionCreateRequest, name return false; } - let Some(value_name) = name.strip_suffix(DATA_SCHEMA_VALUE_INT_COLUMN_SUFFIX) else { + let Some(value_name) = name.strip_prefix(DATA_SCHEMA_VALUE_INT_COLUMN_PREFIX) else { return false; }; let find_column = |name| { @@ -615,7 +622,15 @@ fn is_valid_physical_metric_value_int_column(request: &RegionCreateRequest, name && int_column.column_schema.data_type == ConcreteDataType::int64_datatype() } -fn append_metric_value_int_columns(column_metadatas: &mut Vec) { +fn configure_metric_value_int_columns( + column_metadatas: &mut Vec, + value_split_enabled: bool, +) { + if !value_split_enabled { + column_metadatas + .retain(|metadata| !is_metric_engine_value_int_column(&metadata.column_schema.name)); + } + let mut next_column_id = column_metadatas .iter() .map(|metadata| metadata.column_id) @@ -642,6 +657,10 @@ fn append_metric_value_int_columns(column_metadatas: &mut Vec) { metadata.column_schema.set_nullable(); + if !value_split_enabled { + continue; + } + let int_column_name = metric_engine_value_int_column_name(&metadata.column_schema.name); if existing_names.contains(&int_column_name) { continue; @@ -873,7 +892,7 @@ mod test { requirements: Default::default(), }; MetricEngineInner::verify_region_create_request(&request).unwrap(); - append_metric_value_int_columns(&mut request.column_metadatas); + configure_metric_value_int_columns(&mut request.column_metadatas, true); assert!(request.column_metadatas[2].column_schema.is_nullable()); assert!(request.column_metadatas[3].column_schema.is_nullable()); @@ -1244,6 +1263,28 @@ mod test { ); assert!(value_int_metadata.column_schema.is_nullable()); + let mut last_non_null_request = request.clone(); + last_non_null_request + .options + .insert(MERGE_MODE_KEY.to_string(), "last_non_null".to_string()); + let last_non_null_data_region_request = + engine_inner.create_request_for_data_region(&last_non_null_request); + assert!( + last_non_null_data_region_request + .column_metadatas + .iter() + .all(|metadata| !is_metric_engine_value_int_column(&metadata.column_schema.name)) + ); + assert!( + last_non_null_data_region_request + .column_metadatas + .iter() + .find(|metadata| metadata.column_schema.name == "value") + .unwrap() + .column_schema + .is_nullable() + ); + let table_id_metadata = data_region_request .column_metadatas .iter() diff --git a/src/mito2/src/sst/parquet/metric_value_split.rs b/src/mito2/src/sst/parquet/metric_value_split.rs index f5ae9f4566..8dd6f5b43f 100644 --- a/src/mito2/src/sst/parquet/metric_value_split.rs +++ b/src/mito2/src/sst/parquet/metric_value_split.rs @@ -314,7 +314,7 @@ mod tests { true, ), ( - "greptime_value__metric_int", + "__metric_int_greptime_value", Arc::new(Int64Array::from(vec![None; 9])) as ArrayRef, true, ), diff --git a/src/query/src/promql/planner.rs b/src/query/src/promql/planner.rs index 7dae30e613..adae31ed84 100644 --- a/src/query/src/promql/planner.rs +++ b/src/query/src/promql/planner.rs @@ -3782,18 +3782,13 @@ impl PromPlanner { operator: SCALAR_FUNCTION }, ); - let scalar_tags = if self.ctx.use_tsid { - Vec::new() - } else { - self.ctx - .tag_columns - .iter() - .filter(|tag| { - Self::qualified_column_if_available(input.schema(), None, tag).is_some() - }) - .cloned() - .collect() - }; + let scalar_tags = self + .ctx + .tag_columns + .iter() + .filter(|tag| Self::qualified_column_if_available(input.schema(), None, tag).is_some()) + .cloned() + .collect::>(); let scalar_plan = LogicalPlan::Extension(Extension { node: Arc::new( ScalarCalculate::new( @@ -5646,7 +5641,7 @@ mod test { assert!(plan_str.contains("PromSeriesDivide: tags=[\"__tsid\"]")); assert!(plan_str.contains("TableScan: phy")); assert!(!plan_str.contains("TableScan: some_metric")); - assert!(!plan_str.contains("field_0__metric_int")); + assert!(!plan_str.contains("__metric_int_field_0")); assert!(!plan_str.contains("coalesce")); } @@ -6518,6 +6513,16 @@ mod test { assert!(!plan_str.contains("PromInstantManipulate: range=[99999000..99999000]")); } + #[tokio::test] + async fn scalar_tsid_input_preserves_available_series_labels() { + let plan_str = build_optimized_tsid_plan("scalar(some_metric)", 2, 1, 100_000, 1).await; + + assert!( + plan_str.contains("ScalarCalculate: tags=[\"tag_0\", \"tag_1\"]"), + "{plan_str}" + ); + } + #[tokio::test] async fn scalar_count_count_rewrite_applies_inside_binary_expr_for_tsid_input() { let plan_str = build_optimized_tsid_plan( diff --git a/src/store-api/src/metric_engine_consts.rs b/src/store-api/src/metric_engine_consts.rs index fb84fc4498..126f4342a9 100644 --- a/src/store-api/src/metric_engine_consts.rs +++ b/src/store-api/src/metric_engine_consts.rs @@ -33,7 +33,7 @@ pub const METADATA_SCHEMA_VALUE_COLUMN_INDEX: usize = 2; /// Column name of internal column `__metric` that stores the original metric name pub const DATA_SCHEMA_TABLE_ID_COLUMN_NAME: &str = "__table_id"; pub const DATA_SCHEMA_TSID_COLUMN_NAME: &str = "__tsid"; -pub const DATA_SCHEMA_VALUE_INT_COLUMN_SUFFIX: &str = "__metric_int"; +pub const DATA_SCHEMA_VALUE_INT_COLUMN_PREFIX: &str = "__metric_int_"; pub const METADATA_REGION_SUBDIR: &str = "metadata"; pub const DATA_REGION_SUBDIR: &str = "data"; @@ -91,12 +91,12 @@ pub fn is_metric_engine_internal_column(name: &str) -> bool { /// Returns the physical integer companion column name for a metric value column. pub fn metric_engine_value_int_column_name(value_column_name: &str) -> String { - format!("{value_column_name}{DATA_SCHEMA_VALUE_INT_COLUMN_SUFFIX}") + format!("{DATA_SCHEMA_VALUE_INT_COLUMN_PREFIX}{value_column_name}") } /// Returns true if the column is a physical integer companion for a metric value column. pub fn is_metric_engine_value_int_column(name: &str) -> bool { - name.ends_with(DATA_SCHEMA_VALUE_INT_COLUMN_SUFFIX) + name.starts_with(DATA_SCHEMA_VALUE_INT_COLUMN_PREFIX) } /// Returns true if it's metric engine diff --git a/tests/cases/standalone/common/alter/alter_metric_table.result b/tests/cases/standalone/common/alter/alter_metric_table.result index 8acd39e1ab..b9f6bd6d4d 100644 --- a/tests/cases/standalone/common/alter/alter_metric_table.result +++ b/tests/cases/standalone/common/alter/alter_metric_table.result @@ -50,17 +50,17 @@ DESC TABLE t2; DESC TABLE phy; -+-----------------+----------------------+-----+------+---------+---------------+ -| Column | Type | Key | Null | Default | Semantic Type | -+-----------------+----------------------+-----+------+---------+---------------+ -| ts | TimestampMillisecond | PRI | NO | | TIMESTAMP | -| val | Float64 | | YES | | FIELD | -| val__metric_int | Int64 | | YES | | FIELD | -| __table_id | UInt32 | PRI | NO | | TAG | -| __tsid | UInt64 | PRI | NO | | TAG | -| host | String | PRI | YES | | TAG | -| job | String | PRI | YES | | TAG | -+-----------------+----------------------+-----+------+---------+---------------+ ++------------------+----------------------+-----+------+---------+---------------+ +| Column | Type | Key | Null | Default | Semantic Type | ++------------------+----------------------+-----+------+---------+---------------+ +| ts | TimestampMillisecond | PRI | NO | | TIMESTAMP | +| val | Float64 | | YES | | FIELD | +| __metric_int_val | Int64 | | YES | | FIELD | +| __table_id | UInt32 | PRI | NO | | TAG | +| __tsid | UInt64 | PRI | NO | | TAG | +| host | String | PRI | YES | | TAG | +| job | String | PRI | YES | | TAG | ++------------------+----------------------+-----+------+---------+---------------+ ALTER TABLE t1 ADD COLUMN k STRING PRIMARY KEY; @@ -94,18 +94,18 @@ DESC TABLE t2; DESC TABLE phy; -+-----------------+----------------------+-----+------+---------+---------------+ -| Column | Type | Key | Null | Default | Semantic Type | -+-----------------+----------------------+-----+------+---------+---------------+ -| ts | TimestampMillisecond | PRI | NO | | TIMESTAMP | -| val | Float64 | | YES | | FIELD | -| val__metric_int | Int64 | | YES | | FIELD | -| __table_id | UInt32 | PRI | NO | | TAG | -| __tsid | UInt64 | PRI | NO | | TAG | -| host | String | PRI | YES | | TAG | -| job | String | PRI | YES | | TAG | -| k | String | PRI | YES | | TAG | -+-----------------+----------------------+-----+------+---------+---------------+ ++------------------+----------------------+-----+------+---------+---------------+ +| Column | Type | Key | Null | Default | Semantic Type | ++------------------+----------------------+-----+------+---------+---------------+ +| ts | TimestampMillisecond | PRI | NO | | TIMESTAMP | +| val | Float64 | | YES | | FIELD | +| __metric_int_val | Int64 | | YES | | FIELD | +| __table_id | UInt32 | PRI | NO | | TAG | +| __tsid | UInt64 | PRI | NO | | TAG | +| host | String | PRI | YES | | TAG | +| job | String | PRI | YES | | TAG | +| k | String | PRI | YES | | TAG | ++------------------+----------------------+-----+------+---------+---------------+ DROP TABLE t1; diff --git a/tests/cases/standalone/common/create/create_metric_table.result b/tests/cases/standalone/common/create/create_metric_table.result index e27b15040b..ac69af9bff 100644 --- a/tests/cases/standalone/common/create/create_metric_table.result +++ b/tests/cases/standalone/common/create/create_metric_table.result @@ -65,17 +65,17 @@ SELECT table_catalog, table_schema, table_name, table_type, engine FROM informat DESC TABLE phy; -+-----------------+----------------------+-----+------+---------+---------------+ -| Column | Type | Key | Null | Default | Semantic Type | -+-----------------+----------------------+-----+------+---------+---------------+ -| ts | TimestampMillisecond | PRI | NO | | TIMESTAMP | -| val | Float64 | | YES | | FIELD | -| val__metric_int | Int64 | | YES | | FIELD | -| __table_id | UInt32 | PRI | NO | | TAG | -| __tsid | UInt64 | PRI | NO | | TAG | -| host | String | PRI | YES | | TAG | -| job | String | PRI | YES | | TAG | -+-----------------+----------------------+-----+------+---------+---------------+ ++------------------+----------------------+-----+------+---------+---------------+ +| Column | Type | Key | Null | Default | Semantic Type | ++------------------+----------------------+-----+------+---------+---------------+ +| ts | TimestampMillisecond | PRI | NO | | TIMESTAMP | +| val | Float64 | | YES | | FIELD | +| __metric_int_val | Int64 | | YES | | FIELD | +| __table_id | UInt32 | PRI | NO | | TAG | +| __tsid | UInt64 | PRI | NO | | TAG | +| host | String | PRI | YES | | TAG | +| job | String | PRI | YES | | TAG | ++------------------+----------------------+-----+------+---------+---------------+ SHOW CREATE TABLE phy; @@ -126,17 +126,17 @@ Error: 1004(InvalidArguments), Physical region is busy, there are still some log -- metadata should be restored DESC TABLE phy; -+-----------------+----------------------+-----+------+---------+---------------+ -| Column | Type | Key | Null | Default | Semantic Type | -+-----------------+----------------------+-----+------+---------+---------------+ -| ts | TimestampMillisecond | PRI | NO | | TIMESTAMP | -| val | Float64 | | YES | | FIELD | -| val__metric_int | Int64 | | YES | | FIELD | -| __table_id | UInt32 | PRI | NO | | TAG | -| __tsid | UInt64 | PRI | NO | | TAG | -| host | String | PRI | YES | | TAG | -| job | String | PRI | YES | | TAG | -+-----------------+----------------------+-----+------+---------+---------------+ ++------------------+----------------------+-----+------+---------+---------------+ +| Column | Type | Key | Null | Default | Semantic Type | ++------------------+----------------------+-----+------+---------+---------------+ +| ts | TimestampMillisecond | PRI | NO | | TIMESTAMP | +| val | Float64 | | YES | | FIELD | +| __metric_int_val | Int64 | | YES | | FIELD | +| __table_id | UInt32 | PRI | NO | | TAG | +| __tsid | UInt64 | PRI | NO | | TAG | +| host | String | PRI | YES | | TAG | +| job | String | PRI | YES | | TAG | ++------------------+----------------------+-----+------+---------+---------------+ DROP TABLE t1; diff --git a/tests/cases/standalone/common/insert/logical_metric_table.result b/tests/cases/standalone/common/insert/logical_metric_table.result index f5c1e1cead..4851284c28 100644 --- a/tests/cases/standalone/common/insert/logical_metric_table.result +++ b/tests/cases/standalone/common/insert/logical_metric_table.result @@ -51,17 +51,17 @@ Affected Rows: 0 DESC TABLE phy; -+-----------------+----------------------+-----+------+---------+---------------+ -| Column | Type | Key | Null | Default | Semantic Type | -+-----------------+----------------------+-----+------+---------+---------------+ -| ts | TimestampMillisecond | PRI | NO | | TIMESTAMP | -| val | Float64 | | YES | | FIELD | -| val__metric_int | Int64 | | YES | | FIELD | -| __table_id | UInt32 | PRI | NO | | TAG | -| __tsid | UInt64 | PRI | NO | | TAG | -| host | String | PRI | YES | | TAG | -| job | String | PRI | YES | | TAG | -+-----------------+----------------------+-----+------+---------+---------------+ ++------------------+----------------------+-----+------+---------+---------------+ +| Column | Type | Key | Null | Default | Semantic Type | ++------------------+----------------------+-----+------+---------+---------------+ +| ts | TimestampMillisecond | PRI | NO | | TIMESTAMP | +| val | Float64 | | YES | | FIELD | +| __metric_int_val | Int64 | | YES | | FIELD | +| __table_id | UInt32 | PRI | NO | | TAG | +| __tsid | UInt64 | PRI | NO | | TAG | +| host | String | PRI | YES | | TAG | +| job | String | PRI | YES | | TAG | ++------------------+----------------------+-----+------+---------+---------------+ SELECT ts, val, __tsid, host, job FROM phy; @@ -78,6 +78,86 @@ DROP TABLE phy; Affected Rows: 0 +-- Value splitting is disabled for last_non_null so the two physical value +-- columns cannot be merged independently into a stale logical value. +CREATE TABLE phy_last_non_null ( + ts TIMESTAMP TIME INDEX, + val DOUBLE +) ENGINE = metric WITH ( + "physical_metric_table" = "", + "merge_mode" = "last_non_null" +); + +Affected Rows: 0 + +CREATE TABLE metric_last_non_null ( + ts TIMESTAMP TIME INDEX, + val DOUBLE, + host STRING PRIMARY KEY +) ENGINE = metric WITH ("on_physical_table" = "phy_last_non_null"); + +Affected Rows: 0 + +INSERT INTO metric_last_non_null VALUES ('host1', 1, 1.0); + +Affected Rows: 1 + +ADMIN flush_table('phy_last_non_null'); + ++----------------------------------------+ +| ADMIN flush_table('phy_last_non_null') | ++----------------------------------------+ +| 0 | ++----------------------------------------+ + +INSERT INTO metric_last_non_null VALUES ('host1', 1, 1.5); + +Affected Rows: 1 + +SELECT host, ts, val FROM metric_last_non_null; + ++-------+-------------------------+-----+ +| host | ts | val | ++-------+-------------------------+-----+ +| host1 | 1970-01-01T00:00:00.001 | 1.5 | ++-------+-------------------------+-----+ + +ADMIN flush_table('phy_last_non_null'); + ++----------------------------------------+ +| ADMIN flush_table('phy_last_non_null') | ++----------------------------------------+ +| 0 | ++----------------------------------------+ + +SELECT host, ts, val FROM metric_last_non_null; + ++-------+-------------------------+-----+ +| host | ts | val | ++-------+-------------------------+-----+ +| host1 | 1970-01-01T00:00:00.001 | 1.5 | ++-------+-------------------------+-----+ + +DESC TABLE phy_last_non_null; + ++------------+----------------------+-----+------+---------+---------------+ +| Column | Type | Key | Null | Default | Semantic Type | ++------------+----------------------+-----+------+---------+---------------+ +| ts | TimestampMillisecond | PRI | NO | | TIMESTAMP | +| val | Float64 | | YES | | FIELD | +| __table_id | UInt32 | PRI | NO | | TAG | +| __tsid | UInt64 | PRI | NO | | TAG | +| host | String | PRI | YES | | TAG | ++------------+----------------------+-----+------+---------+---------------+ + +DROP TABLE metric_last_non_null; + +Affected Rows: 0 + +DROP TABLE phy_last_non_null; + +Affected Rows: 0 + CREATE TABLE phy_default (ts timestamp time index, val double default 42) engine=metric with ("physical_metric_table" = ""); Affected Rows: 0 @@ -205,17 +285,17 @@ Affected Rows: 0 DESC TABLE phy; -+-----------------+----------------------+-----+------+---------+---------------+ -| Column | Type | Key | Null | Default | Semantic Type | -+-----------------+----------------------+-----+------+---------+---------------+ -| ts | TimestampMillisecond | PRI | NO | | TIMESTAMP | -| val | Float64 | | YES | | FIELD | -| val__metric_int | Int64 | | YES | | FIELD | -| __table_id | UInt32 | PRI | NO | | TAG | -| __tsid | UInt64 | PRI | NO | | TAG | -| host | String | PRI | YES | | TAG | -| job | String | PRI | YES | | TAG | -+-----------------+----------------------+-----+------+---------+---------------+ ++------------------+----------------------+-----+------+---------+---------------+ +| Column | Type | Key | Null | Default | Semantic Type | ++------------------+----------------------+-----+------+---------+---------------+ +| ts | TimestampMillisecond | PRI | NO | | TIMESTAMP | +| val | Float64 | | YES | | FIELD | +| __metric_int_val | Int64 | | YES | | FIELD | +| __table_id | UInt32 | PRI | NO | | TAG | +| __tsid | UInt64 | PRI | NO | | TAG | +| host | String | PRI | YES | | TAG | +| job | String | PRI | YES | | TAG | ++------------------+----------------------+-----+------+---------+---------------+ DROP TABLE phy; diff --git a/tests/cases/standalone/common/insert/logical_metric_table.sql b/tests/cases/standalone/common/insert/logical_metric_table.sql index 9501c9bb16..707b887396 100644 --- a/tests/cases/standalone/common/insert/logical_metric_table.sql +++ b/tests/cases/standalone/common/insert/logical_metric_table.sql @@ -24,6 +24,36 @@ SELECT ts, val, __tsid, host, job FROM phy; DROP TABLE phy; +-- Value splitting is disabled for last_non_null so the two physical value +-- columns cannot be merged independently into a stale logical value. +CREATE TABLE phy_last_non_null ( + ts TIMESTAMP TIME INDEX, + val DOUBLE +) ENGINE = metric WITH ( + "physical_metric_table" = "", + "merge_mode" = "last_non_null" +); + +CREATE TABLE metric_last_non_null ( + ts TIMESTAMP TIME INDEX, + val DOUBLE, + host STRING PRIMARY KEY +) ENGINE = metric WITH ("on_physical_table" = "phy_last_non_null"); + +INSERT INTO metric_last_non_null VALUES ('host1', 1, 1.0); +ADMIN flush_table('phy_last_non_null'); +INSERT INTO metric_last_non_null VALUES ('host1', 1, 1.5); + +SELECT host, ts, val FROM metric_last_non_null; + +ADMIN flush_table('phy_last_non_null'); + +SELECT host, ts, val FROM metric_last_non_null; +DESC TABLE phy_last_non_null; + +DROP TABLE metric_last_non_null; +DROP TABLE phy_last_non_null; + CREATE TABLE phy_default (ts timestamp time index, val double default 42) engine=metric with ("physical_metric_table" = ""); CREATE TABLE t_default (ts timestamp time index, val double default 42, host string primary key) engine = metric with ("on_physical_table" = "phy_default");