feat(mito2): split integral metric values in parquet

Signed-off-by: Ruihang Xia <waynestxia@gmail.com>
This commit is contained in:
Ruihang Xia
2026-07-10 13:09:14 +08:00
parent e6472fd12a
commit d28bca855f
6 changed files with 467 additions and 10 deletions
+16 -5
View File
@@ -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<FlatSources> {
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());
+1
View File
@@ -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;
+9 -1
View File
@@ -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<SequenceNumber>,
metric_value_split_columns: Vec<MetricValueSplitColumn>,
}
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<RecordBatch> {
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();
@@ -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<MetricValueSplitColumn> {
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<RecordBatch> {
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::<Float64Array>()
.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::<Int64Array>()
.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<bool> {
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<f64> {
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<i64> {
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>, 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::<DictionaryArray<UInt32Type>>()
{
let values = dict
.values()
.as_any()
.downcast_ref::<BinaryArray>()
.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::<Vec<_>>();
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::<BinaryArray>()
.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::<UInt32Type>::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::<Float64Array>()
.unwrap();
let int_values = batch
.column(1)
.as_any()
.downcast_ref::<Int64Array>()
.unwrap();
assert_eq!(
float_values.iter().collect::<Vec<_>>(),
vec![None, None, Some(1.5), Some(2.0)]
);
assert_eq!(
int_values.iter().collect::<Vec<_>>(),
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);
}
}
+16 -3
View File
@@ -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(
+14 -1
View File
@@ -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