From d28bca855f83ad7369554fc27284e7788cdf861c Mon Sep 17 00:00:00 2001 From: Ruihang Xia Date: Fri, 10 Jul 2026 13:09:14 +0800 Subject: [PATCH] feat(mito2): split integral metric values in parquet Signed-off-by: Ruihang Xia --- src/mito2/src/flush.rs | 21 +- src/mito2/src/sst/parquet.rs | 1 + src/mito2/src/sst/parquet/flat_format.rs | 10 +- .../src/sst/parquet/metric_value_split.rs | 411 ++++++++++++++++++ src/mito2/src/sst/parquet/writer.rs | 19 +- src/store-api/src/metric_engine_consts.rs | 15 +- 6 files changed, 467 insertions(+), 10 deletions(-) create mode 100644 src/mito2/src/sst/parquet/metric_value_split.rs diff --git a/src/mito2/src/flush.rs b/src/mito2/src/flush.rs index f087d8a554..510f71db04 100644 --- a/src/mito2/src/flush.rs +++ b/src/mito2/src/flush.rs @@ -63,6 +63,7 @@ use crate::request::{ use crate::schedule::scheduler::{Job, SchedulerRef}; use crate::sst::file::FileMeta; use crate::sst::parquet::metadata::extract_primary_key_range; +use crate::sst::parquet::metric_value_split::metric_value_split_columns; use crate::sst::parquet::{ DEFAULT_READ_BATCH_SIZE, DEFAULT_ROW_GROUP_SIZE, SstInfo, WriteOptions, flat_format, }; @@ -644,11 +645,14 @@ impl RegionFlushTask { ); let field_column_start = flat_format::field_column_start(&version.metadata, batch_schema.fields().len()); + let disable_encoded_ranges = + !metric_value_split_columns(&version.metadata, &batch_schema).is_empty(); let flat_sources = memtable_flat_sources( batch_schema, mem_ranges, &version.options, field_column_start, + disable_encoded_ranges, )?; let mut tasks = Vec::with_capacity(flat_sources.encoded.len() + flat_sources.sources.len()); let num_encoded = flat_sources.encoded.len(); @@ -820,6 +824,7 @@ fn memtable_flat_sources( mem_ranges: MemtableRanges, options: &RegionOptions, field_column_start: usize, + disable_encoded_ranges: bool, ) -> Result { let MemtableRanges { ranges } = mem_ranges; let mut flat_sources = FlatSources { @@ -832,7 +837,7 @@ fn memtable_flat_sources( let only_range = ranges.into_values().next().unwrap(); let max_sequence = only_range.stats().max_sequence(); - if let Some(encoded) = only_range.encoded() { + if !disable_encoded_ranges && let Some(encoded) = only_range.encoded() { flat_sources.encoded.push((encoded, max_sequence)); } else { let iter = only_range.build_record_batch_iter(None, None)?; @@ -876,7 +881,7 @@ fn memtable_flat_sources( }; for (_range_id, range) in ranges { - if let Some(encoded) = range.encoded() { + if !disable_encoded_ranges && let Some(encoded) = range.encoded() { let max_sequence = range.stats().max_sequence(); flat_sources.encoded.push((encoded, max_sequence)); continue; @@ -1932,6 +1937,7 @@ mod tests { mem_ranges, &options, metadata.primary_key.len(), + false, ) .unwrap(); assert!(flat_sources.encoded.is_empty()); @@ -1958,9 +1964,14 @@ mod tests { ..Default::default() }; - let flat_sources = - memtable_flat_sources(schema, mem_ranges, &options, metadata.primary_key.len()) - .unwrap(); + let flat_sources = memtable_flat_sources( + schema, + mem_ranges, + &options, + metadata.primary_key.len(), + false, + ) + .unwrap(); assert!(flat_sources.encoded.is_empty()); assert_eq!(1, flat_sources.sources.len()); diff --git a/src/mito2/src/sst/parquet.rs b/src/mito2/src/sst/parquet.rs index 62874207dd..202476ff40 100644 --- a/src/mito2/src/sst/parquet.rs +++ b/src/mito2/src/sst/parquet.rs @@ -29,6 +29,7 @@ pub mod flat_format; pub mod format; pub(crate) mod helper; pub mod metadata; +pub(crate) mod metric_value_split; pub mod prefilter; pub mod push_decoder; pub mod read_columns; diff --git a/src/mito2/src/sst/parquet/flat_format.rs b/src/mito2/src/sst/parquet/flat_format.rs index 7e84d5ba61..50af9152ee 100644 --- a/src/mito2/src/sst/parquet/flat_format.rs +++ b/src/mito2/src/sst/parquet/flat_format.rs @@ -57,6 +57,9 @@ use crate::sst::parquet::format::{ FIXED_POS_COLUMN_NUM, FormatProjection, INTERNAL_COLUMN_NUM, PrimaryKeyArray, PrimaryKeyReadFormat, StatValues, column_null_counts, column_values, }; +use crate::sst::parquet::metric_value_split::{ + MetricValueSplitColumn, metric_value_split_columns, split_metric_value_columns, +}; use crate::sst::parquet::read_columns::ParquetReadColumns; use crate::sst::{ FlatSchemaOptions, flat_sst_arrow_schema_column_num, tag_maybe_to_dictionary_field, @@ -68,15 +71,18 @@ pub(crate) struct FlatWriteFormat { /// SST file schema. arrow_schema: SchemaRef, override_sequence: Option, + metric_value_split_columns: Vec, } impl FlatWriteFormat { /// Creates a new helper. pub(crate) fn new(metadata: RegionMetadataRef, options: &FlatSchemaOptions) -> FlatWriteFormat { let arrow_schema = to_flat_sst_arrow_schema(&metadata, options); + let metric_value_split_columns = metric_value_split_columns(&metadata, &arrow_schema); FlatWriteFormat { arrow_schema, override_sequence: None, + metric_value_split_columns, } } @@ -99,8 +105,10 @@ impl FlatWriteFormat { pub(crate) fn convert_batch(&self, batch: &RecordBatch) -> Result { debug_assert_eq!(batch.num_columns(), self.arrow_schema.fields().len()); + let batch = split_metric_value_columns(batch, &self.metric_value_split_columns)?; + let Some(override_sequence) = self.override_sequence else { - return Ok(batch.clone()); + return Ok(batch); }; let mut columns = batch.columns().to_vec(); diff --git a/src/mito2/src/sst/parquet/metric_value_split.rs b/src/mito2/src/sst/parquet/metric_value_split.rs new file mode 100644 index 0000000000..82c7ebe1b8 --- /dev/null +++ b/src/mito2/src/sst/parquet/metric_value_split.rs @@ -0,0 +1,411 @@ +// Copyright 2023 Greptime Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use std::collections::HashMap; +use std::sync::Arc; + +use api::v1::SemanticType; +use datatypes::arrow::array::{ + Array, ArrayRef, BinaryArray, DictionaryArray, Float64Array, Float64Builder, Int64Array, + Int64Builder, +}; +use datatypes::arrow::datatypes::{SchemaRef, UInt32Type}; +use datatypes::arrow::record_batch::RecordBatch; +use datatypes::prelude::ConcreteDataType; +use snafu::{OptionExt, ResultExt, ensure}; +use store_api::metadata::RegionMetadata; +use store_api::metric_engine_consts::{ + DATA_SCHEMA_TABLE_ID_COLUMN_NAME, DATA_SCHEMA_TSID_COLUMN_NAME, + metric_engine_value_int_column_name, +}; +use store_api::storage::consts::ReservedColumnId; + +use crate::error::{InvalidRecordBatchSnafu, NewRecordBatchSnafu, Result}; +use crate::sst::parquet::flat_format::primary_key_column_index; + +#[derive(Debug, Clone)] +pub(crate) struct MetricValueSplitColumn { + pub(crate) float_index: usize, + pub(crate) int_index: usize, +} + +pub(crate) fn metric_value_split_columns( + metadata: &RegionMetadata, + arrow_schema: &SchemaRef, +) -> Vec { + if !is_metric_engine_data_region(metadata) { + return vec![]; + } + + metadata + .field_columns() + .filter(|column| column.column_schema.data_type == ConcreteDataType::float64_datatype()) + .filter_map(|float_column| { + let int_name = metric_engine_value_int_column_name(&float_column.column_schema.name); + let int_column = metadata.column_by_name(&int_name)?; + if int_column.semantic_type != SemanticType::Field + || int_column.column_schema.data_type != ConcreteDataType::int64_datatype() + { + return None; + } + + let float_index = arrow_schema + .index_of(&float_column.column_schema.name) + .ok()?; + let int_index = arrow_schema.index_of(&int_name).ok()?; + Some(MetricValueSplitColumn { + float_index, + int_index, + }) + }) + .collect() +} + +fn is_metric_engine_data_region(metadata: &RegionMetadata) -> bool { + let has_internal_tag = |name, column_id| { + metadata.column_by_name(name).is_some_and(|column| { + column.semantic_type == SemanticType::Tag + && column.column_id == column_id + && metadata.primary_key.contains(&column_id) + }) + }; + + has_internal_tag( + DATA_SCHEMA_TABLE_ID_COLUMN_NAME, + ReservedColumnId::table_id(), + ) && has_internal_tag(DATA_SCHEMA_TSID_COLUMN_NAME, ReservedColumnId::tsid()) +} + +pub(crate) fn split_metric_value_columns( + batch: &RecordBatch, + split_columns: &[MetricValueSplitColumn], +) -> Result { + if split_columns.is_empty() { + return Ok(batch.clone()); + } + + let (series_ids, num_series) = primary_key_series_ids(batch)?; + let mut columns = batch.columns().to_vec(); + for split_column in split_columns { + let float_array = batch + .column(split_column.float_index) + .as_any() + .downcast_ref::() + .with_context(|| InvalidRecordBatchSnafu { + reason: format!( + "expected Float64 metric value column at index {}, got {:?}", + split_column.float_index, + batch.column(split_column.float_index).data_type() + ), + })?; + let int_array = batch + .column(split_column.int_index) + .as_any() + .downcast_ref::() + .with_context(|| InvalidRecordBatchSnafu { + reason: format!( + "expected Int64 metric value column {} at index {}, got {:?}", + batch.schema().field(split_column.int_index).name(), + split_column.int_index, + batch.column(split_column.int_index).data_type() + ), + })?; + + let integer_series = integer_series_flags(&series_ids, num_series, float_array, int_array); + let (float_output, int_output) = + split_one_value_column(&series_ids, &integer_series, float_array, int_array); + columns[split_column.float_index] = float_output; + columns[split_column.int_index] = int_output; + } + + RecordBatch::try_new(batch.schema(), columns).context(NewRecordBatchSnafu) +} + +fn integer_series_flags( + series_ids: &[usize], + num_series: usize, + float_array: &Float64Array, + int_array: &Int64Array, +) -> Vec { + let mut integer_series = vec![true; num_series]; + for (row, series_id) in series_ids.iter().copied().enumerate() { + let is_integer = logical_value(float_array, int_array, row) + .is_none_or(|value| integer_value(value).is_some()); + integer_series[series_id] &= is_integer; + } + integer_series +} + +fn split_one_value_column( + series_ids: &[usize], + integer_series: &[bool], + float_array: &Float64Array, + int_array: &Int64Array, +) -> (ArrayRef, ArrayRef) { + let mut float_builder = Float64Builder::with_capacity(float_array.len()); + let mut int_builder = Int64Builder::with_capacity(float_array.len()); + + for (row, series_id) in series_ids.iter().copied().enumerate() { + let Some(value) = logical_value(float_array, int_array, row) else { + float_builder.append_null(); + int_builder.append_null(); + continue; + }; + + if integer_series[series_id] { + float_builder.append_null(); + int_builder.append_value(integer_value(value).unwrap()); + } else { + float_builder.append_value(value); + int_builder.append_null(); + } + } + + ( + Arc::new(float_builder.finish()), + Arc::new(int_builder.finish()), + ) +} + +fn logical_value(float_array: &Float64Array, int_array: &Int64Array, row: usize) -> Option { + if !int_array.is_null(row) { + Some(int_array.value(row) as f64) + } else if !float_array.is_null(row) { + Some(float_array.value(row)) + } else { + None + } +} + +fn integer_value(value: f64) -> Option { + if !value.is_finite() { + return None; + } + if value < i64::MIN as f64 || value >= i64::MAX as f64 { + return None; + } + + let int_value = value as i64; + ((int_value as f64) == value).then_some(int_value) +} + +fn primary_key_series_ids(batch: &RecordBatch) -> Result<(Vec, usize)> { + let primary_key_column = batch.column(primary_key_column_index(batch.num_columns())); + if let Some(dict) = primary_key_column + .as_any() + .downcast_ref::>() + { + let values = dict + .values() + .as_any() + .downcast_ref::() + .with_context(|| InvalidRecordBatchSnafu { + reason: "primary key dictionary values are not binary".to_string(), + })?; + ensure!( + dict.null_count() == 0 && values.null_count() == 0, + InvalidRecordBatchSnafu { + reason: "primary key dictionary contains null".to_string(), + } + ); + let series_ids = dict + .keys() + .values() + .iter() + .map(|key| *key as usize) + .collect::>(); + ensure!( + series_ids.iter().all(|series_id| *series_id < values.len()), + InvalidRecordBatchSnafu { + reason: "primary key dictionary key is out of bounds".to_string(), + } + ); + return Ok((series_ids, values.len())); + } + + let binary = primary_key_column + .as_any() + .downcast_ref::() + .with_context(|| InvalidRecordBatchSnafu { + reason: format!( + "primary key column is not dictionary or binary, got {:?}", + primary_key_column.data_type() + ), + })?; + ensure!( + binary.null_count() == 0, + InvalidRecordBatchSnafu { + reason: "primary key binary column contains null".to_string(), + } + ); + let mut series_by_key = HashMap::<&[u8], usize>::new(); + let mut series_ids = Vec::with_capacity(binary.len()); + for key in binary.iter().flatten() { + let series_id = match series_by_key.get(key) { + Some(series_id) => *series_id, + None => { + let series_id = series_by_key.len(); + series_by_key.insert(key, series_id); + series_id + } + }; + series_ids.push(series_id); + } + Ok((series_ids, series_by_key.len())) +} + +#[cfg(test)] +mod tests { + use datatypes::arrow::array::{ + BinaryArray, DictionaryArray, Float64Array, Int64Array, TimestampMillisecondArray, + UInt8Array, UInt32Array, UInt64Array, + }; + use datatypes::schema::ColumnSchema; + use store_api::metadata::{ColumnMetadata, RegionMetadataBuilder}; + use store_api::storage::RegionId; + use store_api::storage::consts::{ + OP_TYPE_COLUMN_NAME, PRIMARY_KEY_COLUMN_NAME, SEQUENCE_COLUMN_NAME, + }; + + use super::*; + + fn column_metadata( + column_id: u32, + semantic_type: SemanticType, + name: &str, + data_type: ConcreteDataType, + ) -> ColumnMetadata { + ColumnMetadata { + column_id, + semantic_type, + column_schema: ColumnSchema::new(name, data_type, true), + } + } + + #[test] + fn test_split_requires_metric_region() { + let value_int_name = metric_engine_value_int_column_name("value"); + let mut builder = RegionMetadataBuilder::new(RegionId::new(1, 1)); + builder + .push_column_metadata(column_metadata( + 0, + SemanticType::Timestamp, + "ts", + ConcreteDataType::timestamp_millisecond_datatype(), + )) + .push_column_metadata(column_metadata( + 1, + SemanticType::Field, + "value", + ConcreteDataType::float64_datatype(), + )) + .push_column_metadata(column_metadata( + 2, + SemanticType::Field, + &value_int_name, + ConcreteDataType::int64_datatype(), + )); + let metadata = builder.build_without_validation().unwrap(); + + assert!(!is_metric_engine_data_region(&metadata)); + } + + #[test] + fn test_split_metric_value_columns_by_series() { + let batch = RecordBatch::try_from_iter_with_nullable([ + ( + "greptime_value", + Arc::new(Float64Array::from(vec![ + Some(1.0), + Some(2.0), + Some(1.5), + Some(2.0), + ])) as ArrayRef, + true, + ), + ( + "greptime_value__metric_int", + Arc::new(Int64Array::from(vec![None, None, None, None])) as ArrayRef, + true, + ), + ( + "greptime_timestamp", + Arc::new(TimestampMillisecondArray::from(vec![0, 1, 0, 1])) as ArrayRef, + false, + ), + ( + PRIMARY_KEY_COLUMN_NAME, + Arc::new(DictionaryArray::::new( + UInt32Array::from(vec![0, 0, 1, 1]), + Arc::new(BinaryArray::from_iter_values([b"a", b"b"])), + )) as ArrayRef, + false, + ), + ( + SEQUENCE_COLUMN_NAME, + Arc::new(UInt64Array::from(vec![1, 2, 3, 4])) as ArrayRef, + false, + ), + ( + OP_TYPE_COLUMN_NAME, + Arc::new(UInt8Array::from(vec![0, 0, 0, 0])) as ArrayRef, + false, + ), + ]) + .unwrap(); + + let batch = split_metric_value_columns( + &batch, + &[MetricValueSplitColumn { + float_index: 0, + int_index: 1, + }], + ) + .unwrap(); + + let float_values = batch + .column(0) + .as_any() + .downcast_ref::() + .unwrap(); + let int_values = batch + .column(1) + .as_any() + .downcast_ref::() + .unwrap(); + + assert_eq!( + float_values.iter().collect::>(), + vec![None, None, Some(1.5), Some(2.0)] + ); + assert_eq!( + int_values.iter().collect::>(), + vec![Some(1), Some(2), None, None] + ); + } + + #[test] + fn test_integer_value_classification() { + assert_eq!(integer_value(1.0), Some(1)); + assert_eq!(integer_value(-1.0), Some(-1)); + assert_eq!(integer_value(-0.0), Some(0)); + assert_eq!(integer_value(i64::MIN as f64), Some(i64::MIN)); + + assert_eq!(integer_value(1.5), None); + assert_eq!(integer_value(f64::NAN), None); + assert_eq!(integer_value(f64::INFINITY), None); + assert_eq!(integer_value(f64::NEG_INFINITY), None); + assert_eq!(integer_value(i64::MAX as f64), None); + } +} diff --git a/src/mito2/src/sst/parquet/writer.rs b/src/mito2/src/sst/parquet/writer.rs index e97965d8e7..9c5135939f 100644 --- a/src/mito2/src/sst/parquet/writer.rs +++ b/src/mito2/src/sst/parquet/writer.rs @@ -57,6 +57,7 @@ use crate::sst::file::RegionFileId; use crate::sst::index::{IndexOutput, Indexer, IndexerBuilder}; use crate::sst::parquet::flat_format::{FlatWriteFormat, time_index_column_index}; use crate::sst::parquet::format::PrimaryKeyWriteFormat; +use crate::sst::parquet::metric_value_split::metric_value_split_columns; use crate::sst::parquet::{PARQUET_METADATA_KEY, SstInfo, WriteOptions}; use crate::sst::{ DEFAULT_WRITE_BUFFER_SIZE, DEFAULT_WRITE_CONCURRENCY, FlatSchemaOptions, SeriesEstimator, @@ -381,6 +382,7 @@ where fn customize_column_config( builder: WriterPropertiesBuilder, region_metadata: &RegionMetadataRef, + schema: &SchemaRef, ) -> WriterPropertiesBuilder { let ts_col = ColumnPath::new(vec![ region_metadata @@ -392,12 +394,22 @@ where let seq_col = ColumnPath::new(vec![SEQUENCE_COLUMN_NAME.to_string()]); let op_type_col = ColumnPath::new(vec![OP_TYPE_COLUMN_NAME.to_string()]); - builder + let builder = builder .set_column_encoding(seq_col.clone(), Encoding::DELTA_BINARY_PACKED) .set_column_dictionary_enabled(seq_col, false) .set_column_encoding(ts_col.clone(), Encoding::DELTA_BINARY_PACKED) .set_column_dictionary_enabled(ts_col, false) - .set_column_compression(op_type_col, Compression::UNCOMPRESSED) + .set_column_compression(op_type_col, Compression::UNCOMPRESSED); + + metric_value_split_columns(region_metadata, schema) + .into_iter() + .fold(builder, |builder, split_column| { + let int_col = + ColumnPath::new(vec![schema.field(split_column.int_index).name().clone()]); + builder + .set_column_encoding(int_col.clone(), Encoding::DELTA_BINARY_PACKED) + .set_column_dictionary_enabled(int_col, false) + }) } async fn write_next_flat_batch( @@ -445,7 +457,8 @@ where .set_column_index_truncate_length(None) .set_statistics_truncate_length(None); - let props_builder = Self::customize_column_config(props_builder, &self.metadata); + let props_builder = + Self::customize_column_config(props_builder, &self.metadata, schema); let writer_props = props_builder.build(); let sst_file_path = self.path_provider.build_sst_file_path(RegionFileId::new( diff --git a/src/store-api/src/metric_engine_consts.rs b/src/store-api/src/metric_engine_consts.rs index 16d7a67773..fb84fc4498 100644 --- a/src/store-api/src/metric_engine_consts.rs +++ b/src/store-api/src/metric_engine_consts.rs @@ -33,6 +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 METADATA_REGION_SUBDIR: &str = "metadata"; pub const DATA_REGION_SUBDIR: &str = "data"; @@ -83,7 +84,19 @@ pub const MANIFEST_INFO_EXTENSION_KEY: &str = "MANIFEST_INFO"; /// Returns true if it's a internal column of the metric engine. pub fn is_metric_engine_internal_column(name: &str) -> bool { - name == DATA_SCHEMA_TABLE_ID_COLUMN_NAME || name == DATA_SCHEMA_TSID_COLUMN_NAME + name == DATA_SCHEMA_TABLE_ID_COLUMN_NAME + || name == DATA_SCHEMA_TSID_COLUMN_NAME + || is_metric_engine_value_int_column(name) +} + +/// 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}") +} + +/// 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) } /// Returns true if it's metric engine