diff --git a/src/frontend/src/instance/otlp.rs b/src/frontend/src/instance/otlp.rs index e82aa343870..5132b9f8a8d 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 { - mut requests, + requests, rows, semantic_index, resource_info, @@ -170,33 +170,21 @@ impl OpenTelemetryProtocolHandler for Instance { .extension(PHYSICAL_TABLE_PARAM) .unwrap_or(GREPTIME_PHYSICAL_TABLE) .to_string(); + // The bulk path converts each request's time index unit to the + // destination table's unit during batch alignment, so no pre-gate + // alignment is needed here. 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() { - // 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 + let batcher = if batcher.is_some() + && self .inserter .can_batch_metric_rows(&requests, &ctx, &physical_table) .await .map_err(BoxedError::new) .context(error::ExecuteGrpcQuerySnafu)? - { - batcher - } else { - None - } + { + batcher } else { None }; diff --git a/src/operator/src/insert.rs b/src/operator/src/insert.rs index 24acdeaba28..a0ebe236767 100644 --- a/src/operator/src/insert.rs +++ b/src/operator/src/insert.rs @@ -172,14 +172,19 @@ impl Inserter { return Ok(false); } for request in &requests.inserts { - // The logical bulk encoder only supports scalar metric schemas. + // The logical bulk encoder only supports scalar metric schemas; + // any time index unit is accepted — requests are converted to the + // destination table's unit during batch alignment. // Check new tables too, before catalog lookup or schema changes. if request.rows.as_ref().is_some_and(|rows| { rows.schema.iter().any(|column| { column.datatype_extension.is_some() || !matches!( ColumnDataType::try_from(column.datatype), - Ok(ColumnDataType::TimestampMillisecond + Ok(ColumnDataType::TimestampSecond + | ColumnDataType::TimestampMillisecond + | ColumnDataType::TimestampMicrosecond + | ColumnDataType::TimestampNanosecond | ColumnDataType::Float64 | ColumnDataType::String) ) @@ -1376,75 +1381,6 @@ 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, @@ -2373,69 +2309,6 @@ mod tests { ); } - #[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 fcd4e1ea38f..7e0cea44b6a 100644 --- a/src/servers/src/batcher/logical_table.rs +++ b/src/servers/src/batcher/logical_table.rs @@ -21,7 +21,7 @@ mod tables; #[cfg(test)] mod test_util; -use std::collections::HashSet; +use std::collections::{HashMap, HashSet}; use std::future::{Future, ready}; use std::num::NonZeroUsize; use std::sync::Arc; @@ -42,6 +42,7 @@ use meter_macros::write_meter; use partition::manager::PartitionRuleManagerRef; use session::context::QueryContextRef; use snafu::ResultExt; +use store_api::metric_engine_consts::{LOGICAL_TABLE_METADATA_KEY, METRIC_ENGINE_NAME}; use tokio::sync::{Semaphore, broadcast, mpsc, oneshot}; use crate::batcher::flow_notifier::{FlowNotifier, start_flow_notification_worker}; @@ -173,79 +174,17 @@ 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 { + /// Returns the physical metric table's time index unit resolved from + /// `ctx`, defaulting to millisecond when the table does not exist yet + /// (the schema alterer auto-creates it as millisecond). + async fn physical_time_index_unit_or_default(&self, ctx: &QueryContextRef) -> TimeUnit { 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; + return TimeUnit::Millisecond; }; table .table_info() @@ -253,8 +192,77 @@ impl LogicalTablePendingRowsBatcher { .schema .timestamp_column() .and_then(|col| col.data_type.as_timestamp().map(|ts| ts.unit())) - .map(|unit| unit == TimeUnit::Millisecond) - .unwrap_or(true) + .unwrap_or(TimeUnit::Millisecond) + } + + /// Returns whether the bulk path can accept `batches`: every existing + /// destination table must be a metric logical table bound to the + /// physical table selected by its context. A destination bound to + /// another physical table would be flushed through the selected + /// physical's regions, silently misplacing its rows, so such requests + /// must stay on the ordinary insert path (which routes per destination). + /// New tables are always fine: they are created on the selected physical + /// table. Time index units need no check here — the bulk encode converts + /// each request to its destination's unit. Destinations are resolved + /// once per distinct (schema, table). + pub(crate) async fn accepts_bulk_destinations( + &self, + batches: impl Iterator, + ) -> bool { + // One request can select different physical tables per batch (e.g. + // per-series physical-table labels), so the dedupe key includes the + // selected physical table: every distinct (schema, table, physical) + // triple is validated against the table's actual binding. + let mut checked = HashSet::new(); + // For missing tables, one request must not select two different + // physical tables: the batcher would create the table through one + // selection and flush its rows through the other's regions. + let mut missing_selections: HashMap<(String, String), String> = HashMap::new(); + for (ctx, requests) in batches { + let physical_table = batch_key_from_ctx(ctx).physical_table; + let schema = ctx.current_schema(); + for request in &requests.inserts { + if !checked.insert(( + schema.clone(), + request.table_name.clone(), + physical_table.clone(), + )) { + continue; + } + let Ok(Some(table)) = self + .catalog_manager + .table(ctx.current_catalog(), &schema, &request.table_name, None) + .await + else { + // New table: created on the selected physical table, but + // a conflicting selection within the same request cannot + // be batched. + if missing_selections + .insert( + (schema.clone(), request.table_name.clone()), + physical_table.clone(), + ) + .is_some_and(|previous| previous != physical_table) + { + return false; + } + continue; + }; + let info = table.table_info(); + if info.meta.engine != METRIC_ENGINE_NAME + || info + .meta + .options + .extra_options + .get(LOGICAL_TABLE_METADATA_KEY) + .map(String::as_str) + != Some(physical_table.as_str()) + { + return false; + } + } + } + true } /// Submits with request-level accounting after schema preparation and before diff --git a/src/servers/src/batcher/logical_table/batch_convert.rs b/src/servers/src/batcher/logical_table/batch_convert.rs index 28464eb2d3b..37080dbe258 100644 --- a/src/servers/src/batcher/logical_table/batch_convert.rs +++ b/src/servers/src/batcher/logical_table/batch_convert.rs @@ -17,7 +17,7 @@ use std::sync::Arc; use std::time::{Duration, Instant}; use arrow::compute::concat_batches; -use arrow::datatypes::{DataType as ArrowDataType, Schema as ArrowSchema, TimeUnit}; +use arrow::datatypes::{DataType as ArrowDataType, Schema as ArrowSchema}; use arrow::record_batch::RecordBatch; use common_query::prelude::{greptime_timestamp, greptime_value}; use metric_engine::batch_modifier::{TagColumnInfo, modify_batch_sparse}; @@ -130,7 +130,7 @@ pub(in crate::batcher::logical_table) fn columns_taxonomy( essential_column_indices.push(index); } } - ArrowDataType::Timestamp(TimeUnit::Millisecond, _) => { + ArrowDataType::Timestamp(_, _) => { ensure!( timestamp_index.replace(index).is_none(), error::InvalidPromRemoteRequestSnafu { diff --git a/src/servers/src/batcher/logical_table/tables.rs b/src/servers/src/batcher/logical_table/tables.rs index 87110366a8c..140dd2d4f23 100644 --- a/src/servers/src/batcher/logical_table/tables.rs +++ b/src/servers/src/batcher/logical_table/tables.rs @@ -15,9 +15,10 @@ use std::collections::{HashMap, HashSet}; use std::sync::Arc; -use api::v1::{ColumnSchema, RowInsertRequests, Rows}; +use api::v1::{ColumnSchema, RowInsertRequests, Rows, SemanticType}; use arrow::datatypes::Schema as ArrowSchema; use async_trait::async_trait; +use common_time::timestamp::TimeUnit; use session::context::QueryContextRef; use snafu::OptionExt; @@ -91,6 +92,17 @@ impl LogicalTablePendingRowsBatcher { .plan_table_resolution(&catalog, &schema, ctx, &unique_tables) .await?; + // New tables are created on the request's selected physical table; + // their time index must use the physical table's unit (a missing + // physical table is auto-created as millisecond by the schema + // alterer, matching the default here). + if !plan.tables_to_create.is_empty() { + let physical_unit = self.physical_time_index_unit_or_default(ctx).await; + for (_, request_schema) in &mut plan.tables_to_create { + align_create_schema_time_index(request_schema, physical_unit); + } + } + self.create_missing_tables_and_refresh_schemas( &catalog, &schema, @@ -109,6 +121,17 @@ impl LogicalTablePendingRowsBatcher { } } +/// Rewrites the time index column of a create-table schema to `unit`, if the +/// schema carries a timestamp column in another unit. +fn align_create_schema_time_index(request_schema: &mut [ColumnSchema], unit: TimeUnit) { + for column in request_schema { + if column.semantic_type == SemanticType::Timestamp as i32 { + column.datatype = api::helper::timestamp_datatype(unit) as i32; + column.datatype_extension = None; + } + } +} + impl LogicalTablePendingRowsBatcher { /// Extracts non-empty `(table_name, rows)` pairs and computes total row /// count across the retained entries. diff --git a/src/servers/src/http/prom_store.rs b/src/servers/src/http/prom_store.rs index 44abb18a5be..68542a6f642 100644 --- a/src/servers/src/http/prom_store.rs +++ b/src/servers/src/http/prom_store.rs @@ -408,11 +408,11 @@ async fn write_prometheus_rows_with_progress( error, rows_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 { + // Destinations bound to another physical table must stay on the + // ordinary insert path, which routes per destination; time index + // units need no check — the bulk encode converts each request to + // its destination's unit. + if batcher.accepts_bulk_destinations(batches.iter()).await { let mut rows_written = 0; for (temp_ctx, reqs) in batches { let rows = @@ -543,11 +543,11 @@ async fn write_prometheus_v2_rows_with_progress( 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 { + // Destinations bound to another physical table must stay on the + // ordinary insert path, which routes per destination; time index + // units need no check — the bulk encode converts each request to + // its destination's unit. + if batcher.accepts_bulk_destinations(batches.iter()).await { return write_batched_prometheus_v2_rows_with_progress( prom_store_handler, batcher.as_ref(), diff --git a/src/servers/src/prom_row_builder.rs b/src/servers/src/prom_row_builder.rs index 8b4e35bcd32..43b3a464857 100644 --- a/src/servers/src/prom_row_builder.rs +++ b/src/servers/src/prom_row_builder.rs @@ -23,14 +23,17 @@ use api::helper::ColumnDataTypeWrapper; use api::v1::value::ValueData; use api::v1::{ColumnSchema, Rows, SemanticType}; use arrow::array::{ - ArrayRef, Float64Builder, StringBuilder, TimestampMicrosecondBuilder, - TimestampMillisecondBuilder, TimestampNanosecondBuilder, TimestampSecondBuilder, - new_null_array, + ArrayRef, ArrowPrimitiveType, Float64Builder, PrimitiveBuilder, StringBuilder, new_null_array, +}; +use arrow::datatypes::{ + DataType as ArrowDataType, Schema as ArrowSchema, TimestampMicrosecondType, + TimestampMillisecondType, TimestampNanosecondType, TimestampSecondType, }; -use arrow::datatypes::{DataType as ArrowDataType, Schema as ArrowSchema}; use arrow::record_batch::RecordBatch; use arrow_schema::TimeUnit; use common_query::prelude::{greptime_timestamp, greptime_value}; +use common_time::Timestamp; +use common_time::timestamp::TimeUnit as CommonTimeUnit; use datatypes::data_type::DataType; use datatypes::prelude::ConcreteDataType; use snafu::{OptionExt, ResultExt, ensure}; @@ -129,16 +132,7 @@ pub(crate) fn rows_to_aligned_record_batch( ArrowDataType::Float64 => { source_map.insert(&target_field_name, (src_idx, src_arrow_type)); } - ArrowDataType::Timestamp(unit, _) => { - ensure!( - unit == &TimeUnit::Millisecond, - error::InvalidPromRemoteRequestSnafu { - msg: format!( - "Unexpected remote write batch timestamp unit, expect millisecond, got: {}", - unit - ) - } - ); + ArrowDataType::Timestamp(_, _) => { source_map.insert(&target_ts_name, (src_idx, src_arrow_type)); } ArrowDataType::Utf8 => { @@ -156,15 +150,28 @@ pub(crate) fn rows_to_aligned_record_batch( } } - // Build columns in target schema order + // Build columns in target schema order. The timestamp column is built in + // the TARGET schema's unit; `build_arrow_array` converts the request's + // encoding unit (flooring on narrowing) into it. let mut columns = Vec::with_capacity(target_schema.fields().len()); for target_field in target_schema.fields() { if let Some((src_idx, src_arrow_type)) = source_map.get(target_field.name().as_str()) { + let target_type = if matches!( + (src_arrow_type, target_field.data_type()), + ( + ArrowDataType::Timestamp(_, _), + ArrowDataType::Timestamp(_, _) + ) + ) { + target_field.data_type().clone() + } else { + src_arrow_type.clone() + }; let array = build_arrow_array( rows, *src_idx, &rows.schema[*src_idx].column_name, - src_arrow_type.clone(), + target_type, row_count, )?; columns.push(array); @@ -199,6 +206,65 @@ pub(crate) fn identify_missing_columns_from_proto( Ok(missing) } +/// Builds a timestamp array of `T`'s unit from a proto column, appending +/// each value directly into the builder. Values already in the target unit +/// are appended without conversion (the unchanged-unit fast path); others +/// are converted via `Timestamp::convert_to`, flooring on narrowing. +fn build_timestamp_array>( + rows: &Rows, + col_idx: usize, + column_name: &str, + row_count: usize, + target_unit: CommonTimeUnit, +) -> Result { + let mut builder = PrimitiveBuilder::::with_capacity(row_count); + for row in &rows.rows { + let Some(value) = row.values[col_idx].value_data.as_ref() else { + builder.append_null(); + continue; + }; + let (source_unit, raw) = match value { + ValueData::TimestampSecondValue(v) => (CommonTimeUnit::Second, *v), + ValueData::TimestampMillisecondValue(v) => (CommonTimeUnit::Millisecond, *v), + ValueData::DatetimeValue(v) | ValueData::TimestampMicrosecondValue(v) => { + (CommonTimeUnit::Microsecond, *v) + } + ValueData::TimestampNanosecondValue(v) => (CommonTimeUnit::Nanosecond, *v), + v => { + return error::InvalidPromRemoteRequestSnafu { + msg: format!("Unexpected value: {:?}", v), + } + .fail(); + } + }; + if source_unit == target_unit { + builder.append_value(raw); + } else { + let timestamp = Timestamp::new(raw, source_unit); + let Some(converted) = timestamp.convert_to(target_unit) else { + return error::InvalidPromRemoteRequestSnafu { + msg: format!( + "Timestamp value in column '{column_name}' overflows when converting to unit {target_unit:?}" + ), + } + .fail(); + }; + builder.append_value(converted.value()); + } + } + Ok(Arc::new(builder.finish()) as ArrayRef) +} + +/// Converts an arrow time unit to the common time unit. +fn arrow_time_unit(unit: TimeUnit) -> CommonTimeUnit { + match unit { + TimeUnit::Second => CommonTimeUnit::Second, + TimeUnit::Millisecond => CommonTimeUnit::Millisecond, + TimeUnit::Microsecond => CommonTimeUnit::Microsecond, + TimeUnit::Nanosecond => CommonTimeUnit::Nanosecond, + } +} + /// Build a `Vec` suitable for creating a new Prometheus logical table /// directly from the proto `rows.schema`, avoiding the round-trip through Arrow schema. pub fn build_prom_create_table_schema_from_proto( @@ -207,13 +273,22 @@ pub fn build_prom_create_table_schema_from_proto( rows_schema .iter() .map(|col| { - let semantic_type = if col.datatype == api::v1::ColumnDataType::TimestampMillisecond as i32 { + let datatype = api::v1::ColumnDataType::try_from(col.datatype).map_err(|_| { + error::InvalidPromRemoteRequestSnafu { + msg: format!( + "Failed to build create table schema, column '{}' has unknown datatype {}", + col.column_name, col.datatype + ), + } + .build() + })?; + let semantic_type = if api::helper::timestamp_unit(datatype).is_some() { SemanticType::Timestamp - } else if col.datatype == api::v1::ColumnDataType::Float64 as i32 { + } else if datatype == api::v1::ColumnDataType::Float64 { SemanticType::Field } else { // tag columns must be String type - ensure!(col.datatype == api::v1::ColumnDataType::String as i32, error::InvalidPromRemoteRequestSnafu{ + ensure!(datatype == api::v1::ColumnDataType::String, error::InvalidPromRemoteRequestSnafu{ msg: format!( "Failed to build create table schema, tag column '{}' must be String but got datatype {}", col.column_name, col.datatype @@ -268,25 +343,45 @@ fn build_arrow_array( StringBuilder::with_capacity(row_count, 0), ValueData::StringValue(v) => v ), - arrow::datatypes::DataType::Timestamp(u, _) => match u { - TimeUnit::Second => build_array!( - TimestampSecondBuilder::with_capacity(row_count), - ValueData::TimestampSecondValue(v) => *v - ), - TimeUnit::Millisecond => build_array!( - TimestampMillisecondBuilder::with_capacity(row_count), - ValueData::TimestampMillisecondValue(v) => *v - ), - TimeUnit::Microsecond => build_array!( - TimestampMicrosecondBuilder::with_capacity(row_count), - ValueData::DatetimeValue(v) => *v, - ValueData::TimestampMicrosecondValue(v) => *v - ), - TimeUnit::Nanosecond => build_array!( - TimestampNanosecondBuilder::with_capacity(row_count), - ValueData::TimestampNanosecondValue(v) => *v - ), - }, + arrow::datatypes::DataType::Timestamp(u, _) => { + // Accept any timestamp encoding and append directly into the + // target-unit builder — no intermediate column allocation. + // Values already in the target unit (the common unchanged + // millisecond case) are appended as-is; others are converted + // first, flooring on narrowing (same semantics as + // `Timestamp::convert_to` on the ordinary insert path). + let target_unit = arrow_time_unit(u); + match u { + TimeUnit::Second => build_timestamp_array::( + rows, + col_idx, + column_name, + row_count, + target_unit, + )?, + TimeUnit::Millisecond => build_timestamp_array::( + rows, + col_idx, + column_name, + row_count, + target_unit, + )?, + TimeUnit::Microsecond => build_timestamp_array::( + rows, + col_idx, + column_name, + row_count, + target_unit, + )?, + TimeUnit::Nanosecond => build_timestamp_array::( + rows, + col_idx, + column_name, + row_count, + target_unit, + )?, + } + } ty => { return error::InvalidPromRemoteRequestSnafu { msg: format!( @@ -305,7 +400,9 @@ fn build_arrow_array( mod tests { use api::v1::value::ValueData; use api::v1::{ColumnDataType, ColumnSchema, Row, Rows, SemanticType, Value}; - use arrow::array::{Array, Float64Array, StringArray, TimestampMillisecondArray}; + use arrow::array::{ + Array, Float64Array, StringArray, TimestampMicrosecondArray, TimestampMillisecondArray, + }; use arrow::datatypes::{DataType, Field, Schema as ArrowSchema, TimeUnit}; use super::{ @@ -313,6 +410,105 @@ mod tests { rows_to_aligned_record_batch, }; + #[test] + fn test_rows_to_aligned_record_batch_converts_time_units() { + // Microsecond-encoded rows aligned to a millisecond schema: narrowing + // floors towards negative infinity (-1001us -> -2ms), matching + // `Timestamp::convert_to` on the ordinary insert path. + let rows = Rows { + schema: vec![ + ColumnSchema { + column_name: "greptime_timestamp".to_string(), + datatype: ColumnDataType::TimestampMicrosecond as i32, + semantic_type: SemanticType::Timestamp as i32, + ..Default::default() + }, + ColumnSchema { + column_name: "greptime_value".to_string(), + datatype: ColumnDataType::Float64 as i32, + semantic_type: SemanticType::Field as i32, + ..Default::default() + }, + ], + rows: vec![ + Row { + values: vec![ + Value { + value_data: Some(ValueData::TimestampMicrosecondValue(1000)), + }, + Value { + value_data: Some(ValueData::F64Value(1.0)), + }, + ], + }, + Row { + values: vec![ + Value { + value_data: Some(ValueData::TimestampMicrosecondValue(-1001)), + }, + Value { + value_data: Some(ValueData::F64Value(2.0)), + }, + ], + }, + ], + }; + let target = ArrowSchema::new(vec![ + Field::new( + "greptime_timestamp", + DataType::Timestamp(TimeUnit::Millisecond, None), + false, + ), + Field::new("greptime_value", DataType::Float64, true), + ]); + + let (batch, _) = rows_to_aligned_record_batch(&rows, &target) + .unwrap() + .into_parts(); + let ts = batch + .column(0) + .as_any() + .downcast_ref::() + .unwrap(); + assert_eq!(ts.value(0), 1); + assert_eq!(ts.value(1), -2); + + // Millisecond-encoded rows aligned to a microsecond schema: widening + // is lossless. + let mut rows = Rows { + schema: vec![rows.schema[0].clone(), rows.schema[1].clone()], + rows: vec![Row { + values: vec![ + Value { + value_data: Some(ValueData::TimestampMillisecondValue(123)), + }, + Value { + value_data: Some(ValueData::F64Value(4.0)), + }, + ], + }], + }; + rows.schema[0].datatype = ColumnDataType::TimestampMillisecond as i32; + let target = ArrowSchema::new(vec![ + Field::new( + "greptime_timestamp", + DataType::Timestamp(TimeUnit::Microsecond, None), + false, + ), + Field::new("greptime_value", DataType::Float64, true), + ]); + + let (batch, _) = rows_to_aligned_record_batch(&rows, &target) + .unwrap() + .into_parts(); + let ts = batch + .column(0) + .as_any() + .downcast_ref::() + .unwrap(); + assert_eq!(ts.value(0), 123_000); + } + #[test] fn test_rows_to_aligned_record_batch_renames_and_reorders() { let rows = Rows { diff --git a/tests-integration/src/otlp.rs b/tests-integration/src/otlp.rs index 59ce38919f7..29c8891a695 100644 --- a/tests-integration/src/otlp.rs +++ b/tests-integration/src/otlp.rs @@ -16,6 +16,7 @@ mod test { use std::sync::Arc; + use TimeUnit as ArrowTimeUnit; use client::{DEFAULT_CATALOG_NAME, OutputData}; use common_recordbatch::RecordBatches; use datatypes::arrow::array::AsArray; @@ -508,6 +509,268 @@ WITH( Ok(()) } + #[tokio::test(flavor = "multi_thread")] + async fn test_otlp_logical_batcher_non_millisecond_physical_table() { + // Microsecond and nanosecond physical tables must both use the + // logical batcher: requests are converted to the physical table's + // unit during batch alignment + // (). + run_otlp_logical_batcher_non_millisecond_physical_table( + "us", + "TIMESTAMP(6)", + ArrowTimeUnit::Microsecond, + [60_000_000, 120_000_000], + ) + .await; + run_otlp_logical_batcher_non_millisecond_physical_table( + "ns", + "TIMESTAMP(9)", + ArrowTimeUnit::Nanosecond, + [60_000_000_000, 120_000_000_000], + ) + .await; + } + + async fn run_otlp_logical_batcher_non_millisecond_physical_table( + suite: &str, + sql_ts_type: &str, + unit: ArrowTimeUnit, + expected: [i64; 2], + ) { + use std::time::Duration; + + use common_base::Plugins; + use datatypes::arrow::array::AsArray; + use datatypes::arrow::datatypes::{TimestampMicrosecondType, TimestampNanosecondType}; + use frontend::server::Services; + use frontend::service_config::pending_rows_batcher::BatcherOptions; + use prost::Message; + use servers::batcher::BatchingProtocol; + use servers::http::test_helpers::TestClient; + use session::protocol_ctx::{OtlpMetricCtx, ProtocolCtx}; + + let standalone = + GreptimeDbStandaloneBuilder::new(&format!("otlp_logical_{suite}_physical")) + .with_logical_batcher(BatcherOptions { + protocols: vec![BatchingProtocol::Otlp], + pending_rows_flush_interval: Duration::from_millis(5), + ..Default::default() + }) + .build() + .await; + let instance = standalone.fe_instance(); + let options = standalone.opts.clone(); + let services = Services::new(options.clone(), instance.clone(), Plugins::default()); + let server = services + .http_server_builder( + &options.frontend_options(), + services.server_memory_limiter.clone(), + ) + .build(); + let client = TestClient::new(server.build(server.make_app()).unwrap()).await; + + let mut ctx = QueryContext::with(DEFAULT_CATALOG_NAME, "public"); + ctx.set_logical_batching_enabled(true); + ctx.set_protocol_ctx(ProtocolCtx::OtlpMetric(OtlpMetricCtx { + with_metric_engine: true, + ..Default::default() + })); + let ctx = Arc::new(ctx); + + // Pre-create the physical metric table with a non-millisecond time + // index BEFORE any ingestion, so the batcher's bulk path must + // handle it. + 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()); + + // submit_build_and_align counts every batcher submission in both + // acknowledgement modes, so the assertion detects a regression that + // disables non-millisecond batching. + let submissions = servers::metrics::PENDING_ROWS_BATCH_INGEST_STAGE_ELAPSED + .with_label_values(&["submit_build_and_align"]); + let before = submissions.get_sample_count(); + for (ts, value) in [(60, 10), (120, 20)] { + let request = build_sum_request( + "non.ms.alignment", + AggregationTemporality::Cumulative, + &[(ts, value)], + ); + let response = client + .post("/v1/otlp/v1/metrics") + .header("content-type", "application/x-protobuf") + .body(request.encode_to_vec()) + .send() + .await; + assert_eq!(response.status().as_u16(), 200); + } + assert_eq!( + submissions.get_sample_count() - before, + 2, + "non-millisecond physical tables must use the logical batcher" + ); + + // The rows land on the non-millisecond physical table, converted + // from the nanosecond encoding. + let sql = "SELECT greptime_timestamp, greptime_value FROM non_ms_alignment_total ORDER BY greptime_timestamp"; + let batches = tokio::time::timeout(Duration::from_secs(5), async { + loop { + let output = instance.do_query(sql, ctx.clone()).await.remove(0).unwrap(); + let OutputData::Stream(stream) = output.data else { + panic!("expected stream") + }; + let batches = RecordBatches::try_collect(stream).await.unwrap(); + if batches.iter().map(|batch| batch.num_rows()).sum::() == 2 { + break batches; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .unwrap(); + let batch = &batches.take()[0]; + let stored = |row: usize| match unit { + ArrowTimeUnit::Microsecond => batch + .column(0) + .as_primitive::() + .value(row), + ArrowTimeUnit::Nanosecond => batch + .column(0) + .as_primitive::() + .value(row), + _ => unreachable!("covered units only"), + }; + assert_eq!(stored(0), expected[0]); + assert_eq!(stored(1), expected[1]); + let values = batch + .column(1) + .as_any() + .downcast_ref::() + .unwrap_or_else(|| panic!("expected f64 values")); + assert_eq!(values.value(0), 10.0); + assert_eq!(values.value(1), 20.0); + } + + #[tokio::test(flavor = "multi_thread")] + async fn test_otlp_logical_batcher_fallback_for_cross_physical_destination() { + use std::time::Duration; + + use common_base::Plugins; + use datatypes::arrow::array::AsArray; + use datatypes::arrow::datatypes::TimestampMicrosecondType; + use frontend::server::Services; + use frontend::service_config::pending_rows_batcher::BatcherOptions; + use prost::Message; + use servers::batcher::BatchingProtocol; + use servers::http::test_helpers::TestClient; + use session::protocol_ctx::{OtlpMetricCtx, ProtocolCtx}; + + // With the logical batcher enabled, an OTLP request targeting an + // existing logical table bound to ANOTHER physical table must fall + // back to the ordinary insert path (the bulk eligibility check + // rejects the destination binding), which converts the request to + // the destination's unit. This covers the case previously guarded + // by the removed OTLP pre-gate alignment. + let standalone = GreptimeDbStandaloneBuilder::new("otlp_logical_cross_physical") + .with_logical_batcher(BatcherOptions { + protocols: vec![BatchingProtocol::Otlp], + pending_rows_flush_interval: Duration::from_millis(5), + ..Default::default() + }) + .build() + .await; + let instance = standalone.fe_instance(); + let options = standalone.opts.clone(); + let services = Services::new(options.clone(), instance.clone(), Plugins::default()); + let server = services + .http_server_builder( + &options.frontend_options(), + services.server_memory_limiter.clone(), + ) + .build(); + let client = TestClient::new(server.build(server.make_app()).unwrap()).await; + + let mut ctx = QueryContext::with(DEFAULT_CATALOG_NAME, "public"); + ctx.set_logical_batching_enabled(true); + ctx.set_protocol_ctx(ProtocolCtx::OtlpMetric(OtlpMetricCtx { + with_metric_engine: true, + ..Default::default() + })); + let ctx = Arc::new(ctx); + + // A physical table with a microsecond time index and a logical table + // bound to it; OTLP requests always select the default + // (millisecond) physical table, so the destination binding differs. + for sql in [ + "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')", + "CREATE TABLE cross_alignment_total (\ + greptime_timestamp TIMESTAMP(6) NOT NULL, greptime_value DOUBLE NULL, \ + \"stream\" STRING NULL, TIME INDEX (greptime_timestamp), PRIMARY KEY (\"stream\")) \ + ENGINE = metric WITH ('on_physical_table' = 'phy_us')", + ] { + let mut output = instance.do_query(sql, ctx.clone()).await; + let result = output.remove(0); + assert!(result.is_ok(), "setup ddl failed: {result:?} — {sql}"); + } + + let submissions = servers::metrics::PENDING_ROWS_BATCH_INGEST_STAGE_ELAPSED + .with_label_values(&["submit_build_and_align"]); + let before = submissions.get_sample_count(); + let request = build_sum_request( + "cross.alignment", + AggregationTemporality::Cumulative, + &[(60, 10)], + ); + let response = client + .post("/v1/otlp/v1/metrics") + .header("content-type", "application/x-protobuf") + .body(request.encode_to_vec()) + .send() + .await; + assert_eq!(response.status().as_u16(), 200); + // The request must NOT enter the logical batcher. + assert_eq!( + submissions.get_sample_count() - before, + 0, + "cross-physical destinations must fall back to the ordinary insert path" + ); + + // The row lands on the microsecond logical table, converted from the + // nanosecond encoding: 60s -> 60_000_000us. + let sql = "SELECT greptime_timestamp, greptime_value FROM cross_alignment_total"; + let batches = tokio::time::timeout(Duration::from_secs(5), async { + loop { + let output = instance.do_query(sql, ctx.clone()).await.remove(0).unwrap(); + let OutputData::Stream(stream) = output.data else { + panic!("expected stream") + }; + let batches = RecordBatches::try_collect(stream).await.unwrap(); + if batches.iter().map(|batch| batch.num_rows()).sum::() == 1 { + break batches; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .unwrap(); + let batch = &batches.take()[0]; + let timestamps = batch.column(0).as_primitive::(); + assert_eq!(timestamps.value(0), 60_000_000); + } + #[tokio::test(flavor = "multi_thread")] async fn test_otlp_logical_batcher_alignment() { use std::time::Duration; diff --git a/tests-integration/tests/http.rs b/tests-integration/tests/http.rs index 1b68d3720a7..780140e3500 100644 --- a/tests-integration/tests/http.rs +++ b/tests-integration/tests/http.rs @@ -175,6 +175,8 @@ macro_rules! http_tests { 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_write_batched_microsecond_physical_table, + test_prometheus_remote_write_batched_conflicting_physical_selections, test_prometheus_remote_special_labels, test_prometheus_remote_schema_labels, test_prometheus_remote_write_with_pipeline, @@ -3832,8 +3834,9 @@ pub async fn test_prometheus_remote_write_batched_mixed_time_index_units(store_t }; // 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. + // path handles the microsecond selected physical table: the new table + // is created on it and the samples are widened to its unit during + // batch alignment. let res = client .post("/v1/prometheus/write?physical_table=phy_us") .header("Content-Encoding", "snappy") @@ -3843,10 +3846,11 @@ pub async fn test_prometheus_remote_write_batched_mixed_time_index_units(store_t 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. + // (millisecond) physical table: the destination is bound to another + // physical table, so the bulk eligibility check must reject it and fall + // back to the ordinary insert path — without the destination binding + // check the bulk flush would write the rows through the selected + // physical's regions, silently misplacing them. let res = client .post("/v1/prometheus/write") .header("Content-Encoding", "snappy") @@ -4069,6 +4073,224 @@ pub async fn test_prometheus_remote_write_v2_batched_interceptor_time_index_unit /// bypasses `PromStoreProtocolHandler::write`. Verifies the metric table is created /// asynchronously and still carries the Prometheus semantic identity stamped on the /// shared request context. +/// Regression test for : +/// prometheus remote write uses the logical batcher against a physical +/// metric table pre-created with a microsecond time index; the millisecond +/// samples are widened to the physical table's unit during batch alignment. +pub async fn test_prometheus_remote_write_batched_microsecond_physical_table( + store_type: StorageType, +) { + common_telemetry::init_default_ut_logging(); + let (app, mut guard) = setup_test_prom_app_with_frontend_batched( + store_type, + "prometheus_remote_write_batched_us_physical", + ) + .await; + let client = TestClient::new(app).await; + + // Pre-create the default physical metric table with a microsecond time + // index before any remote write, so the batched bulk path must handle it. + let res = client + .get( + "/v1/sql?db=public&sql=CREATE TABLE greptime_physical_table \ + (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_request = WriteRequest { + timeseries: vec![prom_store::mock_timeseries()[0].clone()], + ..Default::default() + }; + let serialized_request = write_request.encode_to_vec(); + let compressed_request = + prom_store::snappy_compress(&serialized_request).expect("failed to encode snappy"); + + // submit_build_and_align counts every batcher submission in both + // acknowledgement modes, so the stored-value assertions below cannot be + // satisfied by a silent fallback to ordinary insertion. + let submissions = servers::metrics::PENDING_ROWS_BATCH_INGEST_STAGE_ELAPSED + .with_label_values(&["submit_build_and_align"]); + let before = submissions.get_sample_count(); + + let res = client + .post("/v1/prometheus/write") + .header("Content-Encoding", "snappy") + .body(compressed_request) + .send() + .await; + assert_eq!(res.status(), StatusCode::NO_CONTENT); + assert_eq!( + submissions.get_sample_count() - before, + 1, + "non-millisecond physical tables must use the logical batcher" + ); + + // metric1 samples are 1.0@1000ms and 2.0@2000ms; on the microsecond + // physical table they must be stored as 1_000_000us and 2_000_000us. + wait_for_data( + &client, + "select greptime_timestamp, greptime_value from metric1 order by greptime_timestamp", + "[[1000000,1.0],[2000000,2.0]]", + ) + .await; + + guard.remove_all().await; +} + +/// Regression test for the per-series physical-table selection conflict: +/// one remote-write request containing the same metric under two different +/// physical-table selections must fall back to the ordinary insert path — +/// both when the destination table already exists (bound to one of the two +/// physicals) and when it is missing (conflicting creations). A consistent +/// selection must still use the batcher. +pub async fn test_prometheus_remote_write_batched_conflicting_physical_selections( + 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_conflicting_physicals", + ) + .await; + let client = TestClient::new(app).await; + + let res = client + .get( + "/v1/sql?db=public&sql=CREATE TABLE p1 (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); + // An existing logical table bound to p1. + let res = client + .get( + "/v1/sql?db=public&sql=CREATE TABLE conflict_existing (greptime_timestamp TIMESTAMP(6) NOT NULL, greptime_value DOUBLE NULL, \"job\" STRING NULL, TIME INDEX (greptime_timestamp), PRIMARY KEY (\"job\")) ENGINE = metric WITH ('on_physical_table' = 'p1')", + ) + .send() + .await; + assert_eq!(res.status(), StatusCode::OK); + + let write = |metric: &str, job: &str, physical: Option<&str>, value: f64| { + let mut labels = vec![ + Label { + name: prom_store::METRIC_NAME_LABEL.to_string(), + value: metric.to_string(), + }, + Label { + name: "job".to_string(), + value: job.to_string(), + }, + ]; + if let Some(physical) = physical { + labels.push(Label { + name: "x_greptime_physical_table".to_string(), + value: physical.to_string(), + }); + } + TimeSeries { + labels, + samples: vec![Sample { + value, + timestamp: 1000, + }], + ..Default::default() + } + }; + let send = |timeseries: Vec| { + let write_request = WriteRequest { + timeseries, + ..Default::default() + }; + prom_store::snappy_compress(&write_request.encode_to_vec()).unwrap() + }; + + let submissions = servers::metrics::PENDING_ROWS_BATCH_INGEST_STAGE_ELAPSED + .with_label_values(&["submit_build_and_align"]); + + // Existing destination under two different selections: the batch must + // fall back (the batch for the unbound selection would otherwise write + // through the default physical's regions), and both series land on the + // bound physical. + let before = submissions.get_sample_count(); + let res = client + .post("/v1/prometheus/write") + .header("Content-Encoding", "snappy") + .body(send(vec![ + write("conflict_existing", "a", Some("p1"), 1.0), + write("conflict_existing", "b", None, 2.0), + ])) + .send() + .await; + assert_eq!(res.status(), StatusCode::NO_CONTENT); + assert_eq!( + submissions.get_sample_count() - before, + 0, + "conflicting selections for an existing table must fall back" + ); + wait_for_data( + &client, + "select count(*), sum(greptime_value) from conflict_existing", + "[[2,3.0]]", + ) + .await; + + // Missing destination under two different selections: conflicting + // creations must also fall back, and both series must be visible + // through the created table. + let before = submissions.get_sample_count(); + let res = client + .post("/v1/prometheus/write") + .header("Content-Encoding", "snappy") + .body(send(vec![ + write("conflict_fresh", "a", Some("p1"), 1.0), + write("conflict_fresh", "b", None, 2.0), + ])) + .send() + .await; + assert_eq!(res.status(), StatusCode::NO_CONTENT); + assert_eq!( + submissions.get_sample_count() - before, + 0, + "conflicting selections for a missing table must fall back" + ); + wait_for_data( + &client, + "select count(*), sum(greptime_value) from conflict_fresh", + "[[2,3.0]]", + ) + .await; + + // A consistent selection still uses the batcher. + let before = submissions.get_sample_count(); + let res = client + .post("/v1/prometheus/write") + .header("Content-Encoding", "snappy") + .body(send(vec![ + write("consistent_fresh", "a", Some("p1"), 1.0), + write("consistent_fresh", "b", Some("p1"), 2.0), + ])) + .send() + .await; + assert_eq!(res.status(), StatusCode::NO_CONTENT); + assert_eq!( + submissions.get_sample_count() - before, + 1, + "a consistent selection must use the logical batcher" + ); + wait_for_data( + &client, + "select count(*), sum(greptime_value) from consistent_fresh", + "[[2,3.0]]", + ) + .await; + + guard.remove_all().await; +} + pub async fn test_prometheus_remote_write_batched(store_type: StorageType) { common_telemetry::init_default_ut_logging(); let (app, mut guard) =