diff --git a/src/api/src/helper.rs b/src/api/src/helper.rs index 40f9bdfb6d6..c23ae2587d5 100644 --- a/src/api/src/helper.rs +++ b/src/api/src/helper.rs @@ -87,6 +87,28 @@ impl ColumnDataTypeWrapper { } } +/// Returns the time unit if `datatype` is a timestamp type. +pub fn timestamp_unit(datatype: ColumnDataType) -> Option { + match datatype { + ColumnDataType::TimestampSecond => Some(TimeUnit::Second), + ColumnDataType::TimestampMillisecond => Some(TimeUnit::Millisecond), + ColumnDataType::TimestampMicrosecond => Some(TimeUnit::Microsecond), + ColumnDataType::TimestampNanosecond => Some(TimeUnit::Nanosecond), + _ => None, + } +} + +/// Returns the timestamp [ColumnDataType] for the given time unit. +/// This is the inverse of [timestamp_unit]. +pub fn timestamp_datatype(unit: TimeUnit) -> ColumnDataType { + match unit { + TimeUnit::Second => ColumnDataType::TimestampSecond, + TimeUnit::Millisecond => ColumnDataType::TimestampMillisecond, + TimeUnit::Microsecond => ColumnDataType::TimestampMicrosecond, + TimeUnit::Nanosecond => ColumnDataType::TimestampNanosecond, + } +} + impl From for ConcreteDataType { fn from(datatype_wrapper: ColumnDataTypeWrapper) -> Self { match datatype_wrapper.datatype { @@ -1215,6 +1237,21 @@ mod tests { use super::*; use crate::v1::Column; + #[test] + fn test_timestamp_unit_roundtrip() { + for unit in [ + TimeUnit::Second, + TimeUnit::Millisecond, + TimeUnit::Microsecond, + TimeUnit::Nanosecond, + ] { + assert_eq!(timestamp_unit(timestamp_datatype(unit)), Some(unit)); + } + // Non-timestamp types have no time unit. + assert_eq!(timestamp_unit(ColumnDataType::String), None); + assert_eq!(timestamp_unit(ColumnDataType::Datetime), None); + } + #[test] fn test_values_with_capacity() { let values = values_with_capacity(ColumnDataType::Int8, 2); diff --git a/src/frontend/src/instance/otlp.rs b/src/frontend/src/instance/otlp.rs index 7a0f706c58d..e82aa343870 100644 --- a/src/frontend/src/instance/otlp.rs +++ b/src/frontend/src/instance/otlp.rs @@ -124,7 +124,7 @@ impl OpenTelemetryProtocolHandler for Instance { metric_ctx.resource_info = self.otlp_resource_info; let otlp::metrics::MetricsConversion { - requests, + mut requests, rows, semantic_index, resource_info, @@ -173,15 +173,30 @@ impl OpenTelemetryProtocolHandler for Instance { let batcher = self.logical_batcher().filter(|_| { ctx.logical_batching_enabled() && !metric_ctx.is_legacy && metric_ctx.with_metric_engine }); - let batcher = if batcher.is_some() - && self + let batcher = if batcher.is_some() { + // Align the requests' time index unit with the physical table's + // before the bulk eligibility check: the OTLP encoder keeps + // nanosecond precision on the metric engine path, while the bulk + // path assumes the physical table's unit (today millisecond). + // Without this, nanosecond requests would silently skip the + // batcher, and a non-millisecond physical table must not enter + // the bulk path either (the eligibility check rejects its unit). + self.inserter + .align_metric_row_inserts_time_unit(&ctx, &physical_table, &mut requests) + .await + .map_err(BoxedError::new) + .context(error::ExecuteGrpcQuerySnafu)?; + if self .inserter .can_batch_metric_rows(&requests, &ctx, &physical_table) .await .map_err(BoxedError::new) .context(error::ExecuteGrpcQuerySnafu)? - { - batcher + { + batcher + } else { + None + } } else { None }; diff --git a/src/metric-engine/src/engine/put.rs b/src/metric-engine/src/engine/put.rs index 56d4af2284f..fdf9b02a5ff 100644 --- a/src/metric-engine/src/engine/put.rs +++ b/src/metric-engine/src/engine/put.rs @@ -791,7 +791,9 @@ mod tests { use common_function::utils::partition_expr_version; use common_query::prelude::{greptime_native_histogram, greptime_timestamp, greptime_value}; use common_recordbatch::RecordBatches; - use datatypes::arrow::array::{Float64Array, TimestampMillisecondArray}; + use datatypes::arrow::array::{ + Float64Array, TimestampMicrosecondArray, TimestampMillisecondArray, + }; use datatypes::prelude::ConcreteDataType; use datatypes::schema::{ColumnDefaultConstraint, ColumnSchema}; use datatypes::value::Value as PartitionValue; @@ -1116,6 +1118,91 @@ mod tests { expr.as_json_str().unwrap() } + #[tokio::test] + async fn test_put_and_scan_microsecond_physical_region() { + let env = TestEnv::new().await; + let engine = env.metric(); + let physical_region_id = env.default_physical_region_id(); + let logical_region_id = env.default_logical_region_id(); + env.create_physical_region_with_ts_type( + physical_region_id, + &TestEnv::default_table_dir(), + vec![], + ConcreteDataType::timestamp_microsecond_datatype(), + ) + .await; + + // A logical region with a matching microsecond time index is accepted. + let region_create_request = test_util::create_logical_region_request_with_ts_type( + &["job"], + physical_region_id, + &table_dir("test", logical_region_id.table_id()), + ConcreteDataType::timestamp_microsecond_datatype(), + ); + engine + .handle_request( + logical_region_id, + RegionRequest::Create(region_create_request), + ) + .await + .unwrap(); + + // Writing microsecond rows works. + let affected_rows = engine + .handle_request( + logical_region_id, + RegionRequest::Put(RegionPutRequest { + rows: Rows { + schema: test_util::row_schema_with_tags_and_ts_datatype( + &["job"], + ColumnDataType::TimestampMicrosecond, + ), + rows: test_util::build_rows_with_ts_datatype( + 1, + 2, + ColumnDataType::TimestampMicrosecond, + ), + }, + hint: None, + partition_expr_version: None, + skip_wal: false, + }), + ) + .await + .unwrap(); + assert_eq!(affected_rows.affected_rows, 2); + + // The scan returns timestamps in the physical region's unit. + let stream = engine + .scan_to_stream(logical_region_id, ScanRequest::default()) + .await + .unwrap(); + let batches = RecordBatches::try_collect(stream).await.unwrap(); + let mut rows = Vec::new(); + for batch in batches.iter() { + let batch = batch.df_record_batch(); + let timestamps = batch + .column(batch.schema().index_of(greptime_timestamp()).unwrap()) + .as_any() + .downcast_ref::() + .unwrap(); + let values = batch + .column(batch.schema().index_of(greptime_value()).unwrap()) + .as_any() + .downcast_ref::() + .unwrap(); + rows.extend( + timestamps + .values() + .iter() + .copied() + .zip(values.values().iter().copied()), + ); + } + rows.sort_unstable_by_key(|(timestamp, _)| *timestamp); + assert_eq!(rows, vec![(0, 0.0), (1, 1.0)]); + } + async fn create_logical_region_with_tags( env: &TestEnv, physical_region_id: RegionId, diff --git a/src/metric-engine/src/test_util.rs b/src/metric-engine/src/test_util.rs index 1f177370988..0671351ddb3 100644 --- a/src/metric-engine/src/test_util.rs +++ b/src/metric-engine/src/test_util.rs @@ -158,6 +158,24 @@ impl TestEnv { physical_region_id: RegionId, table_dir: &str, options: Vec<(String, String)>, + ) { + self.create_physical_region_with_ts_type( + physical_region_id, + table_dir, + options, + ConcreteDataType::timestamp_millisecond_datatype(), + ) + .await + } + + /// Create regions in [MetricEngine] with specific `physical_region_id` + /// and time index column type. + pub async fn create_physical_region_with_ts_type( + &self, + physical_region_id: RegionId, + table_dir: &str, + options: Vec<(String, String)>, + ts_type: ConcreteDataType, ) { let region_create_request = RegionCreateRequest { engine: METRIC_ENGINE_NAME.to_string(), @@ -165,11 +183,7 @@ impl TestEnv { ColumnMetadata { column_id: 0, semantic_type: SemanticType::Timestamp, - column_schema: ColumnSchema::new( - greptime_timestamp(), - ConcreteDataType::timestamp_millisecond_datatype(), - false, - ), + column_schema: ColumnSchema::new(greptime_timestamp(), ts_type, false), }, ColumnMetadata { column_id: 1, @@ -330,16 +344,28 @@ pub fn create_logical_region_request( tags: &[&str], physical_region_id: RegionId, table_dir: &str, +) -> RegionCreateRequest { + create_logical_region_request_with_ts_type( + tags, + physical_region_id, + table_dir, + ConcreteDataType::timestamp_millisecond_datatype(), + ) +} + +/// Generate a [RegionCreateRequest] for logical region with the given time +/// index column type. +pub fn create_logical_region_request_with_ts_type( + tags: &[&str], + physical_region_id: RegionId, + table_dir: &str, + ts_type: ConcreteDataType, ) -> RegionCreateRequest { let mut column_metadatas = vec![ ColumnMetadata { column_id: 0, semantic_type: SemanticType::Timestamp, - column_schema: ColumnSchema::new( - greptime_timestamp(), - ConcreteDataType::timestamp_millisecond_datatype(), - false, - ), + column_schema: ColumnSchema::new(greptime_timestamp(), ts_type, false), }, ColumnMetadata { column_id: 1, @@ -407,10 +433,20 @@ pub fn alter_logical_region_request(tags: &[&str]) -> RegionAlterRequest { /// /// The result will also contains default timestamp and value column at beginning. pub fn row_schema_with_tags(tags: &[&str]) -> Vec { + row_schema_with_tags_and_ts_datatype(tags, ColumnDataType::TimestampMillisecond) +} + +/// Generate a row schema with given tag columns and time index datatype. +/// +/// The result will also contains default timestamp and value column at beginning. +pub fn row_schema_with_tags_and_ts_datatype( + tags: &[&str], + ts_datatype: ColumnDataType, +) -> Vec { let mut schema = vec![ PbColumnSchema { column_name: greptime_timestamp().to_string(), - datatype: ColumnDataType::TimestampMillisecond as i32, + datatype: ts_datatype as i32, semantic_type: SemanticType::Timestamp as _, datatype_extension: None, options: None, @@ -440,11 +476,33 @@ pub fn row_schema_with_tags(tags: &[&str]) -> Vec { /// The schema is generated by [row_schema_with_tags]. `num_tags` doesn't need to be precise, /// it's used to determine the column id for new columns. pub fn build_rows(num_tags: usize, num_rows: usize) -> Vec { + build_rows_with_ts_datatype(num_tags, num_rows, ColumnDataType::TimestampMillisecond) +} + +/// Build [Row]s whose time index values use the given timestamp datatype. +/// +/// The schema is generated by [row_schema_with_tags_and_ts_datatype]. +/// `num_tags` doesn't need to be precise, it's used to determine the column id +/// for new columns. +pub fn build_rows_with_ts_datatype( + num_tags: usize, + num_rows: usize, + ts_datatype: ColumnDataType, +) -> Vec { + let ts_value = |i: usize| -> ValueData { + match ts_datatype { + ColumnDataType::TimestampSecond => ValueData::TimestampSecondValue(i as _), + ColumnDataType::TimestampMillisecond => ValueData::TimestampMillisecondValue(i as _), + ColumnDataType::TimestampMicrosecond => ValueData::TimestampMicrosecondValue(i as _), + ColumnDataType::TimestampNanosecond => ValueData::TimestampNanosecondValue(i as _), + _ => unreachable!("not a timestamp datatype"), + } + }; let mut rows = vec![]; for i in 0..num_rows { let mut values = vec![ Value { - value_data: Some(ValueData::TimestampMillisecondValue(i as _)), + value_data: Some(ts_value(i)), }, Value { value_data: Some(ValueData::F64Value(i as f64)), diff --git a/src/operator/src/insert.rs b/src/operator/src/insert.rs index 49215f89713..24acdeaba28 100644 --- a/src/operator/src/insert.rs +++ b/src/operator/src/insert.rs @@ -22,6 +22,7 @@ use api::v1::region::{ InsertRequest as RegionInsertRequest, InsertRequests as RegionInsertRequests, RegionRequestHeader, }; +use api::v1::value::ValueData; use api::v1::{ AlterTableExpr, ColumnDataType, ColumnSchema, CreateTableExpr, InsertRequests, RowInsertRequest, RowInsertRequests, Rows, SemanticType, @@ -46,6 +47,8 @@ use common_query::native_histogram::{is_native_histogram_value_type, native_hist use common_query::prelude::{greptime_timestamp, greptime_value}; use common_telemetry::tracing_context::TracingContext; use common_telemetry::{debug, error, warn}; +use common_time::Timestamp; +use common_time::timestamp::TimeUnit; use datatypes::schema::SkippingIndexOptions; use futures_util::future; use meter_core::data::MeterRecord; @@ -429,6 +432,7 @@ impl Inserter { statement_executor, accommodate_existing_schema, is_single_value, + None, ) .await?; @@ -565,10 +569,21 @@ impl Inserter { validate_column_count_match(&requests)?; // check and create physical table - self.create_physical_table_on_demand(&ctx, physical_table.clone(), statement_executor) + let physical_table_ref = self + .create_physical_table_on_demand(&ctx, physical_table.clone(), statement_executor) .await?; - // check and create logical tables + // check and create logical tables; `create_or_alter_tables_on_demand` + // aligns each request's time index unit with the unit of the table it + // targets, inside its existing table lookups: existing tables keep + // their own unit (which matches the physical table they are bound + // to), and new tables use the selected physical table's unit (from + // `physical_table_ref`). Ingestion endpoints encode timestamps in a + // fixed unit (prometheus remote write always uses millisecond; OTLP + // keeps nanosecond precision on the metric engine path), and the + // metric engine requires each logical table's requests to match its + // time index unit. Narrowing conversions truncate the sub-unit part + // (floor), following `Timestamp::convert_to`. let CreateAlterTableResult { instant_table_ids, table_infos, @@ -580,6 +595,7 @@ impl Inserter { statement_executor, true, true, + table_time_index_unit(&physical_table_ref), ) .await?; let name_to_info = table_infos @@ -913,6 +929,7 @@ impl Inserter { statement_executor, false, false, + None, ) .await?; Ok(()) @@ -930,6 +947,13 @@ impl Inserter { /// custom schema, and then inserts data with endpoints that have default schema setting, like prometheus /// remote write. This will modify the `RowInsertRequests` in place. /// `is_single_value` indicates whether the default schema only contains single value column so we can accommodate it. + /// + /// `align_time_index_unit` is the selected physical metric table's time + /// index unit; passing `Some` (metric engine path only) rewrites each + /// request's time index column, inside this function's existing table + /// lookups (no extra catalog access): existing destination tables are + /// converted to their own unit, and new tables to the given unit. + #[allow(clippy::too_many_arguments)] async fn create_or_alter_tables_on_demand( &self, requests: &mut RowInsertRequests, @@ -938,6 +962,7 @@ impl Inserter { statement_executor: &StatementExecutor, accommodate_existing_schema: bool, is_single_value: bool, + align_time_index_unit: Option, ) -> Result { let _timer = crate::metrics::CREATE_ALTER_ON_DEMAND .with_label_values(&[auto_create_table_type.as_str()]) @@ -958,7 +983,7 @@ impl Inserter { && !has_auto_create_exempt_table { let mut instant_table_ids = HashSet::new(); - for req in &requests.inserts { + for req in &mut requests.inserts { let table = match self.get_table(catalog, &schema, &req.table_name).await? { Some(table) => table, // System-defined table: created canonically by the system, @@ -978,6 +1003,15 @@ impl Inserter { .fail(); } }; + // Metric path: an existing destination table keeps its own + // time index unit (it may be bound to another physical + // table than the one selected by this request). + if align_time_index_unit.is_some() + && let Some(rows) = req.rows.as_mut() + && let Some(target_unit) = table_time_index_unit(&table) + { + convert_rows_time_unit(rows, target_unit)?; + } let table_info = table.table_info(); if matches!(auto_create_table_type, AutoCreateTableType::Trace { .. }) { validate_trace_table_model(&table_info, ctx)?; @@ -1013,6 +1047,15 @@ impl Inserter { if table_info.is_ttl_instant_table() { instant_table_ids.insert(table_info.table_id()); } + // Metric path: an existing destination table keeps its + // own time index unit (it may be bound to another + // physical table than the one selected by this request). + if align_time_index_unit.is_some() + && let Some(rows) = req.rows.as_mut() + && let Some(target_unit) = table_time_index_unit(&table) + { + convert_rows_time_unit(rows, target_unit)?; + } if auto_create_allowed && let Some(alter_expr) = self.get_alter_table_expr_on_demand( req, @@ -1058,6 +1101,14 @@ impl Inserter { .fail(); } None => { + // Metric path: a new table uses the selected physical + // table's unit; convert before the create expression is + // derived from the request schema. + if let Some(physical_unit) = align_time_index_unit + && let Some(rows) = req.rows.as_mut() + { + convert_rows_time_unit(rows, physical_unit)?; + } let semantic_index = per_table_semantics .get_or_insert_with(|| parse_per_table_semantic_index(ctx)) .as_ref(); @@ -1256,17 +1307,16 @@ impl Inserter { ctx: &QueryContextRef, physical_table: String, statement_executor: &StatementExecutor, - ) -> Result<()> { + ) -> Result { let catalog_name = ctx.current_catalog(); let schema_name = ctx.current_schema(); // check if exist - if self + if let Some(table) = self .get_table(catalog_name, &schema_name, &physical_table) .await? - .is_some() { - return Ok(()); + return Ok(table); } // Gate here too, otherwise a disabled switch would still leak the physical table. @@ -1318,7 +1368,7 @@ impl Inserter { .await; match res { - Ok(_) => Ok(()), + Ok(table) => Ok(table), Err(err) => { error!(err; "Failed to create table {table_reference}"); Err(err) @@ -1326,6 +1376,75 @@ impl Inserter { } } + /// Aligns each request's time index unit with the unit of the table the + /// request targets, so that write-path gates (e.g. the logical batcher + /// eligibility check) observe the destination table's unit instead of the + /// ingestion endpoint's encoding unit (Prometheus remote write always + /// uses millisecond; OTLP keeps nanosecond precision on the metric engine + /// path). Existing destination tables keep their own unit — they may be + /// bound to a different physical table than `physical_table` — while new + /// tables use `physical_table`'s unit. A missing physical table is + /// treated as the millisecond unit of the auto-created default; the + /// actual creation (if any) happens later in + /// [`Inserter::handle_metric_row_inserts`]. + /// + /// This is only needed on paths that choose a write path before + /// `handle_metric_row_inserts` (whose own table lookups perform the + /// alignment again, as a no-op after this). + pub async fn align_metric_row_inserts_time_unit( + &self, + ctx: &QueryContextRef, + physical_table: &str, + requests: &mut RowInsertRequests, + ) -> Result<()> { + // The unit conversion indexes rows by the time index position, which + // requires well-formed requests. + validate_column_count_match(requests)?; + let physical_unit = match self + .get_table(ctx.current_catalog(), &ctx.current_schema(), physical_table) + .await? + { + Some(table) => table_time_index_unit(&table), + // A missing physical table is auto-created with the millisecond + // unit later. + None => Some(TimeUnit::Millisecond), + }; + self.align_metric_rows_per_destination(ctx, physical_unit, requests) + .await + } + + /// Converts each request's time index column to the unit of the table the + /// request targets: the existing table's own unit when the table exists, + /// and `physical_unit` (of the request's selected physical table) for new + /// tables. + async fn align_metric_rows_per_destination( + &self, + ctx: &QueryContextRef, + physical_unit: Option, + requests: &mut RowInsertRequests, + ) -> Result<()> { + for request in &mut requests.inserts { + let Some(rows) = request.rows.as_mut() else { + continue; + }; + let target_unit = match self + .get_table( + ctx.current_catalog(), + &ctx.current_schema(), + &request.table_name, + ) + .await? + { + Some(table) => table_time_index_unit(&table), + None => physical_unit, + }; + if let Some(target_unit) = target_unit { + convert_rows_time_unit(rows, target_unit)?; + } + } + Ok(()) + } + async fn get_table( &self, catalog: &str, @@ -1592,6 +1711,84 @@ fn request_is_native_histogram(request_schema: &[ColumnSchema]) -> bool { ) } +/// Returns the table's time index unit, if any. A metric table without a +/// timestamp column is left alone by the unit alignment: the metric engine +/// rejects it anyway. +fn table_time_index_unit(table: &TableRef) -> Option { + table + .table_info() + .meta + .schema + .timestamp_column() + .and_then(|col| col.data_type.as_timestamp().map(|ts| ts.unit())) +} + +fn convert_rows_time_unit(rows: &mut Rows, target_unit: TimeUnit) -> Result<()> { + let Some(ts_index) = rows + .schema + .iter() + .position(|col| col.semantic_type == SemanticType::Timestamp as i32) + else { + return Ok(()); + }; + let Some(source_unit) = ColumnDataType::try_from(rows.schema[ts_index].datatype) + .ok() + .and_then(api::helper::timestamp_unit) + else { + return Ok(()); + }; + if source_unit == target_unit { + return Ok(()); + } + + rows.schema[ts_index].datatype = api::helper::timestamp_datatype(target_unit) as i32; + // Timestamp columns never carry a datatype extension. + rows.schema[ts_index].datatype_extension = None; + + // Note: the schema is rewritten before the rows are converted, so an + // overflow error mid-batch leaves this request half-converted. That is + // harmless: the error aborts the whole insert request. + // + // `validate_column_count_match` guarantees every row carries exactly one + // value per schema column, so the time index position is directly in + // bounds; no per-value search is needed. + for row in &mut rows.rows { + debug_assert_eq!(row.values.len(), rows.schema.len()); + let value = &mut row.values[ts_index]; + let Some(value_data) = value.value_data.take() else { + continue; + }; + value.value_data = + convert_timestamp_value_data(value_data, source_unit, target_unit, ts_index)?; + } + Ok(()) +} + +fn convert_timestamp_value_data( + value_data: ValueData, + source_unit: TimeUnit, + target_unit: TimeUnit, + column_index: usize, +) -> Result> { + let timestamp = match value_data { + ValueData::TimestampSecondValue(v) => Timestamp::new_second(v), + ValueData::TimestampMillisecondValue(v) => Timestamp::new_millisecond(v), + ValueData::TimestampMicrosecondValue(v) => Timestamp::new_microsecond(v), + ValueData::TimestampNanosecondValue(v) => Timestamp::new_nanosecond(v), + // Null or non-timestamp value; nothing to convert. + other => return Ok(Some(other)), + }; + let converted = timestamp + .convert_to(target_unit) + .with_context(|| InvalidInsertRequestSnafu { + reason: format!( + "timestamp column {column_index} value {} in unit {source_unit:?} overflows when converting to unit {target_unit:?}", + timestamp.value() + ), + })?; + Ok(api::helper::to_grpc_value(datatypes::value::Value::Timestamp(converted)).value_data) +} + fn table_is_native_histogram(table: &TableRef) -> bool { let mut fields = table.field_columns(); let Some(col) = fields.next() else { @@ -1909,7 +2106,7 @@ mod tests { use api::helper::ColumnDataTypeWrapper; use api::v1::helper::{field_column_schema, time_index_column_schema}; - use api::v1::{RowInsertRequest, Rows, Value}; + use api::v1::{Row, RowInsertRequest, Rows, Value}; use common_catalog::consts::{DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME}; use common_meta::cache::new_table_flownode_set_cache; use common_meta::ddl::test_util::datanode_handler::NaiveDatanodeHandler; @@ -1976,6 +2173,269 @@ mod tests { )) } + fn make_metric_physical_table_ref_with_time_unit(unit: TimeUnit) -> TableRef { + let schema = datatypes::schema::SchemaBuilder::try_from_columns(vec![ + ColumnSchema::new( + greptime_timestamp(), + ConcreteDataType::timestamp_datatype(unit), + false, + ) + .with_time_index(true), + ColumnSchema::new(greptime_value(), ConcreteDataType::float64_datatype(), true), + ]) + .unwrap() + .build() + .unwrap(); + let meta = TableMetaBuilder::empty() + .schema(Arc::new(schema)) + .primary_key_indices(vec![]) + .value_indices(vec![1]) + .engine("metric") + .next_column_id(0) + .options(Default::default()) + .created_on(Default::default()) + .build() + .unwrap(); + let info = Arc::new( + TableInfoBuilder::default() + .table_id(1) + .table_version(0) + .name("greptime_physical_table") + .schema_name(DEFAULT_SCHEMA_NAME) + .catalog_name(DEFAULT_CATALOG_NAME) + .desc(None) + .table_type(TableType::Base) + .meta(meta) + .build() + .unwrap(), + ); + Arc::new(table::Table::new( + info, + table::metadata::FilterPushDownType::Unsupported, + Arc::new(DummyDataSource), + )) + } + + fn ms_row_insert_request(timestamp_ms: i64) -> RowInsertRequest { + ms_row_insert_request_named("my_metric", timestamp_ms) + } + + fn ms_row_insert_request_named(table: &str, timestamp_ms: i64) -> RowInsertRequest { + RowInsertRequest { + table_name: table.to_string(), + rows: Some(Rows { + schema: vec![ + time_index_column_schema( + greptime_timestamp(), + ColumnDataType::TimestampMillisecond, + ), + field_column_schema(greptime_value(), ColumnDataType::Float64), + ], + rows: vec![Row { + values: vec![ + Value { + value_data: Some(ValueData::TimestampMillisecondValue(timestamp_ms)), + }, + Value { + value_data: Some(ValueData::F64Value(1.0)), + }, + ], + }], + }), + } + } + + /// Converts the time index column of each request to `target_unit`. + fn convert_row_insert_requests_time_unit( + requests: &mut RowInsertRequests, + target_unit: TimeUnit, + ) -> Result<()> { + for request in &mut requests.inserts { + let Some(rows) = request.rows.as_mut() else { + continue; + }; + convert_rows_time_unit(rows, target_unit)?; + } + Ok(()) + } + + #[test] + fn test_convert_row_insert_requests_time_unit_noop_when_matching() { + let mut requests = RowInsertRequests { + inserts: vec![ms_row_insert_request(123)], + }; + convert_row_insert_requests_time_unit(&mut requests, TimeUnit::Millisecond).unwrap(); + let rows = requests.inserts[0].rows.as_ref().unwrap(); + assert_eq!( + rows.schema[0].datatype, + ColumnDataType::TimestampMillisecond as i32 + ); + assert!(matches!( + rows.rows[0].values[0].value_data, + Some(ValueData::TimestampMillisecondValue(123)) + )); + } + + #[test] + fn test_convert_row_insert_requests_time_unit_widens_losslessly() { + let mut requests = RowInsertRequests { + inserts: vec![ms_row_insert_request(123)], + }; + convert_row_insert_requests_time_unit(&mut requests, TimeUnit::Microsecond).unwrap(); + let rows = requests.inserts[0].rows.as_ref().unwrap(); + assert_eq!( + rows.schema[0].datatype, + ColumnDataType::TimestampMicrosecond as i32 + ); + assert!(matches!( + rows.rows[0].values[0].value_data, + Some(ValueData::TimestampMicrosecondValue(123_000)) + )); + } + + #[test] + fn test_convert_row_insert_requests_time_unit_truncates_on_narrowing() { + // 123_456_789 ns floors to 123_456 us and 123 ms; negative values + // floor towards negative infinity, matching `Timestamp::convert_to`. + let requests = |value_ns: i64| RowInsertRequests { + inserts: vec![RowInsertRequest { + table_name: "my_metric".to_string(), + rows: Some(Rows { + schema: vec![ + time_index_column_schema( + greptime_timestamp(), + ColumnDataType::TimestampNanosecond, + ), + field_column_schema(greptime_value(), ColumnDataType::Float64), + ], + rows: vec![Row { + values: vec![ + Value { + value_data: Some(ValueData::TimestampNanosecondValue(value_ns)), + }, + Value { + value_data: Some(ValueData::F64Value(1.0)), + }, + ], + }], + }), + }], + }; + + let mut reqs = requests(123_456_789); + convert_row_insert_requests_time_unit(&mut reqs, TimeUnit::Microsecond).unwrap(); + let rows = reqs.inserts[0].rows.as_ref().unwrap(); + assert!(matches!( + rows.rows[0].values[0].value_data, + Some(ValueData::TimestampMicrosecondValue(123_456)) + )); + + let mut reqs = requests(123_456_789); + convert_row_insert_requests_time_unit(&mut reqs, TimeUnit::Millisecond).unwrap(); + let rows = reqs.inserts[0].rows.as_ref().unwrap(); + assert!(matches!( + rows.rows[0].values[0].value_data, + Some(ValueData::TimestampMillisecondValue(123)) + )); + + let mut reqs = requests(-123_456_789); + convert_row_insert_requests_time_unit(&mut reqs, TimeUnit::Microsecond).unwrap(); + let rows = reqs.inserts[0].rows.as_ref().unwrap(); + assert!(matches!( + rows.rows[0].values[0].value_data, + Some(ValueData::TimestampMicrosecondValue(-123_457)) + )); + } + + #[test] + fn test_convert_row_insert_requests_time_unit_overflow() { + let mut requests = RowInsertRequests { + inserts: vec![ms_row_insert_request(i64::MAX)], + }; + let err = + convert_row_insert_requests_time_unit(&mut requests, TimeUnit::Nanosecond).unwrap_err(); + assert!(err.to_string().contains("overflows"), "{err}"); + } + + #[test] + fn test_table_time_index_unit() { + assert_eq!( + table_time_index_unit(&make_metric_physical_table_ref_with_time_unit( + TimeUnit::Microsecond + )), + Some(TimeUnit::Microsecond) + ); + assert_eq!( + table_time_index_unit(&make_metric_physical_table_ref_with_time_unit( + TimeUnit::Millisecond + )), + Some(TimeUnit::Millisecond) + ); + } + + #[tokio::test] + async fn test_align_metric_rows_per_destination() { + use catalog::RegisterTableRequest; + use catalog::memory::MemoryCatalogManager; + + // An existing millisecond logical table `existing`, plus a + // microsecond physical table selected by the request: the existing + // table's request keeps the millisecond unit (it is bound to another + // physical table), while the new table's request is converted to the + // selected physical table's microsecond unit. + let mut inserter = batcher_test_inserter().await; + let existing = make_metric_physical_table_ref_with_time_unit(TimeUnit::Millisecond); + let phy_us = make_metric_physical_table_ref_with_time_unit(TimeUnit::Microsecond); + let catalog = MemoryCatalogManager::with_default_setup(); + for (table_name, table_id, table) in [("existing", 1, existing), ("phy_us", 2, phy_us)] { + catalog + .register_table_sync(RegisterTableRequest { + catalog: DEFAULT_CATALOG_NAME.to_string(), + schema: DEFAULT_SCHEMA_NAME.to_string(), + table_name: table_name.to_string(), + table_id, + table, + }) + .unwrap(); + } + inserter.catalog_manager = catalog; + + let ctx = Arc::new(QueryContext::with( + DEFAULT_CATALOG_NAME, + DEFAULT_SCHEMA_NAME, + )); + let mut requests = RowInsertRequests { + inserts: vec![ + ms_row_insert_request_named("existing", 123), + ms_row_insert_request_named("fresh", 123), + ], + }; + inserter + .align_metric_row_inserts_time_unit(&ctx, "phy_us", &mut requests) + .await + .unwrap(); + + let existing_rows = requests.inserts[0].rows.as_ref().unwrap(); + assert_eq!( + existing_rows.schema[0].datatype, + ColumnDataType::TimestampMillisecond as i32 + ); + assert!(matches!( + existing_rows.rows[0].values[0].value_data, + Some(ValueData::TimestampMillisecondValue(123)) + )); + + let fresh_rows = requests.inserts[1].rows.as_ref().unwrap(); + assert_eq!( + fresh_rows.schema[0].datatype, + ColumnDataType::TimestampMicrosecond as i32 + ); + assert!(matches!( + fresh_rows.rows[0].values[0].value_data, + Some(ValueData::TimestampMicrosecondValue(123_000)) + )); + } + #[tokio::test] async fn test_accommodate_existing_schema_and_reject_kind_changes() { let ts_name = "my_ts"; diff --git a/src/servers/src/batcher/logical_table.rs b/src/servers/src/batcher/logical_table.rs index bd7a26da613..fcd4e1ea38f 100644 --- a/src/servers/src/batcher/logical_table.rs +++ b/src/servers/src/batcher/logical_table.rs @@ -21,6 +21,7 @@ mod tables; #[cfg(test)] mod test_util; +use std::collections::HashSet; use std::future::{Future, ready}; use std::num::NonZeroUsize; use std::sync::Arc; @@ -35,6 +36,7 @@ use common_batcher::worker_registry::WorkerRegistry; use common_meta::cache::TableFlownodeSetCacheRef; use common_meta::node_manager::NodeManagerRef; use common_query::prelude::GREPTIME_PHYSICAL_TABLE; +use common_time::timestamp::TimeUnit; use meter_core::data::MeterRecord; use meter_macros::write_meter; use partition::manager::PartitionRuleManagerRef; @@ -67,15 +69,15 @@ const PHYSICAL_TABLE_KEY: &str = "physical_table"; const WORKER_IDLE_TIMEOUT_MULTIPLIER: u32 = 3; #[derive(Debug, Clone, Hash, Eq, PartialEq)] -struct BatchKey { - catalog: String, - schema: String, - physical_table: String, - skip_wal: bool, +pub(crate) struct BatchKey { + pub(crate) catalog: String, + pub(crate) schema: String, + pub(crate) physical_table: String, + pub(crate) skip_wal: bool, } // Requests can share a batch only when their write target and WAL policy match. -fn batch_key_from_ctx(ctx: &QueryContextRef) -> BatchKey { +pub(crate) fn batch_key_from_ctx(ctx: &QueryContextRef) -> BatchKey { let physical_table = ctx .extension(PHYSICAL_TABLE_KEY) .unwrap_or(GREPTIME_PHYSICAL_TABLE) @@ -171,6 +173,90 @@ impl LogicalTablePendingRowsBatcher { .map(|(rows, ())| rows) } + /// Returns whether the bulk path can accept `batches`: + /// - the physical metric table selected by each context must have a + /// millisecond time index (or not exist yet — it is auto-created with + /// the millisecond unit): new tables are created on it and the bulk + /// encode produces millisecond batches; + /// - every existing destination table must itself have a millisecond time + /// index: it may be bound to another physical table with another unit, + /// and the bulk encoder would build millisecond arrays against its + /// schema. + /// + /// Incompatible requests must stay on the ordinary insert path, which + /// converts the requests to each destination table's unit. Write targets + /// and destination tables are resolved once per distinct name. + pub(crate) async fn accepts_bulk_time_indexes( + &self, + batches: impl Iterator, + ) -> bool { + let mut checked_targets = HashSet::new(); + let mut checked_tables = HashSet::new(); + for (ctx, requests) in batches { + let key = batch_key_from_ctx(ctx); + if checked_targets.insert((key.catalog, key.schema, key.physical_table)) + && !self.accepts_physical_table_time_index(ctx).await + { + return false; + } + for request in &requests.inserts { + if !checked_tables.insert((ctx.current_schema(), request.table_name.clone())) { + continue; + } + let Ok(Some(table)) = self + .catalog_manager + .table( + ctx.current_catalog(), + &ctx.current_schema(), + &request.table_name, + None, + ) + .await + else { + // New table: governed by the selected physical table's + // unit, checked above. + continue; + }; + if table + .table_info() + .meta + .schema + .timestamp_column() + .and_then(|col| col.data_type.as_timestamp().map(|ts| ts.unit())) + .is_none_or(|unit| unit != TimeUnit::Millisecond) + { + return false; + } + } + } + true + } + + /// Returns whether the physical metric table resolved from `ctx` can use + /// the bulk path. The bulk encode produces millisecond timestamp batches + /// only, so a physical table with another time index unit (e.g. created + /// as `TIMESTAMP(6)`) must stay on the ordinary insert path, which + /// converts the requests to the physical table's unit. A missing physical + /// table is accepted: it is auto-created with the millisecond unit. + async fn accepts_physical_table_time_index(&self, ctx: &QueryContextRef) -> bool { + let key = batch_key_from_ctx(ctx); + let Ok(Some(table)) = self + .catalog_manager + .table(&key.catalog, &key.schema, &key.physical_table, None) + .await + else { + return true; + }; + table + .table_info() + .meta + .schema + .timestamp_column() + .and_then(|col| col.data_type.as_timestamp().map(|ts| ts.unit())) + .map(|unit| unit == TimeUnit::Millisecond) + .unwrap_or(true) + } + /// Submits with request-level accounting after schema preparation and before /// queue admission. Acknowledgement follows the global batching policy. pub async fn submit_with( diff --git a/src/servers/src/http/prom_store.rs b/src/servers/src/http/prom_store.rs index 2c4b6927451..44abb18a5be 100644 --- a/src/servers/src/http/prom_store.rs +++ b/src/servers/src/http/prom_store.rs @@ -397,27 +397,45 @@ async fn write_prometheus_rows_with_progress( mut batches: Vec, ) -> std::result::Result { if prom_store_with_metric_engine && let Some(batcher) = pending_rows_batcher { + // Preflight before the bulk eligibility decision: pre_write hooks + // may redirect the contexts (e.g. to a per-tenant schema), so + // eligibility must be evaluated against the prepared batches — + // and the fallback below reuses them without re-running hooks or + // admission. preflight_prometheus_rows(&prom_store_handler, &mut batches) .await .map_err(|error| PromWriteError { error, rows_written: 0, })?; - let mut rows_written = 0; - for (temp_ctx, reqs) in batches { - let rows = batcher - .submit(reqs, temp_ctx) - .await - .map_err(|error| PromWriteError { - error, - rows_written, - })?; - rows_written += rows; + // The bulk encode produces millisecond batches only; write targets + // and existing destination tables with another time index unit stay + // on the ordinary insert path, which converts the requests to each + // destination table's unit. + if batcher.accepts_bulk_time_indexes(batches.iter()).await { + let mut rows_written = 0; + for (temp_ctx, reqs) in batches { + let rows = + batcher + .submit(reqs, temp_ctx) + .await + .map_err(|error| PromWriteError { + error, + rows_written, + })?; + rows_written += rows; + } + return Ok(PromWriteOutcome { + write_cost: 0, + rows_written, + }); } - return Ok(PromWriteOutcome { - write_cost: 0, - rows_written, - }); + return write_prepared_prometheus_rows_with_progress( + prom_store_handler, + batches, + prom_store_with_metric_engine, + ) + .await; } let row_counts = batches @@ -456,6 +474,33 @@ async fn write_prometheus_rows_with_progress( }) } +/// Writes already-preflighted batches through the ordinary (prepared) write +/// path; the pre_write hooks and admission have already run. +async fn write_prepared_prometheus_rows_with_progress( + prom_store_handler: PromStoreProtocolHandlerRef, + batches: Vec, + prom_store_with_metric_engine: bool, +) -> std::result::Result { + let mut write_cost = 0; + let mut rows_written = 0; + for (ctx, request) in batches { + let rows = prom_write_row_count(&request); + let output = prom_store_handler + .write_prepared(request, ctx, prom_store_with_metric_engine) + .await + .map_err(|error| PromWriteError { + error, + rows_written, + })?; + write_cost += output.meta.cost; + rows_written += rows; + } + Ok(PromWriteOutcome { + write_cost, + rows_written, + }) +} + async fn write_prometheus_v2_rows_with_progress( prom_store_handler: PromStoreProtocolHandlerRef, pending_rows_batcher: Option>, @@ -483,20 +528,44 @@ async fn write_prometheus_v2_rows_with_progress( }); } + let sample_batch_count = sample_batches.len(); + let mut batches = sample_batches; + batches.extend(histogram_batches); + if prom_store_with_metric_engine && let Some(batcher) = pending_rows_batcher { - return write_batched_prometheus_v2_rows_with_progress( + // Same ordering as the v1 path: preflight (which may redirect the + // contexts) before the bulk eligibility decision, and the fallback + // reuses the prepared batches without re-running hooks or admission. + preflight_prometheus_rows(&prom_store_handler, &mut batches) + .await + .map_err(|error| PromWriteV2Error { + error, + samples_written: 0, + histograms_written: 0, + })?; + // The bulk encode produces millisecond batches only; write targets + // and existing destination tables with another time index unit stay + // on the ordinary insert path, which converts the requests to each + // destination table's unit. + if batcher.accepts_bulk_time_indexes(batches.iter()).await { + return write_batched_prometheus_v2_rows_with_progress( + prom_store_handler, + batcher.as_ref(), + prom_store_with_metric_engine, + sample_batch_count, + batches, + ) + .await; + } + return write_prepared_prometheus_v2_rows_with_progress( prom_store_handler, - batcher.as_ref(), + batches, + sample_batch_count, prom_store_with_metric_engine, - sample_batches, - histogram_batches, ) .await; } - let sample_batch_count = sample_batches.len(); - let mut batches = sample_batches; - batches.extend(histogram_batches); let row_counts = batches .iter() .map(|(_, request)| prom_write_row_count(request)) @@ -544,24 +613,49 @@ async fn write_prometheus_v2_rows_with_progress( }) } +/// Writes already-preflighted batches through the ordinary (prepared) write +/// path with v2 partial-progress accounting; the pre_write hooks and +/// admission have already run. +async fn write_prepared_prometheus_v2_rows_with_progress( + prom_store_handler: PromStoreProtocolHandlerRef, + batches: Vec, + sample_batch_count: usize, + prom_store_with_metric_engine: bool, +) -> std::result::Result { + let mut write_cost = 0; + let mut samples_written = 0; + let mut histograms_written = 0; + for (index, (ctx, request)) in batches.into_iter().enumerate() { + let rows = prom_write_row_count(&request); + let output = prom_store_handler + .write_prepared(request, ctx, prom_store_with_metric_engine) + .await + .map_err(|error| PromWriteV2Error { + error, + samples_written, + histograms_written, + })?; + write_cost += output.meta.cost; + if index < sample_batch_count { + samples_written += rows; + } else { + histograms_written += rows; + } + } + Ok(PromWriteV2Outcome { + write_cost, + samples_written, + histograms_written, + }) +} + async fn write_batched_prometheus_v2_rows_with_progress( prom_store_handler: PromStoreProtocolHandlerRef, batcher: &B, prom_store_with_metric_engine: bool, - sample_batches: Vec, - histogram_batches: Vec, + sample_batch_count: usize, + batches: Vec, ) -> std::result::Result { - let sample_batch_count = sample_batches.len(); - let mut batches = sample_batches; - batches.extend(histogram_batches); - preflight_prometheus_rows(&prom_store_handler, &mut batches) - .await - .map_err(|error| PromWriteV2Error { - error, - samples_written: 0, - histograms_written: 0, - })?; - let mut samples_written = 0; let mut histograms_written = 0; let mut write_cost = 0; @@ -862,14 +956,18 @@ mod tests { events: events.clone(), }; - let Ok(outcome) = write_batched_prometheus_v2_rows_with_progress( - handler, - &batcher, - true, - vec![test_prom_write_batch("sample")], - vec![test_prom_write_batch("histogram")], - ) - .await + let mut batches = vec![ + test_prom_write_batch("sample"), + test_prom_write_batch("histogram"), + ]; + // The caller preflights before choosing the bulk path; the batched + // writer consumes the prepared batches. + preflight_prometheus_rows(&handler, &mut batches) + .await + .unwrap(); + let Ok(outcome) = + write_batched_prometheus_v2_rows_with_progress(handler, &batcher, true, 1, batches) + .await else { panic!("mixed remote write should succeed") }; diff --git a/src/servers/src/otlp/logs.rs b/src/servers/src/otlp/logs.rs index f32816fbe36..53a2189044d 100644 --- a/src/servers/src/otlp/logs.rs +++ b/src/servers/src/otlp/logs.rs @@ -15,7 +15,7 @@ use std::collections::BTreeMap; use ahash::{HashMap, HashMapExt}; -use api::helper::ColumnDataTypeWrapper; +use api::helper::{ColumnDataTypeWrapper, timestamp_unit}; use api::v1::column_data_type_extension::TypeExt; use api::v1::column_def::options_from_column_schema; use api::v1::value::ValueData; @@ -898,16 +898,6 @@ fn is_timestamp_type(datatype: ColumnDataType) -> bool { timestamp_unit(datatype).is_some() } -fn timestamp_unit(datatype: ColumnDataType) -> Option { - match datatype { - ColumnDataType::TimestampSecond => Some(TimeUnit::Second), - ColumnDataType::TimestampMillisecond => Some(TimeUnit::Millisecond), - ColumnDataType::TimestampMicrosecond => Some(TimeUnit::Microsecond), - ColumnDataType::TimestampNanosecond => Some(TimeUnit::Nanosecond), - _ => None, - } -} - fn parse_export_logs_service_request_to_rows( request: ExportLogsServiceRequest, select_info: Box, diff --git a/src/servers/src/otlp/metrics.rs b/src/servers/src/otlp/metrics.rs index e518d732479..fd496db5a98 100644 --- a/src/servers/src/otlp/metrics.rs +++ b/src/servers/src/otlp/metrics.rs @@ -998,9 +998,15 @@ fn write_timestamp( table: &mut TableData, row: &mut Vec, time_nano: i64, - legacy_mode: bool, + metric_ctx: &OtlpMetricCtx, ) -> Result<()> { - if legacy_mode { + // Keep the full nanosecond precision whenever the request is headed for + // the metric engine: `Inserter::handle_metric_row_inserts` converts the + // timestamps to the physical table's time index unit (which may be + // micro/nanosecond). Only the non-metric prometheus-compatible path is + // fixed to milliseconds, to keep auto-created mito tables on the + // millisecond time index. + if metric_ctx.is_legacy || metric_ctx.with_metric_engine { row_writer::write_ts_to_nanos( table, greptime_timestamp(), @@ -1102,7 +1108,7 @@ fn write_tags_and_timestamp( )?; } - write_timestamp(table, row, timestamp_nanos, metric_ctx.is_legacy)?; + write_timestamp(table, row, timestamp_nanos, metric_ctx)?; Ok(()) } @@ -1541,6 +1547,7 @@ fn encode_summary( #[cfg(test)] mod tests { + use api::v1::ColumnDataType; use common_query::prelude::set_default_prefix; use otel_arrow_rust::proto::opentelemetry::common::v1::AnyValue; use otel_arrow_rust::proto::opentelemetry::common::v1::any_value::Value as Val; @@ -2625,6 +2632,81 @@ mod tests { } } + #[test] + fn test_metric_engine_path_keeps_nanosecond_precision() { + let time_unix_nano = 1_704_067_200_123_456_789u64; + let request = metrics_request(vec![Metric { + name: "my_gauge".to_string(), + data: Some(metric::Data::Gauge(Gauge { + data_points: vec![NumberDataPoint { + time_unix_nano, + value: Some(Value::AsDouble(1.0)), + ..Default::default() + }], + })), + ..Default::default() + }]); + + // The metric engine path keeps the full nanosecond precision here; + // `Inserter::handle_metric_row_inserts` converts the timestamps to + // the physical table's time index unit afterwards. + let mut metric_ctx = OtlpMetricCtx { + with_metric_engine: true, + ..Default::default() + }; + let MetricsConversion { requests, .. } = + to_grpc_insert_requests(request, &mut metric_ctx).unwrap(); + + let rows = requests.inserts[0].rows.as_ref().unwrap(); + let ts_index = rows + .schema + .iter() + .position(|column| column.column_name == greptime_timestamp()) + .unwrap(); + assert_eq!( + rows.schema[ts_index].datatype, + ColumnDataType::TimestampNanosecond as i32 + ); + assert!(matches!( + rows.rows[0].values[ts_index].value_data, + Some(ValueData::TimestampNanosecondValue( + 1_704_067_200_123_456_789 + )) + )); + + // The non-metric prometheus-compatible path stays millisecond so + // auto-created mito tables keep the millisecond time index. + let mut compat_ctx = OtlpMetricCtx::default(); + let request = metrics_request(vec![Metric { + name: "my_gauge".to_string(), + data: Some(metric::Data::Gauge(Gauge { + data_points: vec![NumberDataPoint { + time_unix_nano, + value: Some(Value::AsDouble(1.0)), + ..Default::default() + }], + })), + ..Default::default() + }]); + let MetricsConversion { requests, .. } = + to_grpc_insert_requests(request, &mut compat_ctx).unwrap(); + + let rows = requests.inserts[0].rows.as_ref().unwrap(); + let ts_index = rows + .schema + .iter() + .position(|column| column.column_name == greptime_timestamp()) + .unwrap(); + assert_eq!( + rows.schema[ts_index].datatype, + ColumnDataType::TimestampMillisecond as i32 + ); + assert!(matches!( + rows.rows[0].values[ts_index].value_data, + Some(ValueData::TimestampMillisecondValue(1_704_067_200_123)) + )); + } + #[test] fn test_exponential_histogram_rejection_metrics() { // Other conversion tests update the same process-global counters. diff --git a/src/servers/src/prom_store.rs b/src/servers/src/prom_store.rs index 9182f64ad00..e202070d98d 100644 --- a/src/servers/src/prom_store.rs +++ b/src/servers/src/prom_store.rs @@ -18,14 +18,19 @@ use std::cmp::Ordering; use std::collections::HashMap; use std::collections::hash_map::DefaultHasher; use std::hash::{Hash, Hasher}; +use std::sync::Arc; use api::prom_store::remote::label_matcher::Type as MatcherType; use api::prom_store::remote::{Label, Query, Sample, TimeSeries, WriteRequest}; use api::v1::RowInsertRequests; use arrow::array::{ - Array, AsArray, DictionaryArray, LargeStringArray, StringArray, StringViewArray, + Array, ArrayRef, AsArray, DictionaryArray, LargeStringArray, StringArray, StringViewArray, +}; +use arrow::compute::kernels::cast as casts; +use arrow::datatypes::{ + DataType, Float64Type, TimeUnit, TimestampMicrosecondType, TimestampMillisecondType, + TimestampNanosecondType, UInt32Type, }; -use arrow::datatypes::{Float64Type, TimestampMillisecondType, UInt32Type}; use common_grpc::precision::Precision; use common_query::prelude::{greptime_timestamp, greptime_value}; use common_recordbatch::{RecordBatch, RecordBatches}; @@ -398,6 +403,50 @@ fn append_recordbatch_to_timeseries( .with_context(|| error::InvalidPromRemoteReadQueryResultSnafu { msg: format!("missing timestamp column '{timestamp_column_name}' in query result"), })?; + // The Prometheus remote read wire format carries millisecond timestamps, + // while the table's time index can use any time unit (e.g. a metric + // physical table created with TIMESTAMP(6)). Narrowing floors towards + // negative infinity, consistent with `Timestamp::convert_to` on the + // ingestion path; arrow's cast would truncate towards zero, so a + // pre-epoch sample (e.g. -1001us) would round to -1ms instead of -2ms + // and disagree with what ingesting the same instant into a millisecond + // table would have stored. + let ts_column: ArrayRef = match ts_column.data_type() { + DataType::Timestamp(TimeUnit::Millisecond, _) => ts_column.clone(), + // Second -> millisecond is a widening (exact) conversion. + DataType::Timestamp(TimeUnit::Second, _) => casts::cast( + ts_column, + &DataType::Timestamp(TimeUnit::Millisecond, None), + ) + .map_err(|e| { + error::InvalidPromRemoteReadQueryResultSnafu { + msg: format!( + "failed to cast timestamp column '{timestamp_column_name}' of datatype {:?} to millisecond: {e}", + ts_column.data_type() + ), + } + .build() + })?, + DataType::Timestamp(TimeUnit::Microsecond, _) => Arc::new( + ts_column + .as_primitive::() + .unary::<_, TimestampMillisecondType>(|v| v.div_euclid(1_000)), + ), + DataType::Timestamp(TimeUnit::Nanosecond, _) => Arc::new( + ts_column + .as_primitive::() + .unary::<_, TimestampMillisecondType>(|v| v.div_euclid(1_000_000)), + ), + _ => { + return error::InvalidPromRemoteReadQueryResultSnafu { + msg: format!( + "Expect timestamp column of datatype Timestamp(Millisecond), actual {:?}", + ts_column.data_type() + ), + } + .fail(); + } + }; let ts_column = ts_column .as_primitive_opt::() .with_context(|| error::InvalidPromRemoteReadQueryResultSnafu { diff --git a/tests-integration/src/otlp.rs b/tests-integration/src/otlp.rs index 67a4e512f7d..59ce38919f7 100644 --- a/tests-integration/src/otlp.rs +++ b/tests-integration/src/otlp.rs @@ -19,6 +19,10 @@ mod test { use client::{DEFAULT_CATALOG_NAME, OutputData}; use common_recordbatch::RecordBatches; use datatypes::arrow::array::AsArray; + use datatypes::arrow::datatypes::{ + DataType, TimeUnit, TimestampMicrosecondType, TimestampMillisecondType, + TimestampNanosecondType, TimestampSecondType, + }; use frontend::instance::Instance; use opentelemetry_proto::tonic::collector::trace::v1::ExportTraceServiceRequest; use otel_arrow_rust::proto::opentelemetry::collector::metrics::v1::ExportMetricsServiceRequest; @@ -37,6 +41,7 @@ mod test { use servers::query_handler::OpenTelemetryProtocolHandler; use servers::query_handler::sql::SqlQueryHandler; use session::context::QueryContext; + use session::protocol_ctx::{OtlpMetricCtx, ProtocolCtx}; use crate::standalone::GreptimeDbStandaloneBuilder; use crate::tests; @@ -1164,6 +1169,162 @@ WITH( ); } + #[tokio::test(flavor = "multi_thread")] + pub async fn test_otlp_metrics_into_microsecond_physical_table_on_standalone() { + let standalone = GreptimeDbStandaloneBuilder::new("test_otlp_us_physical") + .build() + .await; + + test_otlp_metrics_into_non_millisecond_physical_table( + standalone.fe_instance(), + "TIMESTAMP(6)", + TimeUnit::Microsecond, + 1_704_067_200_123_456, + ) + .await; + } + + #[tokio::test(flavor = "multi_thread")] + pub async fn test_otlp_metrics_into_microsecond_physical_table_on_distributed() { + let instance = tests::create_distributed_instance("test_otlp_us_physical_dist").await; + + test_otlp_metrics_into_non_millisecond_physical_table( + &instance.frontend(), + "TIMESTAMP(6)", + TimeUnit::Microsecond, + 1_704_067_200_123_456, + ) + .await; + } + + #[tokio::test(flavor = "multi_thread")] + pub async fn test_otlp_metrics_into_seconds_physical_table_on_standalone() { + let standalone = GreptimeDbStandaloneBuilder::new("test_otlp_s_physical") + .build() + .await; + + // Narrowing truncates: 1704067200123456789ns -> 1704067200s. + test_otlp_metrics_into_non_millisecond_physical_table( + standalone.fe_instance(), + "TIMESTAMP(0)", + TimeUnit::Second, + 1_704_067_200, + ) + .await; + } + + #[tokio::test(flavor = "multi_thread")] + pub async fn test_otlp_metrics_into_nanoseconds_physical_table_on_standalone() { + let standalone = GreptimeDbStandaloneBuilder::new("test_otlp_ns_physical") + .build() + .await; + + // Same unit: the full nanosecond precision is kept verbatim. + test_otlp_metrics_into_non_millisecond_physical_table( + standalone.fe_instance(), + "TIMESTAMP(9)", + TimeUnit::Nanosecond, + 1_704_067_200_123_456_789, + ) + .await; + } + + /// Regression test for : + /// a physical metric table pre-created with a non-millisecond time index + /// accepts OTLP ingestion; the nanosecond samples are converted to the + /// physical table's unit (kept verbatim for a nanosecond table, + /// truncating the sub-unit part for micro/second tables). + async fn test_otlp_metrics_into_non_millisecond_physical_table( + instance: &Arc, + sql_ts_type: &str, + expected_unit: TimeUnit, + expected_value: i64, + ) { + let db = "otlp_non_ms_physical"; + let mut ctx = QueryContext::with(DEFAULT_CATALOG_NAME, db); + // Route the request to the metric engine, like the OTLP HTTP handler + // does with the default `prom_store.with_metric_engine = true`. + ctx.set_protocol_ctx(ProtocolCtx::OtlpMetric(OtlpMetricCtx { + with_metric_engine: true, + ..Default::default() + })); + let ctx = Arc::new(ctx); + assert!( + SqlQueryHandler::do_query( + instance.as_ref(), + &format!("CREATE DATABASE IF NOT EXISTS {db}"), + ctx.clone(), + ) + .await + .first() + .unwrap() + .is_ok() + ); + + let mut output = instance + .do_query( + &format!( + "CREATE TABLE greptime_physical_table (\ + greptime_timestamp {sql_ts_type} NOT NULL, \ + greptime_value DOUBLE NULL, \ + TIME INDEX (greptime_timestamp)) \ + ENGINE = metric WITH ('physical_metric_table' = 'true')" + ), + ctx.clone(), + ) + .await; + assert!(output.remove(0).is_ok()); + + let request = ExportMetricsServiceRequest { + resource_metrics: vec![ResourceMetrics { + scope_metrics: vec![ScopeMetrics { + metrics: vec![Metric { + name: "my_gauge".to_string(), + data: Some(metric::Data::Gauge(Gauge { + data_points: vec![NumberDataPoint { + attributes: vec![keyvalue("host", "h1")], + time_unix_nano: 1_704_067_200_123_456_789, + value: Some(Value::AsDouble(1.0)), + ..Default::default() + }], + })), + ..Default::default() + }], + ..Default::default() + }], + ..Default::default() + }], + }; + instance.metrics(request, ctx.clone()).await.unwrap(); + + let mut output = instance + .do_query("SELECT greptime_timestamp FROM my_gauge", ctx.clone()) + .await; + let OutputData::Stream(stream) = output.remove(0).unwrap().data else { + unreachable!() + }; + let batches = RecordBatches::try_collect(stream).await.unwrap().take(); + assert_eq!(batches[0].num_rows(), 1); + // The logical table's time index keeps the physical table's unit. + let ts_column = batches[0].column(0); + assert_eq!( + ts_column.data_type(), + &DataType::Timestamp(expected_unit, None), + "unexpected time index type" + ); + let stored = match expected_unit { + TimeUnit::Second => ts_column.as_primitive::().value(0), + TimeUnit::Millisecond => ts_column + .as_primitive::() + .value(0), + TimeUnit::Microsecond => ts_column + .as_primitive::() + .value(0), + TimeUnit::Nanosecond => ts_column.as_primitive::().value(0), + }; + assert_eq!(stored, expected_value); + } + fn build_request() -> ExportMetricsServiceRequest { let data_points = vec![ NumberDataPoint { diff --git a/tests-integration/src/prom_store.rs b/tests-integration/src/prom_store.rs index 0649b2dd778..af901120454 100644 --- a/tests-integration/src/prom_store.rs +++ b/tests-integration/src/prom_store.rs @@ -18,9 +18,19 @@ mod tests { use api::prom_store::remote::label_matcher::Type as MatcherType; use api::prom_store::remote::{ - Label, LabelMatcher, Query, ReadRequest, ReadResponse, Sample, WriteRequest, + Label, LabelMatcher, Query, ReadRequest, ReadResponse, Sample, TimeSeries, WriteRequest, }; + use api::v1::value::ValueData; + use api::v1::{ + ColumnDataType, Row, RowInsertRequest, RowInsertRequests, Rows, SemanticType, Value, + }; + use client::OutputData; use common_catalog::consts::DEFAULT_CATALOG_NAME; + use datatypes::arrow::array::AsArray; + use datatypes::arrow::datatypes::{ + DataType, TimeUnit, TimestampMicrosecondType, TimestampMillisecondType, + TimestampNanosecondType, TimestampSecondType, + }; use frontend::instance::Instance; use prost::Message; use servers::http::prom_store::PHYSICAL_TABLE_PARAM; @@ -83,6 +93,432 @@ mod tests { .await; } + #[tokio::test(flavor = "multi_thread")] + async fn test_standalone_prom_store_remote_rw_microsecond_physical_table() { + common_telemetry::init_default_ut_logging(); + let standalone = + GreptimeDbStandaloneBuilder::new("test_prom_store_remote_rw_us_physical_table") + .build() + .await; + let instance = standalone.fe_instance(); + + test_prom_store_remote_rw_non_millisecond_physical_table( + instance, + "TIMESTAMP(6)", + TimeUnit::Microsecond, + 1_000_000, + 2_000_000, + ) + .await; + } + + #[tokio::test(flavor = "multi_thread")] + async fn test_distributed_prom_store_remote_rw_microsecond_physical_table() { + common_telemetry::init_default_ut_logging(); + let distributed = + tests::create_distributed_instance("test_prom_store_remote_rw_us_physical_table").await; + test_prom_store_remote_rw_non_millisecond_physical_table( + &distributed.frontend(), + "TIMESTAMP(6)", + TimeUnit::Microsecond, + 1_000_000, + 2_000_000, + ) + .await; + } + + #[tokio::test(flavor = "multi_thread")] + async fn test_standalone_prom_store_remote_rw_seconds_physical_table() { + common_telemetry::init_default_ut_logging(); + let standalone = + GreptimeDbStandaloneBuilder::new("test_prom_store_remote_rw_s_physical_table") + .build() + .await; + let instance = standalone.fe_instance(); + + // Narrowing truncates: 1000ms/2000ms -> 1s/2s. + test_prom_store_remote_rw_non_millisecond_physical_table( + instance, + "TIMESTAMP(0)", + TimeUnit::Second, + 1, + 2, + ) + .await; + } + + #[tokio::test(flavor = "multi_thread")] + async fn test_standalone_prom_store_remote_rw_nanoseconds_physical_table() { + common_telemetry::init_default_ut_logging(); + let standalone = + GreptimeDbStandaloneBuilder::new("test_prom_store_remote_rw_ns_physical_table") + .build() + .await; + let instance = standalone.fe_instance(); + + // Lossless widening: 1000ms/2000ms -> 1e9ns/2e9ns. + test_prom_store_remote_rw_non_millisecond_physical_table( + instance, + "TIMESTAMP(9)", + TimeUnit::Nanosecond, + 1_000_000_000, + 2_000_000_000, + ) + .await; + } + + /// Regression test for : + /// prometheus remote write/read against a physical metric table + /// pre-created with a non-millisecond time index. Millisecond samples are + /// converted to the physical table's unit on write (lossless widening, + /// truncating narrowing) and narrowed back to milliseconds on remote read. + async fn test_prom_store_remote_rw_non_millisecond_physical_table( + instance: &Arc, + sql_ts_type: &str, + expected_unit: TimeUnit, + expected_first: i64, + expected_second: i64, + ) { + let db = "prometheus_non_ms"; + let ctx = Arc::new(QueryContext::with(DEFAULT_CATALOG_NAME, db)); + assert!( + SqlQueryHandler::do_query( + instance.as_ref(), + &format!("CREATE DATABASE IF NOT EXISTS {db}"), + ctx.clone(), + ) + .await + .first() + .unwrap() + .is_ok() + ); + + let mut output = instance + .do_query( + &format!( + "CREATE TABLE greptime_physical_table (\ + greptime_timestamp {sql_ts_type} NOT NULL, \ + greptime_value DOUBLE NULL, \ + TIME INDEX (greptime_timestamp)) \ + ENGINE = metric WITH ('physical_metric_table' = 'true')" + ), + ctx.clone(), + ) + .await; + assert!(output.remove(0).is_ok()); + + let write_request = WriteRequest { + timeseries: vec![prom_store::mock_timeseries()[0].clone()], + ..Default::default() + }; + let (row_inserts, _) = to_grpc_row_insert_requests(&write_request).unwrap(); + instance + .write(row_inserts, ctx.clone(), true) + .await + .unwrap(); + + let read_request = ReadRequest { + queries: vec![Query { + start_timestamp_ms: 1000, + end_timestamp_ms: 2000, + matchers: vec![LabelMatcher { + name: prom_store::METRIC_NAME_LABEL.to_string(), + value: "metric1".to_string(), + r#type: 0, + }], + ..Default::default() + }], + ..Default::default() + }; + let resp = instance.read(read_request, ctx.clone()).await.unwrap(); + let body = prom_store::snappy_decompress(&resp.body).unwrap(); + let read_response = ReadResponse::decode(&body[..]).unwrap(); + assert_eq!(1, read_response.results.len()); + assert_eq!(1, read_response.results[0].timeseries.len()); + let timeseries = &read_response.results[0].timeseries[0]; + assert_eq!( + timeseries.samples, + vec![ + Sample { + value: 1.0, + timestamp: 1000, + }, + Sample { + value: 2.0, + timestamp: 2000, + }, + ] + ); + + // The stored timestamps keep the physical table's time index unit. + let mut output = instance + .do_query("SELECT greptime_timestamp FROM metric1", ctx.clone()) + .await; + let OutputData::Stream(stream) = output.remove(0).unwrap().data else { + unreachable!() + }; + let batches = common_recordbatch::RecordBatches::try_collect(stream) + .await + .unwrap() + .take(); + assert_eq!(batches[0].num_rows(), 2); + let ts_column = batches[0].column(0); + assert_eq!( + ts_column.data_type(), + &DataType::Timestamp(expected_unit, None), + "unexpected time index type" + ); + let stored = |row: usize| match expected_unit { + TimeUnit::Second => ts_column.as_primitive::().value(row), + TimeUnit::Millisecond => ts_column + .as_primitive::() + .value(row), + TimeUnit::Microsecond => ts_column + .as_primitive::() + .value(row), + TimeUnit::Nanosecond => ts_column + .as_primitive::() + .value(row), + }; + assert_eq!(stored(0), expected_first); + assert_eq!(stored(1), expected_second); + + // Negative, non-aligned timestamps must floor towards negative + // infinity on remote read, matching `Timestamp::convert_to` on the + // ingestion path: -1001us/-1001000ns -> -2ms, not -1ms. + if matches!(expected_unit, TimeUnit::Microsecond | TimeUnit::Nanosecond) { + let (ts_datatype, negative_value) = if expected_unit == TimeUnit::Microsecond { + ( + ColumnDataType::TimestampMicrosecond, + ValueData::TimestampMicrosecondValue(-1001), + ) + } else { + ( + ColumnDataType::TimestampNanosecond, + ValueData::TimestampNanosecondValue(-1_001_000), + ) + }; + let negative_request = RowInsertRequests { + inserts: vec![RowInsertRequest { + table_name: "metric1".to_string(), + rows: Some(Rows { + schema: vec![ + api::v1::ColumnSchema { + column_name: "greptime_timestamp".to_string(), + datatype: ts_datatype as i32, + semantic_type: SemanticType::Timestamp as i32, + datatype_extension: None, + options: None, + }, + api::v1::ColumnSchema { + column_name: "greptime_value".to_string(), + datatype: ColumnDataType::Float64 as i32, + semantic_type: SemanticType::Field as i32, + datatype_extension: None, + options: None, + }, + api::v1::ColumnSchema { + column_name: "job".to_string(), + datatype: ColumnDataType::String as i32, + semantic_type: SemanticType::Tag as i32, + datatype_extension: None, + options: None, + }, + ], + rows: vec![Row { + values: vec![ + Value { + value_data: Some(negative_value), + }, + Value { + value_data: Some(ValueData::F64Value(3.0)), + }, + Value { + value_data: Some(ValueData::StringValue("spark".to_string())), + }, + ], + }], + }), + }], + }; + instance + .write(negative_request, ctx.clone(), true) + .await + .unwrap(); + + let read_request = ReadRequest { + queries: vec![Query { + start_timestamp_ms: -2, + end_timestamp_ms: 0, + matchers: vec![LabelMatcher { + name: prom_store::METRIC_NAME_LABEL.to_string(), + value: "metric1".to_string(), + r#type: 0, + }], + ..Default::default() + }], + ..Default::default() + }; + let resp = instance.read(read_request, ctx.clone()).await.unwrap(); + let body = prom_store::snappy_decompress(&resp.body).unwrap(); + let read_response = ReadResponse::decode(&body[..]).unwrap(); + assert_eq!(1, read_response.results.len()); + assert_eq!( + read_response.results[0].timeseries[0].samples, + vec![Sample { + value: 3.0, + timestamp: -2, + }], + "remote read must floor negative non-aligned timestamps" + ); + } + } + + #[tokio::test(flavor = "multi_thread")] + async fn test_standalone_prom_store_write_existing_logical_table_of_other_physical_unit() { + common_telemetry::init_default_ut_logging(); + let standalone = GreptimeDbStandaloneBuilder::new( + "test_prom_store_write_existing_logical_other_physical_unit", + ) + .build() + .await; + let instance = standalone.fe_instance(); + + // Two physical metric tables with different time index units. + let db = "prometheus_mixed_units"; + let ctx = Arc::new(QueryContext::with(DEFAULT_CATALOG_NAME, db)); + assert!( + SqlQueryHandler::do_query( + instance.as_ref(), + &format!("CREATE DATABASE IF NOT EXISTS {db}"), + ctx.clone(), + ) + .await + .first() + .unwrap() + .is_ok() + ); + let mut output = instance + .do_query( + "CREATE TABLE phy_us (greptime_timestamp TIMESTAMP(6) NOT NULL, \ + greptime_value DOUBLE NULL, TIME INDEX (greptime_timestamp)) \ + ENGINE = metric WITH ('physical_metric_table' = 'true')", + ctx.clone(), + ) + .await; + assert!(output.remove(0).is_ok()); + + // First create the logical table `shared_metric` on the default + // (millisecond) physical table. + let shared_series = TimeSeries { + labels: vec![ + Label { + name: prom_store::METRIC_NAME_LABEL.to_string(), + value: "shared_metric".to_string(), + }, + Label { + name: "job".to_string(), + value: "demo".to_string(), + }, + ], + samples: vec![Sample { + value: 1.0, + timestamp: 1000, + }], + ..Default::default() + }; + let (row_inserts, _) = to_grpc_row_insert_requests(&WriteRequest { + timeseries: vec![shared_series.clone()], + ..Default::default() + }) + .unwrap(); + instance + .write(row_inserts, ctx.clone(), true) + .await + .unwrap(); + + // Now select the microsecond physical table while writing to BOTH the + // existing millisecond `shared_metric` (bound to the default physical + // table) and a new table: the existing table's request must keep the + // millisecond unit, and only the new table uses microsecond. + let mut hint_ctx = QueryContext::with(DEFAULT_CATALOG_NAME, db); + hint_ctx.set_extension(PHYSICAL_TABLE_PARAM, "phy_us".to_string()); + let hint_ctx = Arc::new(hint_ctx); + let (row_inserts, _) = to_grpc_row_insert_requests(&WriteRequest { + timeseries: vec![ + TimeSeries { + samples: vec![Sample { + value: 1.5, + timestamp: 1500, + }], + ..shared_series + }, + TimeSeries { + labels: vec![ + Label { + name: prom_store::METRIC_NAME_LABEL.to_string(), + value: "fresh_us_metric".to_string(), + }, + Label { + name: "job".to_string(), + value: "demo".to_string(), + }, + ], + samples: vec![Sample { + value: 2.5, + timestamp: 1500, + }], + ..Default::default() + }, + ], + ..Default::default() + }) + .unwrap(); + instance + .write(row_inserts, hint_ctx.clone(), true) + .await + .unwrap(); + + // The existing table keeps its millisecond unit and precision. + let mut output = instance + .do_query( + "SELECT greptime_timestamp FROM shared_metric ORDER BY greptime_timestamp", + ctx.clone(), + ) + .await; + let OutputData::Stream(stream) = output.remove(0).unwrap().data else { + unreachable!() + }; + let batches = common_recordbatch::RecordBatches::try_collect(stream) + .await + .unwrap() + .take(); + let shared_ts = batches[0] + .column(0) + .as_primitive::(); + assert_eq!((shared_ts.value(0), shared_ts.value(1)), (1000, 1500)); + + // The new table is created on the selected physical table with the + // microsecond unit. + let mut output = instance + .do_query( + "SELECT greptime_timestamp FROM fresh_us_metric", + ctx.clone(), + ) + .await; + let OutputData::Stream(stream) = output.remove(0).unwrap().data else { + unreachable!() + }; + let batches = common_recordbatch::RecordBatches::try_collect(stream) + .await + .unwrap() + .take(); + let fresh_ts = batches[0] + .column(0) + .as_primitive::(); + assert_eq!(fresh_ts.value(0), 1_500_000); + } + async fn test_prom_store_remote_rw(instance: &Arc, physical_table: Option) { let write_request = WriteRequest { timeseries: prom_store::mock_timeseries(), diff --git a/tests-integration/tests/http.rs b/tests-integration/tests/http.rs index fcafbc0ce52..1b68d3720a7 100644 --- a/tests-integration/tests/http.rs +++ b/tests-integration/tests/http.rs @@ -15,18 +15,20 @@ use std::collections::BTreeMap; use std::io::Write; use std::str::FromStr; +use std::sync::Arc; use std::time::Duration; use api::greptime_proto::io::prometheus::write::v2::histogram::{Count, ZeroCount}; use api::greptime_proto::io::prometheus::write::v2::metadata::MetricType as RemoteWriteV2MetricType; use api::greptime_proto::io::prometheus::write::v2::{ - BucketSpan, Histogram, Metadata as RemoteWriteV2Metadata, Sample as RemoteWriteV2Sample, - TimeSeries as RemoteWriteV2TimeSeries, + BucketSpan, Histogram, Metadata as RemoteWriteV2Metadata, Request as RemoteWriteV2Request, + Sample as RemoteWriteV2Sample, TimeSeries as RemoteWriteV2TimeSeries, }; use api::prom_store::remote::label_matcher::Type as MatcherType; use api::prom_store::remote::{ Label, LabelMatcher, Query, ReadRequest, ReadResponse, Sample, TimeSeries, WriteRequest, }; +use api::v1::RowInsertRequests; use auth::{UserProviderRef, user_provider_from_option}; use axum::http::{HeaderMap, HeaderName, HeaderValue, StatusCode}; use base64::prelude::{BASE64_STANDARD, Engine as _}; @@ -69,7 +71,9 @@ use servers::http::header::constants::{ }; use servers::http::header::{GREPTIME_DB_HEADER_NAME, GREPTIME_TIMEZONE_HEADER_NAME}; use servers::http::otlp::GoogleRpcStatus; -use servers::http::prometheus::{Column, PrometheusJsonResponse, PrometheusResponse}; +use servers::http::prometheus::{ + Column, PromQueryResult, PrometheusJsonResponse, PrometheusResponse, +}; use servers::http::result::error_result::ErrorResponse; use servers::http::result::greptime_result_v1::GreptimedbV1Response; use servers::http::result::influxdb_result_v1::{InfluxdbOutput, InfluxdbV1Response}; @@ -79,11 +83,13 @@ use servers::prom_remote_write::v2::test_util as remote_write_v2; use servers::prom_remote_write::validation::PromValidationMode; use servers::prom_store::{self, mock_timeseries_new_label}; use servers::request_memory_limiter::ServerMemoryLimiter; +use session::context::QueryContextRef; use standalone::options::StandaloneOptions; use table::table_name::TableName; use tests_integration::test_util::{ - MockInstanceImpl, StorageType, assert_wal_delta, build_test_prom_server, setup_test_http_app, - setup_test_http_app_with_frontend, setup_test_http_app_with_frontend_and_slow_query_threshold, + MockInstanceImpl, StorageType, TestGuard, assert_wal_delta, build_test_prom_server, + setup_test_http_app, setup_test_http_app_with_frontend, + setup_test_http_app_with_frontend_and_slow_query_threshold, setup_test_http_app_with_frontend_and_user_provider, setup_test_prom_app_with_frontend, setup_test_prom_app_with_frontend_batched, }; @@ -157,6 +163,7 @@ macro_rules! http_tests { test_http_analyze_stream_tql, test_http_sql_slow_query, test_prometheus_promql_api, + test_promql_over_non_millisecond_physical_tables, test_prometheus_label_replace_response, test_prom_http_api, test_config_api, @@ -165,6 +172,9 @@ macro_rules! http_tests { test_prometheus_remote_write_v2, test_prometheus_remote_write_v2_native_histogram, test_prometheus_remote_write_batched, + test_prometheus_remote_write_batched_mixed_time_index_units, + test_prometheus_remote_write_batched_interceptor_time_index_units, + test_prometheus_remote_write_v2_batched_interceptor_time_index_units, test_prometheus_remote_special_labels, test_prometheus_remote_schema_labels, test_prometheus_remote_write_with_pipeline, @@ -1713,6 +1723,197 @@ pub async fn test_prom_http_api(store_type: StorageType) { guard.remove_all().await; } +/// PromQL must behave identically regardless of the physical metric table's +/// time index unit: the same samples are remote-written into physical tables +/// pre-created with second/micro/nano time indexes (plus a millisecond +/// baseline), and every PromQL response must match the baseline exactly. +/// Regression test for . +pub async fn test_promql_over_non_millisecond_physical_tables(store_type: StorageType) { + common_telemetry::init_default_ut_logging(); + let (app, mut guard) = + setup_test_prom_app_with_frontend(store_type, "promql_non_ms_units").await; + let client = TestClient::new(app).await; + + let units = [ + ("ms", "timestamp(3)"), + ("second", "timestamp(0)"), + ("micro", "timestamp(6)"), + ("nano", "timestamp(9)"), + ]; + + for (suffix, ts_type) in units { + let db = format!("promql_units_{suffix}"); + let res = client + .get(&format!("/v1/sql?db=public&sql=create database {db}")) + .send() + .await; + assert_eq!(res.status(), StatusCode::OK, "create database {db}"); + + let res = client + .get(&format!( + "/v1/sql?db={db}&sql=CREATE TABLE greptime_physical_table \ + (greptime_timestamp {ts_type} NOT NULL, greptime_value DOUBLE NULL, \ + TIME INDEX (greptime_timestamp)) \ + ENGINE = metric WITH ('physical_metric_table' = 'true')" + )) + .send() + .await; + assert_eq!( + res.status(), + StatusCode::OK, + "create physical table with {ts_type} in {db}" + ); + + // unit_gauge{job="demo"} = 1.0 @ 1000ms, 2.0 @ 2000ms + let write_request = WriteRequest { + timeseries: vec![TimeSeries { + labels: vec![ + Label { + name: prom_store::METRIC_NAME_LABEL.to_string(), + value: "unit_gauge".to_string(), + }, + Label { + name: "job".to_string(), + value: "demo".to_string(), + }, + ], + samples: vec![ + Sample { + value: 1.0, + timestamp: 1000, + }, + Sample { + value: 2.0, + timestamp: 2000, + }, + ], + ..Default::default() + }], + ..Default::default() + }; + let compressed = prom_store::snappy_compress(&write_request.encode_to_vec()).unwrap(); + let write_url = format!("/v1/prometheus/write?db={db}"); + let res = client + .post(write_url.as_str()) + .header("Content-Encoding", "snappy") + .body(compressed) + .send() + .await; + assert_eq!( + res.status(), + StatusCode::NO_CONTENT, + "remote write into {db}" + ); + } + + let query_paths = [ + "query?query=unit_gauge&time=2", + "query_range?query=unit_gauge&start=0&end=5&step=1", + "query?query=avg_over_time(unit_gauge[1m])&time=2", + "query?query=rate(unit_gauge[1m])&time=2", + ]; + for path in query_paths { + let mut baseline: Option<(String, PrometheusResponse)> = None; + for (suffix, _) in units { + let db = format!("promql_units_{suffix}"); + let res = client + .get(&format!("/v1/prometheus/api/v1/{path}&db={db}")) + .send() + .await; + assert_eq!(res.status(), StatusCode::OK, "promql `{path}` on {db}"); + let body = res.json::().await; + assert_eq!(body.status, "success", "promql `{path}` on {db}"); + match &baseline { + None => baseline = Some((suffix.to_string(), body.data)), + Some((baseline_suffix, baseline_data)) => { + assert_eq!( + &body.data, baseline_data, + "promql `{path}` on {db} differs from the {baseline_suffix} baseline" + ); + } + } + } + } + + // Pin the baseline responses so the cross-unit equality above cannot pass + // vacuously on wrong data: the samples are 1.0@1s and 2.0@2s, and + // `rate(unit_gauge[1m])@2s` follows Prometheus's extrapolated rate: + // Δv=1 over a 1s sampled interval extrapolated to 1.5s, divided by the + // 60s window -> 1.5/60 = 0.025 per second. + let res = client + .get("/v1/prometheus/api/v1/query?db=promql_units_ms&query=unit_gauge&time=2") + .send() + .await; + assert_eq!(res.status(), StatusCode::OK); + let body = res.json::().await; + assert_eq!( + body.data, + serde_json::from_value::(json!({ + "resultType": "vector", + "result": [{ + "metric": {"__name__": "unit_gauge", "job": "demo"}, + "value": [2.0, "2.0"] + }] + })) + .unwrap() + ); + + let res = client + .get("/v1/prometheus/api/v1/query_range?db=promql_units_ms&query=unit_gauge&start=0&end=5&step=1") + .send() + .await; + assert_eq!(res.status(), StatusCode::OK); + let body = res.json::().await; + assert_eq!( + body.data, + serde_json::from_value::(json!({ + "resultType": "matrix", + "result": [{ + "metric": {"__name__": "unit_gauge", "job": "demo"}, + "values": [ + [1.0, "1.0"], [2.0, "2.0"], [3.0, "2.0"], [4.0, "2.0"], [5.0, "2.0"] + ] + }] + })) + .unwrap() + ); + + let res = client + .get("/v1/prometheus/api/v1/query?db=promql_units_ms&query=avg_over_time(unit_gauge[1m])&time=2") + .send() + .await; + assert_eq!(res.status(), StatusCode::OK); + let body = res.json::().await; + assert_eq!( + body.data, + serde_json::from_value::(json!({ + "resultType": "vector", + "result": [{ + "metric": {"__name__": "unit_gauge", "job": "demo"}, + "value": [2.0, "1.5"] + }] + })) + .unwrap() + ); + + let res = client + .get("/v1/prometheus/api/v1/query?db=promql_units_ms&query=rate(unit_gauge[1m])&time=2") + .send() + .await; + assert_eq!(res.status(), StatusCode::OK); + let body = res.json::().await; + let rate = match body.data { + PrometheusResponse::PromData(data) => match data.result { + PromQueryResult::Vector(v) => v[0].value.as_ref().unwrap().1.clone(), + other => panic!("expected vector, got {other:?}"), + }, + other => panic!("expected prom data, got {other:?}"), + }; + assert_eq!(rate, "0.025"); + + guard.remove_all().await; +} + pub async fn test_metrics_api(store_type: StorageType) { common_telemetry::init_default_ut_logging(); let (app, mut guard) = setup_test_http_app(store_type, "metrics_api").await; @@ -3593,6 +3794,277 @@ async fn write_prometheus_skip_wal_sample( assert_eq!(request.send().await.status(), StatusCode::NO_CONTENT); } +/// Batched remote write against a logical table bound to a non-millisecond +/// physical table: the bulk guard must reject the destination (the bulk +/// encode only produces millisecond batches) and fall back to the ordinary +/// insert path, which converts the requests to the table's unit. +pub async fn test_prometheus_remote_write_batched_mixed_time_index_units(store_type: StorageType) { + common_telemetry::init_default_ut_logging(); + let (app, mut guard) = + setup_test_prom_app_with_frontend_batched(store_type, "prom_rw_batched_mixed_units").await; + let client = TestClient::new(app).await; + + let res = client + .get("/v1/sql?db=public&sql=CREATE TABLE phy_us (greptime_timestamp TIMESTAMP(6) NOT NULL, greptime_value DOUBLE NULL, TIME INDEX (greptime_timestamp)) ENGINE = metric WITH ('physical_metric_table' = 'true')") + .send() + .await; + assert_eq!(res.status(), StatusCode::OK); + + let write = |metric: &str, value: f64, timestamp: i64| { + let write_request = WriteRequest { + timeseries: vec![TimeSeries { + labels: vec![ + Label { + name: prom_store::METRIC_NAME_LABEL.to_string(), + value: metric.to_string(), + }, + Label { + name: "job".to_string(), + value: "demo".to_string(), + }, + ], + samples: vec![Sample { value, timestamp }], + ..Default::default() + }], + ..Default::default() + }; + prom_store::snappy_compress(&write_request.encode_to_vec()).unwrap() + }; + + // Create the logical table on the microsecond physical table. The bulk + // guard rejects the microsecond selected physical table, so this write + // takes the ordinary insert path. + let res = client + .post("/v1/prometheus/write?physical_table=phy_us") + .header("Content-Encoding", "snappy") + .body(write("us_metric", 2.5, 1500)) + .send() + .await; + assert_eq!(res.status(), StatusCode::NO_CONTENT); + + // Write the same existing table again while selecting the default + // (millisecond) physical table: the guard must reject the existing + // microsecond destination and fall back to the ordinary insert path — + // without the destination check the bulk path would build millisecond + // arrays against the microsecond schema and fail the write. + let res = client + .post("/v1/prometheus/write") + .header("Content-Encoding", "snappy") + .body(write("us_metric", 3.5, 2000)) + .send() + .await; + assert_eq!(res.status(), StatusCode::NO_CONTENT); + + // Both samples are stored on the microsecond time index. + validate_data( + "prom_rw_batched_mixed_units", + &client, + "SELECT COUNT(*), MAX(greptime_value) FROM us_metric", + "[[2,3.5]]", + ) + .await; + + guard.remove_all().await; +} + +/// A Prometheus write interceptor that redirects every remote write to the +/// `tenant_redirect` schema, like a per-tenant redirection: the incoming +/// context targets `public` while the destination tables live elsewhere. +struct PromSchemaRedirectInterceptor; + +impl servers::interceptor::PromStoreProtocolInterceptor for PromSchemaRedirectInterceptor { + type Error = servers::error::Error; + + fn pre_write( + &self, + _write_req: &RowInsertRequests, + ctx: QueryContextRef, + ) -> servers::error::Result<()> { + ctx.set_current_schema("tenant_redirect"); + Ok(()) + } +} + +async fn setup_redirecting_batched_prom_app( + store_type: StorageType, + name: &str, +) -> (TestClient, TestGuard) { + let plugins = Plugins::default(); + plugins.insert::>( + Arc::new(PromSchemaRedirectInterceptor), + ); + let standalone = tests_integration::standalone::GreptimeDbStandaloneBuilder::new(name) + .with_default_store_type(store_type) + .with_plugin(plugins) + .build() + .await; + let server = build_test_prom_server(standalone.fe_instance().clone(), true) + .with_greptime_config_options(standalone.opts.datanode_options().to_toml().unwrap()) + .build(); + let client = TestClient::new(server.build(server.make_app()).unwrap()).await; + (client, standalone.guard) +} + +async fn create_redirect_schema_with_microsecond_table(client: &TestClient, table: &str) { + for sql in [ + "create database if not exists tenant_redirect", + "CREATE TABLE tenant_redirect.phy_us \ + (ts timestamp(6) time index, val double, host string primary key) \ + engine=metric with ('physical_metric_table' = 'true')", + &format!( + "CREATE TABLE tenant_redirect.{table} \ + (ts timestamp(6) time index, val double, host string primary key) \ + engine=metric with ('on_physical_table' = 'phy_us')" + ), + ] { + let res = client + .get(format!("/v1/sql?sql={sql}").as_str()) + .send() + .await; + assert_eq!(res.status(), StatusCode::OK, "setup: {sql}"); + } +} + +/// Regression test: with batching enabled and a pre_write interceptor that +/// redirects the context to a schema holding same-named non-millisecond +/// tables, the bulk eligibility must be evaluated against the redirected +/// destinations (after preflight) and fall back to the ordinary insert path +/// instead of failing in the bulk encode. +pub async fn test_prometheus_remote_write_batched_interceptor_time_index_units( + store_type: StorageType, +) { + common_telemetry::init_default_ut_logging(); + let (client, mut guard) = + setup_redirecting_batched_prom_app(store_type, "prom_rw_batched_interceptor_units").await; + create_redirect_schema_with_microsecond_table(&client, "intercept_metric").await; + + let write_request = WriteRequest { + timeseries: vec![TimeSeries { + labels: vec![ + Label { + name: prom_store::METRIC_NAME_LABEL.to_string(), + value: "intercept_metric".to_string(), + }, + Label { + name: "job".to_string(), + value: "demo".to_string(), + }, + ], + samples: vec![Sample { + value: 1.0, + timestamp: 1000, + }], + ..Default::default() + }], + ..Default::default() + }; + let compressed = prom_store::snappy_compress(&write_request.encode_to_vec()).unwrap(); + let res = client + .post("/v1/prometheus/write") + .header("Content-Encoding", "snappy") + .body(compressed) + .send() + .await; + assert_eq!(res.status(), StatusCode::NO_CONTENT); + + validate_data( + "prom_rw_batched_interceptor_units", + &client, + "SELECT COUNT(*) FROM tenant_redirect.intercept_metric", + "[[1]]", + ) + .await; + + guard.remove_all().await; +} + +/// The v2 variant of the interceptor regression test: a mixed samples-plus- +/// histograms request whose existing sample destination lives in the +/// redirected schema with a microsecond time index must fall back to the +/// ordinary insert path (both series land), not fail in the bulk encode. +pub async fn test_prometheus_remote_write_v2_batched_interceptor_time_index_units( + store_type: StorageType, +) { + common_telemetry::init_default_ut_logging(); + let (client, mut guard) = + setup_redirecting_batched_prom_app(store_type, "prom_rw_v2_batched_interceptor_units") + .await; + create_redirect_schema_with_microsecond_table(&client, "v2_mixed_sample").await; + + // One sample series targeting the existing microsecond table, plus one + // histogram series for a new (millisecond) table. + let mut symbols = vec![String::new()]; + let mut symbol = |value: &str| { + symbols.push(value.to_string()); + (symbols.len() - 1) as u32 + }; + let name_ref = symbol("__name__"); + let sample_ref = symbol("v2_mixed_sample"); + let job_ref = symbol("job"); + let demo_ref = symbol("demo"); + let name_ref2 = symbol("__name__"); + let histo_ref = symbol("v2_mixed_histo"); + let request = RemoteWriteV2Request { + symbols, + timeseries: vec![ + RemoteWriteV2TimeSeries { + labels_refs: vec![name_ref, sample_ref, job_ref, demo_ref], + samples: vec![RemoteWriteV2Sample { + value: 1.0, + timestamp: 1000, + start_timestamp: 0, + }], + ..Default::default() + }, + RemoteWriteV2TimeSeries { + labels_refs: vec![name_ref2, histo_ref, job_ref, demo_ref], + histograms: vec![Histogram { + count: Some(Count::CountInt(1)), + sum: 1.0, + positive_spans: vec![BucketSpan { + offset: 0, + length: 1, + }], + positive_deltas: vec![1], + timestamp: 1000, + start_timestamp: 1000, + ..Default::default() + }], + ..Default::default() + }, + ], + }; + let compressed = prom_store::snappy_compress(&request.encode_to_vec()).unwrap(); + let res = client + .post("/v1/prometheus/write") + .header( + "Content-Type", + "application/x-protobuf;proto=io.prometheus.write.v2.Request", + ) + .header("Content-Encoding", "snappy") + .body(compressed) + .send() + .await; + assert_eq!(res.status(), StatusCode::NO_CONTENT); + + validate_data( + "prom_rw_v2_batched_interceptor_units", + &client, + "SELECT COUNT(*) FROM tenant_redirect.v2_mixed_sample", + "[[1]]", + ) + .await; + validate_data( + "prom_rw_v2_batched_interceptor_units", + &client, + "SELECT COUNT(*) FROM tenant_redirect.v2_mixed_histo", + "[[1]]", + ) + .await; + + guard.remove_all().await; +} + /// Covers the batched (pending-rows-batcher) Prometheus remote write path, which /// bypasses `PromStoreProtocolHandler::write`. Verifies the metric table is created /// asynchronously and still carries the Prometheus semantic identity stamped on the