From fedce5c5ec1e458483901e2dc75dab6d29b741cf Mon Sep 17 00:00:00 2001 From: Ning Sun Date: Mon, 28 Sep 2026 12:27:36 +0000 Subject: [PATCH] feat: support non-millisecond time index units in the logical batcher (#9346) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * feat: allow customized time index unit for metric engine table * test: provide query tests * refactor: revert unnecessary change * refactor: share timestamp unit conversions in api helper Address review feedback on the time index unit changeset: - Add shared timestamp_unit/timestamp_datatype helpers to api::helper (the only crate that sees both proto ColumnDataType and TimeUnit due to layering; common-time and datatypes have no greptime-proto dep). This removes the ColumnDataType -> TimeUnit match duplicated between operator's insert path and the OTLP logs path. - Collapse the two TimeUnit <-> ValueData matches in convert_timestamp_value_data by reusing api::helper::to_grpc_value for the construction side. - Note that convert_rows_time_unit rewrites the schema before the values, so an overflow mid-batch leaves the request half-converted; harmless because the error aborts the whole insert request. Signed-off-by: Ning Sun * fix: align time units per destination table and floor remote-read timestamps Address review feedback on PR #9236: - Align each metric insert request to the unit of the table it actually targets: an existing logical table keeps its own unit (it may be bound to a different physical table than the one selected by the request), and only new tables use the selected physical table's unit. The previous blanket conversion rewrote valid millisecond samples to the selected physical table's unit and the engine rejected them. Regression test: writing an existing millisecond logical table and a new table in one request that selects a microsecond physical table. - Remote read now floors narrowing timestamp conversions towards negative infinity (div_euclid), consistent with Timestamp::convert_to on the ingestion path; arrow's cast truncates towards zero and returned -1ms for a stored -1001us. Widening (second -> millisecond) keeps the exact arrow cast. Regression test: a negative, non-aligned timestamp round-trips as -2ms. Signed-off-by: Ning Sun * perf: fold time unit alignment into existing table lookups Address review feedback on PR #9236: - The per-destination unit alignment no longer runs its own pass of table lookups: create_or_alter_tables_on_demand gains an align_time_index_unit parameter (metric engine path only) and converts each request inside the lookups it already performs — existing tables to their own unit, new tables to the selected physical table's. Default ingest paths now issue zero additional catalog lookups compared to main; the separate alignment pass remains only in the opt-in logical batcher pre-gate, next to the eligibility check that already looks up the same tables. - convert_rows_time_unit indexes the time index position directly (validate_column_count_match guarantees row widths) instead of Optional get_mut; the gate-side alignment validates widths itself. Signed-off-by: Ning Sun * perf: resolve the batcher time index guard once per write target All batches of one remote write request share the same write target (catalog, schema, physical table), so the batcher time index guard now resolves each distinct target once instead of once per batch. Signed-off-by: Ning Sun * feat: support non-millisecond time index units in the logical batcher Make the logical table batcher's bulk encode path unit-aware so physical metric tables with a non-millisecond time index (e.g. TIMESTAMP(6)) can use logical batching instead of falling back to the ordinary insert path. - rows_to_aligned_record_batch builds the time index column in the TARGET schema's unit, converting any timestamp encoding via Timestamp::convert_to (flooring on narrowing, consistent with the ordinary insert path). - New tables created by the batcher use the selected physical table's time index unit (resolved once per submit; a missing physical table keeps the millisecond auto-create default). - columns_taxonomy and the can_batch_metric_rows schema whitelist accept any timestamp unit; the prometheus remote write v1/v2 batcher gates and the OTLP pre-gate alignment are removed together with Inserter::align_metric_row_inserts_time_unit, as the batcher now converts internally. Closes #9342 Signed-off-by: Ning Sun * fix: address review comments * fix: address review issue * refactor: drop the OTLP pre-gate unit alignment made redundant by the bulk path The main merge of #9236 (squash) resurrected the OTLP pre-gate alignment and Inserter::align_metric_row_inserts_time_unit, which this branch had removed. Drop them again: - The pre-gate existed because the #9236-era bulk eligibility gate only accepted millisecond schemas, so nanosecond-encoded OTLP requests had to be converted before the check. This branch makes the bulk path unit-aware (the gate accepts all time index units and batch alignment converts each request to its destination's unit), so the pre-gate is redundant and only added N+1 catalog lookups per batched request — the very lookup-count overhead raised in the #9236 review. - The per-destination unit semantics it implemented remain enforced in the two paths that need them: the ordinary insert path (create_or_alter_tables_on_demand converts inside its existing table lookups) and the batched path (batch alignment resolves each destination schema and converts to it). test_otlp_logical_batcher_alignment (the test the pre-gate originally fixed) and the mixed-physical-table regression both pass without it. Signed-off-by: Ning Sun * test: cover OTLP batcher cross-physical fallback and nanosecond physical Extend the logical batcher integration coverage for the cases previously guarded by the removed OTLP pre-gate alignment: - test_otlp_logical_batcher_fallback_for_cross_physical_destination: with the batcher enabled, an OTLP request targeting an existing logical table bound to another physical table must NOT enter the batcher (the bulk eligibility check rejects the destination binding) and the ordinary insert path must convert it to the destination's unit (60s -> 60_000_000us). - test_otlp_logical_batcher_non_millisecond_physical_table now covers both microsecond and nanosecond physical tables (parameterized), asserting batcher submissions and unit-precise stored values. Signed-off-by: Ning Sun * fix: validate per-batch physical bindings and avoid intermediate timestamp buffers Address review feedback on PR #9346: - accepts_bulk_destinations dedupes on (schema, table, selected physical) instead of (schema, table): one request can select different physical tables per batch (per-series x_greptime_physical_table labels), and the old key let a second selection skip validation and flush rows through the wrong physical's regions. Missing tables additionally reject conflicting physical selections within the same request. Regression test covers an existing destination, a missing destination, and a consistent selection (which must still batch). - The timestamp column builder appends each value directly into the target-unit Arrow builder; values already in the target unit (the unchanged millisecond fast path) are appended without conversion, so the default millisecond physical pays no Timestamp construction or intermediate Vec allocation. - The non-millisecond batching tests assert the submit_build_and_align counter, which increments on every batcher submission in both acknowledgement modes, so a silent fallback to ordinary insertion fails the tests instead of passing on stored values alone. Signed-off-by: Ning Sun --------- Signed-off-by: Ning Sun --- src/frontend/src/instance/otlp.rs | 28 +- src/operator/src/insert.rs | 141 +-------- src/servers/src/batcher/logical_table.rs | 148 +++++----- .../batcher/logical_table/batch_convert.rs | 4 +- .../src/batcher/logical_table/tables.rs | 25 +- src/servers/src/http/prom_store.rs | 20 +- src/servers/src/prom_row_builder.rs | 274 +++++++++++++++--- tests-integration/src/otlp.rs | 263 +++++++++++++++++ tests-integration/tests/http.rs | 234 ++++++++++++++- 9 files changed, 855 insertions(+), 282 deletions(-) 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) =