fix: preserve metric value and scalar semantics

Signed-off-by: Ruihang Xia <waynestxia@gmail.com>
This commit is contained in:
Ruihang Xia
2026-07-10 20:06:41 +08:00
parent 993aa19c89
commit 06631be859
8 changed files with 246 additions and 90 deletions
+47 -6
View File
@@ -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<ColumnMetadata>) {
fn configure_metric_value_int_columns(
column_metadatas: &mut Vec<ColumnMetadata>,
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<ColumnMetadata>) {
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()
@@ -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,
),
+18 -13
View File
@@ -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::<Vec<_>>();
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(
+3 -3
View File
@@ -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
@@ -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;
@@ -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;
@@ -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;
@@ -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");