diff --git a/src/promql/src/extension_plan.rs b/src/promql/src/extension_plan.rs index 29b31b7ca0..a8fcd10ca4 100644 --- a/src/promql/src/extension_plan.rs +++ b/src/promql/src/extension_plan.rs @@ -28,10 +28,16 @@ mod union_distinct_on; pub use absent::{Absent, AbsentExec, AbsentStream}; use common_query::native_histogram::{SUM_FIELD, native_histogram_value_type}; use common_query::prometheus::is_prometheus_stale_nan; -use datafusion::arrow::array::{Array, Float64Array, StructArray}; -use datafusion::arrow::datatypes::{ArrowPrimitiveType, TimestampMillisecondType}; -use datafusion::common::DFSchemaRef; +use datafusion::arrow::array::{ + Array, Float64Array, StructArray, TimestampMicrosecondArray, TimestampMillisecondArray, + TimestampNanosecondArray, TimestampSecondArray, +}; +use datafusion::arrow::datatypes::{ + ArrowPrimitiveType, DataType, TimeUnit, TimestampMillisecondType, +}; +use datafusion::common::{Column, DFSchemaRef}; use datafusion::error::{DataFusionError, Result as DataFusionResult}; +use datafusion::logical_expr::{Expr, Extension, LogicalPlan}; use datatypes::data_type::DataType as _; pub use empty_metric::{EmptyMetric, EmptyMetricExec, EmptyMetricStream, build_special_time_expr}; pub use histogram_fold::{ @@ -47,6 +53,101 @@ pub use union_distinct_on::{UnionDistinctOn, UnionDistinctOnExec, UnionDistinctO pub type Millisecond = ::Native; +/// Borrows timestamp values without reducing their Arrow storage precision. +pub(crate) fn native_timestamp_values(array: &dyn Array) -> datafusion::error::Result<&[i64]> { + let value = match array.data_type() { + DataType::Timestamp(TimeUnit::Second, _) => array + .as_any() + .downcast_ref::() + .map(|a| a.values().as_ref()), + DataType::Timestamp(TimeUnit::Millisecond, _) => array + .as_any() + .downcast_ref::() + .map(|a| a.values().as_ref()), + DataType::Timestamp(TimeUnit::Microsecond, _) => array + .as_any() + .downcast_ref::() + .map(|a| a.values().as_ref()), + DataType::Timestamp(TimeUnit::Nanosecond, _) => array + .as_any() + .downcast_ref::() + .map(|a| a.values().as_ref()), + _ => None, + }; + value.ok_or_else(|| { + datafusion::error::DataFusionError::Execution("Time index column is not a timestamp".into()) + }) +} + +pub(crate) fn timestamp_unit(data_type: &DataType) -> datafusion::error::Result { + match data_type { + DataType::Timestamp(unit, _) => Ok(*unit), + _ => Err(datafusion::error::DataFusionError::Execution( + "Time index column is not a timestamp".into(), + )), + } +} + +pub(crate) fn nanoseconds_per_native_tick(unit: TimeUnit) -> i128 { + match unit { + TimeUnit::Second => 1_000_000_000, + TimeUnit::Millisecond => 1_000_000, + TimeUnit::Microsecond => 1_000, + TimeUnit::Nanosecond => 1, + } +} + +/// Returns the offset of an immediately underlying normalize node when the +/// requested time index retains its logical identity through projections. +pub(crate) fn local_offset(plan: &LogicalPlan, time_index: &str) -> Millisecond { + let Some(index) = plan.schema().index_of_column_by_name(None, time_index) else { + return 0; + }; + let (qualifier, field) = plan.schema().qualified_field(index); + let mut time_index = Column::new(qualifier.cloned(), field.name().clone()); + let mut plan = plan; + + loop { + match plan { + LogicalPlan::Extension(Extension { node }) => { + return node + .as_any() + .downcast_ref::() + .and_then(|normalize| normalize.offset_for_time_index(&time_index)) + .unwrap_or_default(); + } + LogicalPlan::Projection(projection) => { + let Some(output_index) = projection.schema.maybe_index_of_column(&time_index) + else { + return 0; + }; + let expr = &projection.expr[output_index]; + let source = match expr { + Expr::Column(column) => column, + Expr::Alias(alias) => { + let Expr::Column(column) = alias.expr.as_ref() else { + return 0; + }; + if alias.name != column.name { + return 0; + } + column + } + _ => return 0, + }; + let Some(input_index) = projection.input.schema().maybe_index_of_column(source) + else { + return 0; + }; + let (qualifier, field) = projection.input.schema().qualified_field(input_index); + time_index = Column::new(qualifier.cloned(), field.name().clone()); + plan = projection.input.as_ref(); + } + _ => return 0, + } + } +} + const METRIC_NUM_SERIES: &str = "num_series"; fn prometheus_stale_sample_column(column: &dyn Array) -> Option<(&dyn Array, &Float64Array)> { @@ -109,3 +210,87 @@ pub fn resolve_column_names( .map(|idx| resolve_column_name(*idx, schema, context, column_type)) .collect() } + +#[cfg(test)] +mod tests { + use std::sync::Arc; + + use datafusion::arrow::datatypes::{DataType, Field, Schema, TimeUnit}; + use datafusion::common::ToDFSchema; + use datafusion::logical_expr::{EmptyRelation, Extension, LogicalPlan, Projection}; + use datafusion_expr::col; + + use super::*; + + fn input() -> LogicalPlan { + LogicalPlan::EmptyRelation(EmptyRelation { + produce_one_row: false, + schema: Arc::new(Schema::new(vec![ + Field::new( + "timestamp", + DataType::Timestamp(TimeUnit::Millisecond, None), + false, + ), + Field::new( + "other_ts", + DataType::Timestamp(TimeUnit::Millisecond, None), + false, + ), + Field::new("value", DataType::Float64, true), + ])) + .to_dfschema_ref() + .unwrap(), + }) + } + + fn normalized() -> LogicalPlan { + LogicalPlan::Extension(Extension { + node: Arc::new(SeriesNormalize::new( + 1_000, + "timestamp", + false, + Vec::new(), + input(), + )), + }) + } + + #[test] + fn local_offset_tracks_identity_preserving_projections() { + let projection = + Projection::try_new(vec![col("timestamp"), col("value")], Arc::new(normalized())) + .unwrap(); + let projection = Projection::try_new( + vec![col("timestamp").alias("timestamp"), col("value")], + Arc::new(LogicalPlan::Projection(projection)), + ) + .unwrap(); + + assert_eq!( + 1_000, + local_offset(&LogicalPlan::Projection(projection), "timestamp") + ); + } + + #[test] + fn local_offset_rejects_a_different_timestamp_or_manipulator() { + let renamed = Projection::try_new( + vec![col("other_ts").alias("timestamp"), col("value")], + Arc::new(normalized()), + ) + .unwrap(); + assert_eq!( + 0, + local_offset(&LogicalPlan::Projection(renamed), "timestamp") + ); + + let divide = LogicalPlan::Extension(Extension { + node: Arc::new(SeriesDivide::new( + Vec::new(), + "timestamp".to_string(), + normalized(), + )), + }); + assert_eq!(0, local_offset(÷, "timestamp")); + } +} diff --git a/src/promql/src/extension_plan/instant_manipulate.rs b/src/promql/src/extension_plan/instant_manipulate.rs index 25d62e6bee..4951ce41e0 100644 --- a/src/promql/src/extension_plan/instant_manipulate.rs +++ b/src/promql/src/extension_plan/instant_manipulate.rs @@ -13,7 +13,6 @@ // limitations under the License. use std::any::Any; -use std::cmp::Ordering; use std::pin::Pin; use std::sync::Arc; use std::task::{Context, Poll}; @@ -29,6 +28,7 @@ use datafusion::execution::context::TaskContext; use datafusion::logical_expr::{ EmptyRelation, Expr, Extension, LogicalPlan, UserDefinedLogicalNodeCore, }; +use datafusion::physical_expr::EquivalenceProperties; use datafusion::physical_plan::metrics::{ BaselineMetrics, Count, ExecutionPlanMetricsSet, MetricBuilder, MetricValue, MetricsSet, }; @@ -46,8 +46,9 @@ use snafu::ResultExt; use crate::error::{DeserializeSnafu, Result}; use crate::extension_plan::series_divide::SeriesDivide; use crate::extension_plan::{ - METRIC_NUM_SERIES, Millisecond, is_prometheus_stale_sample, prometheus_stale_sample_column, - resolve_column_name, serialize_column_index, + METRIC_NUM_SERIES, Millisecond, is_prometheus_stale_sample, local_offset, + nanoseconds_per_native_tick, native_timestamp_values, prometheus_stale_sample_column, + resolve_column_name, serialize_column_index, timestamp_unit, }; use crate::metrics::PROMQL_SERIES_COUNT; @@ -67,7 +68,7 @@ fn mixed_sample_fields(field: Option<&str>) -> [Option<&str>; 2] { /// This plan will try to align the input time series, for every timestamp between /// `start` and `end` with step `interval`. Find in the `lookback` range if data /// is missing at the given timestamp. -#[derive(Debug, PartialEq, Eq, Hash, PartialOrd)] +#[derive(Debug, PartialEq, Eq, Hash)] pub struct InstantManipulate { start: Millisecond, end: Millisecond, @@ -79,9 +80,37 @@ pub struct InstantManipulate { /// Primary sample column used to derive the columns checked for staleness. field_column: Option, input: LogicalPlan, + output_schema: DFSchemaRef, unfix: Option, } +impl PartialOrd for InstantManipulate { + fn partial_cmp(&self, other: &Self) -> Option { + ( + self.start, + self.end, + self.lookback_delta, + self.interval, + &self.time_index_column, + &self.tag_columns, + &self.field_column, + &self.input, + &self.unfix, + ) + .partial_cmp(&( + other.start, + other.end, + other.lookback_delta, + other.interval, + &other.time_index_column, + &other.tag_columns, + &other.field_column, + &other.input, + &other.unfix, + )) + } +} + #[derive(Debug, PartialEq, Eq, Hash, PartialOrd)] struct UnfixIndices { pub time_index_idx: u64, @@ -98,7 +127,7 @@ impl UserDefinedLogicalNodeCore for InstantManipulate { } fn schema(&self) -> &DFSchemaRef { - self.input.schema() + &self.output_schema } fn expressions(&self) -> Vec { @@ -180,6 +209,7 @@ impl UserDefinedLogicalNodeCore for InstantManipulate { end: self.end, lookback_delta: self.lookback_delta, interval: self.interval, + output_schema: Self::calculate_output_schema(&input, &time_index_column)?, time_index_column, tag_columns: Self::resolve_tag_columns(&input, &self.tag_columns), field_column, @@ -195,6 +225,7 @@ impl UserDefinedLogicalNodeCore for InstantManipulate { time_index_column: self.time_index_column.clone(), tag_columns: Self::resolve_tag_columns(&input, &self.tag_columns), field_column: self.field_column.clone(), + output_schema: Self::calculate_output_schema(&input, &self.time_index_column)?, input, unfix: None, }) @@ -203,6 +234,38 @@ impl UserDefinedLogicalNodeCore for InstantManipulate { } impl InstantManipulate { + fn calculate_output_schema( + input: &LogicalPlan, + time_index_column: &str, + ) -> DataFusionResult { + let input_schema = input.schema(); + let time_index = input_schema + .index_of_column_by_name(None, time_index_column) + .ok_or_else(|| { + DataFusionError::Internal(format!( + "InstantManipulate time index {time_index_column} not found" + )) + })?; + let mut fields = (0..input_schema.fields().len()) + .map(|index| { + let (qualifier, field) = input_schema.qualified_field(index); + (qualifier.cloned(), field.clone()) + }) + .collect::>(); + let (qualifier, field) = input_schema.qualified_field(time_index); + fields[time_index] = ( + qualifier.cloned(), + Arc::new(field.as_ref().clone().with_data_type(DataType::Timestamp( + datafusion::arrow::datatypes::TimeUnit::Millisecond, + None, + ))), + ); + Ok(Arc::new(DFSchema::new_with_metadata( + fields, + input_schema.metadata().clone(), + )?)) + } + #[allow(clippy::too_many_arguments)] pub fn new( start: Millisecond, @@ -219,6 +282,8 @@ impl InstantManipulate { end, lookback_delta, interval, + output_schema: Self::calculate_output_schema(&input, &time_index_column) + .unwrap_or_else(|_| input.schema().clone()), time_index_column, tag_columns, field_column, @@ -274,7 +339,27 @@ impl InstantManipulate { pub fn to_execution_plan(&self, exec_input: Arc) -> Arc { let reuse_tsid_column = matches!(self.tag_columns.as_slice(), [tag] if tag == "__tsid"); + let mut fields = exec_input.schema().fields().to_vec(); + let time_index = exec_input + .schema() + .index_of(&self.time_index_column) + .expect("time index column not found"); + fields[time_index] = Arc::new(fields[time_index].as_ref().clone().with_data_type( + DataType::Timestamp(datafusion::arrow::datatypes::TimeUnit::Millisecond, None), + )); + let output_schema = Arc::new(datafusion::arrow::datatypes::Schema::new_with_metadata( + fields, + exec_input.schema().metadata().clone(), + )); + let input_properties = exec_input.properties(); + let properties = Arc::new(PlanProperties::new( + EquivalenceProperties::new(output_schema.clone()), + input_properties.partitioning.clone(), + input_properties.emission_type, + input_properties.boundedness, + )); Arc::new(InstantManipulateExec { + offset: local_offset(&self.input, &self.time_index_column), start: self.start, end: self.end, lookback_delta: self.lookback_delta, @@ -283,6 +368,8 @@ impl InstantManipulate { field_column: self.field_column.clone(), reuse_tsid_column, input: exec_input, + output_schema, + properties, metric: ExecutionPlanMetricsSet::new(), }) } @@ -311,9 +398,10 @@ impl InstantManipulate { pub fn deserialize(bytes: &[u8]) -> Result { let pb_instant_manipulate = pb::InstantManipulate::decode(bytes).context(DeserializeSnafu)?; + let empty_schema = Arc::new(DFSchema::empty()); let placeholder_plan = LogicalPlan::EmptyRelation(EmptyRelation { produce_one_row: false, - schema: Arc::new(DFSchema::empty()), + schema: empty_schema.clone(), }); let unfix = UnfixIndices { @@ -329,6 +417,7 @@ impl InstantManipulate { time_index_column: String::new(), tag_columns: Vec::new(), field_column: None, + output_schema: empty_schema, input: placeholder_plan, unfix: Some(unfix), }) @@ -337,6 +426,7 @@ impl InstantManipulate { #[derive(Debug)] pub struct InstantManipulateExec { + offset: Millisecond, start: Millisecond, end: Millisecond, lookback_delta: Millisecond, @@ -346,6 +436,8 @@ pub struct InstantManipulateExec { reuse_tsid_column: bool, input: Arc, + output_schema: SchemaRef, + properties: Arc, metric: ExecutionPlanMetricsSet, } @@ -355,11 +447,11 @@ impl ExecutionPlan for InstantManipulateExec { } fn schema(&self) -> SchemaRef { - self.input.schema() + self.output_schema.clone() } fn properties(&self) -> &Arc { - self.input.properties() + &self.properties } fn required_input_distribution(&self) -> Vec { @@ -380,7 +472,16 @@ impl ExecutionPlan for InstantManipulateExec { children: Vec>, ) -> DataFusionResult> { assert!(!children.is_empty()); + let input = children[0].clone(); + let input_properties = input.properties(); + let properties = Arc::new(PlanProperties::new( + EquivalenceProperties::new(self.output_schema.clone()), + input_properties.partitioning.clone(), + input_properties.emission_type, + input_properties.boundedness, + )); Ok(Arc::new(Self { + offset: self.offset, start: self.start, end: self.end, lookback_delta: self.lookback_delta, @@ -388,7 +489,9 @@ impl ExecutionPlan for InstantManipulateExec { time_index_column: self.time_index_column.clone(), field_column: self.field_column.clone(), reuse_tsid_column: self.reuse_tsid_column, - input: children[0].clone(), + input, + output_schema: self.output_schema.clone(), + properties, metric: self.metric.clone(), })) } @@ -413,6 +516,7 @@ impl ExecutionPlan for InstantManipulateExec { .column_with_name(&self.time_index_column) .expect("time index column not found") .0; + let time_unit = timestamp_unit(schema.field(time_index).data_type())?; let field_indices = mixed_sample_fields(self.field_column.as_deref()).map(|field| { field.and_then(|field| schema.column_with_name(field).map(|(index, _)| index)) }); @@ -421,15 +525,17 @@ impl ExecutionPlan for InstantManipulateExec { .filter(|(_, field)| field.data_type() == &DataType::UInt64) .map(|(index, _)| index); Ok(Box::pin(InstantManipulateStream { + offset: self.offset, start: self.start, end: self.end, lookback_delta: self.lookback_delta, interval: self.interval, time_index, + time_unit, field_indices, tsid_index, reuse_tsid_column: self.reuse_tsid_column && tsid_index.is_some(), - schema, + schema: self.output_schema.clone(), input, metric: baseline_metric, num_series, @@ -487,12 +593,14 @@ impl DisplayAs for InstantManipulateExec { } pub struct InstantManipulateStream { + offset: Millisecond, start: Millisecond, end: Millisecond, lookback_delta: Millisecond, interval: Millisecond, // Column index of TIME INDEX column's position in schema time_index: usize, + time_unit: datafusion::arrow::datatypes::TimeUnit, field_indices: [Option; 2], tsid_index: Option, reuse_tsid_column: bool, @@ -516,9 +624,6 @@ impl Stream for InstantManipulateStream { fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { let poll = match ready!(self.input.poll_next_unpin(cx)) { Some(Ok(batch)) => { - if batch.num_rows() == 0 { - return Poll::Pending; - } let timer = std::time::Instant::now(); self.num_series.add(1); let result = Ok(batch).and_then(|batch| self.manipulate(batch)); @@ -543,22 +648,11 @@ impl InstantManipulateStream { /// lookback window `(eval_ts - lookback_delta, eval_ts]`; a sample at exactly /// `eval_ts - lookback_delta` is too old. pub fn manipulate(&self, input: RecordBatch) -> DataFusionResult { - let ts_column = input - .column(self.time_index) - .as_any() - .downcast_ref::() - .ok_or_else(|| { - DataFusionError::Execution( - "Time index Column downcast to TimestampMillisecondArray failed".into(), - ) - })?; - - // Early return for empty input + let ts_column = input.column(self.time_index); if ts_column.is_empty() { - return Ok(input); + return Ok(RecordBatch::new_empty(self.schema.clone())); } - - // Field columns for staleness checks, classified once per batch. + let scale = nanoseconds_per_native_tick(self.time_unit); let stale_sample_columns = self.field_indices.map(|index| { index.and_then(|index| prometheus_stale_sample_column(input.column(index).as_ref())) }); @@ -568,101 +662,72 @@ impl InstantManipulateStream { .flatten() .any(|column| is_prometheus_stale_sample(*column, row)) }; - - // Optimize iteration range based on actual data bounds - let first_ts = ts_column.value(0); - let last_ts = ts_column.value(ts_column.len() - 1); - // A sample at `t` is eligible for eval time `eval_ts` iff: - // t > eval_ts - lookback_delta <=> eval_ts < t + lookback_delta. - // Therefore the last eval timestamp for which the last sample is still eligible is: - // last_ts + lookback_delta - 1 (millisecond granularity). - let last_useful = if self.lookback_delta > 0 { - last_ts + self.lookback_delta - 1 + let timestamps = native_timestamp_values(ts_column.as_ref())?; + let len = timestamps.len(); + let to_nanoseconds = + |timestamp: i64| (timestamp as i128) * scale + (self.offset as i128) * 1_000_000; + let first_ns = to_nanoseconds(timestamps[0]); + let last_ns = to_nanoseconds(timestamps[len - 1]); + // An exact sample remains useful with zero lookback. Otherwise the lower + // boundary is exclusive, so subtract one nanosecond from its final window. + let last_useful = if self.lookback_delta == 0 { + last_ns } else { - last_ts + last_ns + (self.lookback_delta as i128) * 1_000_000 - 1 + }; + let first_ms = (first_ns + 999_999).div_euclid(1_000_000); + let last_ms = last_useful.div_euclid(1_000_000); + let query_start = self.start as i128; + let query_end = self.end as i128; + let interval = self.interval as i128; + let max_start = first_ms.max(query_start); + let min_end = last_ms.min(query_end); + let (aligned_start, aligned_end) = if max_start > min_end { + (1, 0) + } else { + ( + query_start + (max_start - query_start) / interval * interval, + query_end - (query_end - min_end) / interval * interval, + ) }; - - let max_start = first_ts.max(self.start); - let min_end = last_useful.min(self.end); - - let aligned_start = self.start + (max_start - self.start) / self.interval * self.interval; - let aligned_end = self.end - (self.end - min_end) / self.interval * self.interval; - let estimated_points = if aligned_end >= aligned_start { - ((aligned_end - aligned_start) / self.interval).saturating_add(1) as usize + (aligned_end - aligned_start) / interval + 1 } else { 0 }; - if estimated_points > MAX_INSTANT_MANIPULATE_OUTPUT_POINTS { + if estimated_points > MAX_INSTANT_MANIPULATE_OUTPUT_POINTS as i128 { return Err(DataFusionError::Execution(format!( "InstantManipulate output points exceed limit: {estimated_points} > {MAX_INSTANT_MANIPULATE_OUTPUT_POINTS}" ))); } + let estimated_points = estimated_points as usize; + let aligned_start = aligned_start as i64; + let aligned_end = aligned_end as i64; let mut take_indices = Vec::with_capacity(estimated_points); - - let mut cursor = 0; - - let aligned_ts_iter = (aligned_start..=aligned_end).step_by(self.interval as usize); let mut aligned_ts = Vec::with_capacity(estimated_points); - - // calculate the offsets to take - 'next: for expected_ts in aligned_ts_iter { - // first, search toward end to see if there is matched timestamp - while cursor < ts_column.len() { - let curr = ts_column.value(cursor); - match curr.cmp(&expected_ts) { - Ordering::Equal => { - if is_stale(cursor) { - // Ignore the stale marker. - } else { - take_indices.push(cursor as u64); - aligned_ts.push(expected_ts); - } - continue 'next; - } - Ordering::Greater => break, - Ordering::Less => {} + let mut cursor = 0; + for expected_ms in (aligned_start..=aligned_end).step_by(self.interval as usize) { + let expected = (expected_ms as i128) * 1_000_000; + let mut exact_candidate = None; + while cursor < len && to_nanoseconds(timestamps[cursor]) <= expected { + if to_nanoseconds(timestamps[cursor]) == expected && exact_candidate.is_none() { + exact_candidate = Some(cursor); } cursor += 1; } - if cursor == ts_column.len() { - cursor -= 1; - // short cut this loop - if ts_column.value(cursor) + self.lookback_delta <= expected_ts { - break; - } - } - - // then examine the value - let curr_ts = ts_column.value(cursor); - if curr_ts + self.lookback_delta <= expected_ts { + let Some(candidate) = exact_candidate.or_else(|| cursor.checked_sub(1)) else { continue; - } - if curr_ts > expected_ts { - // exceeds current expected timestamp, examine the previous value - if let Some(prev_cursor) = cursor.checked_sub(1) { - let prev_ts = ts_column.value(prev_cursor); - if prev_ts + self.lookback_delta > expected_ts { - // only use the point in the time range - if is_stale(prev_cursor) { - // Do not use a stale marker as the newest value. - continue; - } - // use this point - take_indices.push(prev_cursor as u64); - aligned_ts.push(expected_ts); - } - } - } else if is_stale(cursor) { - // Do not use a stale marker as the newest value. - } else { - // use this point - take_indices.push(cursor as u64); - aligned_ts.push(expected_ts); + }; + let candidate_ts = to_nanoseconds(timestamps[candidate]); + let lower = expected - (self.lookback_delta as i128) * 1_000_000; + if (candidate_ts == expected || candidate_ts > lower) + && candidate_ts <= expected + && !is_stale(candidate) + { + take_indices.push(candidate as u64); + aligned_ts.push(expected_ms); } } - - // take record batch and replace the time index column self.take_record_batch_optional(input, take_indices, aligned_ts) } @@ -696,7 +761,7 @@ impl InstantManipulateStream { arrays.push(compute::take(array, indices_array, None)?); } - let result = RecordBatch::try_new(record_batch.schema(), arrays) + let result = RecordBatch::try_new(self.schema.clone(), arrays) .map_err(|e| DataFusionError::ArrowError(Box::new(e), None))?; Ok(result) } @@ -718,13 +783,17 @@ fn reuse_constant_column(array: &Arc, len: usize) -> DataFusionResult mod test { use common_query::native_histogram::build_histogram_array; use common_query::prometheus::PROMETHEUS_STALE_NAN_BITS; - use datafusion::arrow::array::Float64Array; + use datafusion::arrow::array::{ + Float64Array, TimestampMicrosecondArray, TimestampNanosecondArray, TimestampSecondArray, + }; use datafusion::arrow::buffer::NullBuffer; - use datafusion::arrow::datatypes::{DataType, Field, Schema}; + use datafusion::arrow::datatypes::{DataType, Field, Schema, TimeUnit}; use datafusion::common::ToDFSchema; use datafusion::datasource::memory::MemorySourceConfig; use datafusion::datasource::source::DataSourceExec; - use datafusion::logical_expr::{EmptyRelation, LogicalPlan}; + use datafusion::logical_expr::{ + EmptyRelation, Extension, LogicalPlan, UserDefinedLogicalNodeCore, + }; use datafusion::prelude::SessionContext; use super::*; @@ -746,6 +815,7 @@ mod test { Arc::new(prepare_test_data()) }; let normalize_exec = Arc::new(InstantManipulateExec { + offset: 0, start, end, lookback_delta, @@ -753,6 +823,8 @@ mod test { time_index_column: TIME_INDEX_COLUMN.to_string(), field_column: Some("value".to_string()), reuse_tsid_column: false, + output_schema: memory_exec.schema(), + properties: memory_exec.properties().clone(), input: memory_exec, metric: ExecutionPlanMetricsSet::new(), }); @@ -767,6 +839,255 @@ mod test { assert_eq!(result_literal, expected); } + #[tokio::test] + async fn native_timestamps_select_exact_samples_and_keep_ms_output() { + for (unit, ticks_per_ms) in [ + (TimeUnit::Microsecond, 1_000_i64), + (TimeUnit::Nanosecond, 1_000_000_i64), + ] { + let lower = 1_000 * ticks_per_ms; + let upper = 1_001 * ticks_per_ms; + let stale = f64::from_bits(PROMETHEUS_STALE_NAN_BITS); + for (name, timestamps, values, expected_timestamps, expected_values) in [ + ( + "exact upper sample", + vec![lower + 1, upper], + vec![1.0, 2.0], + vec![1_001], + vec![2.0], + ), + ( + "exclusive lower boundary and future sample", + vec![lower, upper + 1], + vec![1.0, 2.0], + vec![1_000], + vec![1.0], + ), + ( + "one native tick above lower boundary", + vec![lower + 1, upper + 1], + vec![1.0, 2.0], + vec![1_001], + vec![1.0], + ), + ( + "future stale marker does not suppress", + vec![lower + 1, upper + 1], + vec![1.0, stale], + vec![1_001], + vec![1.0], + ), + ( + "latest in-window stale marker suppresses", + vec![lower + 1, lower + 2, upper + 1], + vec![1.0, stale, 3.0], + vec![], + vec![], + ), + ] { + let schema = Arc::new(Schema::new(vec![ + Field::new(TIME_INDEX_COLUMN, DataType::Timestamp(unit, None), false), + Field::new("value", DataType::Float64, true), + ])); + let time: Arc = match unit { + TimeUnit::Microsecond => Arc::new(TimestampMicrosecondArray::from(timestamps)), + TimeUnit::Nanosecond => Arc::new(TimestampNanosecondArray::from(timestamps)), + _ => unreachable!(), + }; + let batch = RecordBatch::try_new( + schema.clone(), + vec![time, Arc::new(Float64Array::from(values))], + ) + .unwrap(); + let logical_input = LogicalPlan::EmptyRelation(EmptyRelation { + produce_one_row: false, + schema: schema.clone().to_dfschema_ref().unwrap(), + }); + let plan = InstantManipulate::new( + 1_000, + 1_001, + 1, + 1, + TIME_INDEX_COLUMN.to_string(), + Vec::new(), + Some("value".to_string()), + logical_input.clone(), + ); + let output_schema = Arc::new(Schema::new(vec![ + Field::new( + TIME_INDEX_COLUMN, + DataType::Timestamp(TimeUnit::Millisecond, None), + false, + ), + Field::new("value", DataType::Float64, true), + ])); + assert_eq!(plan.schema().as_arrow(), output_schema.as_ref()); + + let rebuilt = InstantManipulate::deserialize(&plan.serialize()) + .unwrap() + .with_exprs_and_inputs(vec![], vec![logical_input]) + .unwrap(); + assert_eq!(rebuilt.schema(), plan.schema()); + assert_eq!(rebuilt.input.schema().as_arrow(), schema.as_ref()); + + let input = Arc::new(DataSourceExec::new(Arc::new( + MemorySourceConfig::try_new(&[vec![batch]], schema.clone(), None).unwrap(), + ))); + let exec = rebuilt.to_execution_plan(input); + assert_eq!(exec.schema(), output_schema); + assert_eq!(exec.children()[0].schema(), schema); + + let batches = + datafusion::physical_plan::collect(exec, SessionContext::default().task_ctx()) + .await + .unwrap(); + assert_eq!(batches.len(), 1, "{unit:?}: {name}"); + let output = &batches[0]; + assert_eq!(output.schema(), output_schema); + let timestamps = output + .column(0) + .as_any() + .downcast_ref::() + .unwrap(); + let values = output + .column(1) + .as_any() + .downcast_ref::() + .unwrap(); + assert_eq!( + timestamps.values().as_ref(), + expected_timestamps.as_slice(), + "{unit:?}: {name}" + ); + assert_eq!( + values.values().as_ref(), + expected_values.as_slice(), + "{unit:?}: {name}" + ); + assert_eq!(values.null_count(), 0, "{unit:?}: {name}"); + } + } + } + + #[tokio::test] + async fn logical_normalize_offset_survives_rebuild_and_executes() { + for (name, time_unit, raw, offset, start, lookback_delta) in [ + ( + "millisecond offset", + TimeUnit::Millisecond, + 0, + 1_000, + 1_000, + 0, + ), + ( + "second timestamp with negative fractional offset", + TimeUnit::Second, + 1, + -500, + 500, + 0, + ), + ] { + let schema = Arc::new(Schema::new(vec![ + Field::new( + TIME_INDEX_COLUMN, + DataType::Timestamp(time_unit, None), + false, + ), + Field::new("value", DataType::Float64, true), + ])); + let input = LogicalPlan::EmptyRelation(EmptyRelation { + produce_one_row: false, + schema: schema.clone().to_dfschema_ref().unwrap(), + }); + let normalize = crate::extension_plan::SeriesNormalize::new( + offset, + TIME_INDEX_COLUMN, + false, + Vec::new(), + input.clone(), + ); + let normalize = + crate::extension_plan::SeriesNormalize::deserialize(&normalize.serialize()) + .unwrap() + .with_exprs_and_inputs(vec![], vec![input]) + .unwrap(); + let normalized = LogicalPlan::Extension(Extension { + node: Arc::new(normalize), + }); + let plan = InstantManipulate::new( + start, + start, + lookback_delta, + 1, + TIME_INDEX_COLUMN.to_string(), + Vec::new(), + Some("value".to_string()), + normalized.clone(), + ); + let rebuilt = InstantManipulate::deserialize(&plan.serialize()) + .unwrap() + .with_exprs_and_inputs(vec![], vec![normalized]) + .unwrap(); + let timestamp: Arc = match time_unit { + TimeUnit::Millisecond => Arc::new(TimestampMillisecondArray::from(vec![raw])), + TimeUnit::Second => Arc::new(TimestampSecondArray::from(vec![raw])), + _ => unreachable!(), + }; + let batch = RecordBatch::try_new( + schema.clone(), + vec![timestamp, Arc::new(Float64Array::from(vec![7.0]))], + ) + .unwrap(); + let exec_input = Arc::new(DataSourceExec::new(Arc::new( + MemorySourceConfig::try_new(&[vec![batch]], schema, None).unwrap(), + ))); + let output = datafusion::physical_plan::collect( + rebuilt.to_execution_plan(exec_input), + SessionContext::default().task_ctx(), + ) + .await + .unwrap(); + let output = &output[0]; + assert_eq!(output.num_rows(), 1, "{name}"); + assert_eq!( + output + .column(0) + .as_any() + .downcast_ref::() + .unwrap() + .value(0), + start, + "{name}" + ); + assert_eq!( + output + .column(1) + .as_any() + .downcast_ref::() + .unwrap() + .value(0), + 7.0, + "{name}" + ); + } + } + + #[test] + fn deserialized_ordering_preserves_column_indices() { + let mut wire = pb::InstantManipulate::default(); + let first = InstantManipulate::deserialize(&wire.encode_to_vec()).unwrap(); + wire.time_index_idx = 1; + let second = InstantManipulate::deserialize(&wire.encode_to_vec()).unwrap(); + assert_ne!(first, second); + assert_eq!(first.partial_cmp(&second), Some(std::cmp::Ordering::Less)); + wire.field_index_idx = 2; + let third = InstantManipulate::deserialize(&wire.encode_to_vec()).unwrap(); + assert_ne!(second, third); + assert_eq!(second.partial_cmp(&third), Some(std::cmp::Ordering::Less)); + } + #[test] fn pruning_should_keep_time_and_field_columns_for_exec() { let df_schema = prepare_test_data().schema().to_dfschema_ref().unwrap(); @@ -927,6 +1248,7 @@ mod test { MemorySourceConfig::try_new(&[vec![batch]], schema, None).unwrap(), ))); let normalize_exec = Arc::new(InstantManipulateExec { + offset: 0, start: 0, end: 1_500, lookback_delta: 1_000, @@ -934,6 +1256,8 @@ mod test { time_index_column: TIME_INDEX_COLUMN.to_string(), field_column: Some("value".to_string()), reuse_tsid_column: true, + output_schema: input.schema(), + properties: input.properties().clone(), input, metric: ExecutionPlanMetricsSet::new(), }); @@ -988,6 +1312,7 @@ mod test { MemorySourceConfig::try_new(&[vec![batch]], schema, None).unwrap(), ))); let normalize_exec = Arc::new(InstantManipulateExec { + offset: 0, start: 0, end: 1_500, lookback_delta: 1_000, @@ -995,6 +1320,8 @@ mod test { time_index_column: TIME_INDEX_COLUMN.to_string(), field_column: Some("value".to_string()), reuse_tsid_column: true, + output_schema: input.schema(), + properties: input.properties().clone(), input, metric: ExecutionPlanMetricsSet::new(), }); @@ -1042,6 +1369,7 @@ mod test { ))); let too_many_points = MAX_INSTANT_MANIPULATE_OUTPUT_POINTS as Millisecond + 1; let normalize_exec = Arc::new(InstantManipulateExec { + offset: 0, start: 0, end: too_many_points, lookback_delta: too_many_points + 1, @@ -1049,6 +1377,8 @@ mod test { time_index_column: TIME_INDEX_COLUMN.to_string(), field_column: Some("value".to_string()), reuse_tsid_column: false, + output_schema: input.schema(), + properties: input.properties().clone(), input, metric: ExecutionPlanMetricsSet::new(), }); @@ -1336,6 +1666,241 @@ mod test { .await; } + #[test] + fn exact_ties_select_first_and_lookback_uses_latest() { + for (values, expected_timestamp) in [ + (vec![42.0, f64::from_bits(PROMETHEUS_STALE_NAN_BITS)], 1_000), + (vec![f64::from_bits(PROMETHEUS_STALE_NAN_BITS), 42.0], 1_050), + ] { + let schema = Arc::new(Schema::new(vec![ + Field::new( + TIME_INDEX_COLUMN, + DataType::Timestamp(TimeUnit::Millisecond, None), + false, + ), + Field::new("value", DataType::Float64, true), + ])); + let input = RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(TimestampMillisecondArray::from(vec![1_000, 1_000])), + Arc::new(Float64Array::from(values)), + ], + ) + .unwrap(); + let stream = InstantManipulateStream { + offset: 0, + start: 1_000, + end: 1_050, + lookback_delta: 100, + interval: 50, + time_index: 0, + time_unit: TimeUnit::Millisecond, + field_indices: [Some(1), None], + tsid_index: None, + reuse_tsid_column: false, + schema: schema.clone(), + input: Box::pin( + datafusion::physical_plan::memory::MemoryStream::try_new(vec![], schema, None) + .unwrap(), + ), + metric: BaselineMetrics::new(&ExecutionPlanMetricsSet::new(), 0), + num_series: Count::new(), + }; + + let output = stream.manipulate(input).unwrap(); + let timestamps = output + .column(0) + .as_any() + .downcast_ref::() + .unwrap(); + let values = output + .column(1) + .as_any() + .downcast_ref::() + .unwrap(); + assert_eq!(timestamps.values(), &[expected_timestamp]); + assert_eq!(values.values(), &[42.0]); + } + } + + #[test] + fn empty_batch_uses_declared_output_schema() { + let input_schema = Arc::new(Schema::new(vec![Field::new( + TIME_INDEX_COLUMN, + DataType::Timestamp(TimeUnit::Second, None), + false, + )])); + let output_schema = Arc::new(Schema::new(vec![Field::new( + TIME_INDEX_COLUMN, + DataType::Timestamp(TimeUnit::Millisecond, None), + false, + )])); + let input = RecordBatch::new_empty(input_schema.clone()); + let stream = InstantManipulateStream { + offset: 0, + start: 0, + end: 0, + lookback_delta: 0, + interval: 1, + time_index: 0, + time_unit: TimeUnit::Second, + field_indices: [None, None], + tsid_index: None, + reuse_tsid_column: false, + schema: output_schema.clone(), + input: Box::pin( + datafusion::physical_plan::memory::MemoryStream::try_new( + vec![], + input_schema, + None, + ) + .unwrap(), + ), + metric: BaselineMetrics::new(&ExecutionPlanMetricsSet::new(), 0), + num_series: Count::new(), + }; + + let output = stream.manipulate(input).unwrap(); + assert_eq!(output.schema(), output_schema); + assert_eq!( + output.schema().field(0).data_type(), + &DataType::Timestamp(TimeUnit::Millisecond, None) + ); + } + + #[test] + fn extreme_alignment_retains_exact_sample() { + let schema = Arc::new(Schema::new(vec![ + Field::new( + TIME_INDEX_COLUMN, + DataType::Timestamp(TimeUnit::Millisecond, None), + false, + ), + Field::new("value", DataType::Float64, true), + ])); + let input = RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(TimestampMillisecondArray::from(vec![i64::MAX])), + Arc::new(Float64Array::from(vec![7.0])), + ], + ) + .unwrap(); + let stream = InstantManipulateStream { + offset: 0, + start: i64::MIN + 1, + end: i64::MAX, + lookback_delta: 0, + interval: i64::MAX, + time_index: 0, + time_unit: TimeUnit::Millisecond, + field_indices: [Some(1), None], + tsid_index: None, + reuse_tsid_column: false, + schema: schema.clone(), + input: Box::pin( + datafusion::physical_plan::memory::MemoryStream::try_new(vec![], schema, None) + .unwrap(), + ), + metric: BaselineMetrics::new(&ExecutionPlanMetricsSet::new(), 0), + num_series: Count::new(), + }; + + let output = stream.manipulate(input).unwrap(); + let timestamps = output + .column(0) + .as_any() + .downcast_ref::() + .unwrap(); + let values = output + .column(1) + .as_any() + .downcast_ref::() + .unwrap(); + assert_eq!(timestamps.values(), &[i64::MAX]); + assert_eq!(values.values(), &[7.0]); + } + + #[test] + fn native_nanosecond_offset_uses_wide_shifted_timeline() { + for (raw, offset, eval) in [ + ( + 9_223_112_837_000_000_000_i64, + 259_200_000, + 9_223_372_037_000, + ), + ( + -9_223_112_837_000_000_000_i64, + -259_200_000, + -9_223_372_037_000, + ), + ] { + let schema = Arc::new(Schema::new(vec![ + Field::new( + TIME_INDEX_COLUMN, + DataType::Timestamp(TimeUnit::Nanosecond, None), + false, + ), + Field::new("value", DataType::Float64, true), + ])); + let input = RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(TimestampNanosecondArray::from(vec![raw])), + Arc::new(Float64Array::from(vec![7.0])), + ], + ) + .unwrap(); + let stream = InstantManipulateStream { + offset, + start: eval, + end: eval, + lookback_delta: 300_000, + interval: 1, + time_index: 0, + time_unit: TimeUnit::Nanosecond, + field_indices: [Some(1), None], + tsid_index: None, + reuse_tsid_column: false, + schema: Arc::new(Schema::new(vec![ + Field::new( + TIME_INDEX_COLUMN, + DataType::Timestamp(TimeUnit::Millisecond, None), + false, + ), + Field::new("value", DataType::Float64, true), + ])), + input: Box::pin( + datafusion::physical_plan::memory::MemoryStream::try_new(vec![], schema, None) + .unwrap(), + ), + metric: BaselineMetrics::new(&ExecutionPlanMetricsSet::new(), 0), + num_series: Count::new(), + }; + let output = stream.manipulate(input).unwrap(); + assert_eq!(output.num_rows(), 1); + assert_eq!( + output + .column(0) + .as_any() + .downcast_ref::() + .unwrap() + .value(0), + eval + ); + assert_eq!( + output + .column(1) + .as_any() + .downcast_ref::() + .unwrap() + .value(0), + 7.0 + ); + } + } + #[tokio::test] async fn ordinary_nan_is_selected_for_exact_and_lookback() { let schema = Arc::new(Schema::new(vec![ @@ -1358,6 +1923,7 @@ mod test { MemorySourceConfig::try_new(&[vec![batch]], schema, None).unwrap(), ))); let exec = Arc::new(InstantManipulateExec { + offset: 0, start: 1_000, end: 1_500, lookback_delta: 1_000, @@ -1365,6 +1931,8 @@ mod test { time_index_column: TIME_INDEX_COLUMN.to_string(), field_column: Some("value".to_string()), reuse_tsid_column: false, + output_schema: input.schema(), + properties: input.properties().clone(), input, metric: ExecutionPlanMetricsSet::new(), }); @@ -1437,6 +2005,7 @@ mod test { MemorySourceConfig::try_new(&[vec![batch]], schema, None).unwrap(), ))); let exec = Arc::new(InstantManipulateExec { + offset: 0, start: 750, end: 1_500, lookback_delta: 1_001, @@ -1444,6 +2013,8 @@ mod test { time_index_column: TIME_INDEX_COLUMN.to_string(), field_column: Some("value".to_string()), reuse_tsid_column: false, + output_schema: input.schema(), + properties: input.properties().clone(), input, metric: ExecutionPlanMetricsSet::new(), }); @@ -1500,6 +2071,7 @@ mod test { MemorySourceConfig::try_new(&[vec![batch]], schema, None).unwrap(), ))); let exec = Arc::new(InstantManipulateExec { + offset: 0, start: 1_000, end: 1_500, lookback_delta: 1_001, @@ -1507,6 +2079,8 @@ mod test { time_index_column: TIME_INDEX_COLUMN.to_string(), field_column: Some("value".to_string()), reuse_tsid_column: false, + output_schema: input.schema(), + properties: input.properties().clone(), input, metric: ExecutionPlanMetricsSet::new(), }); @@ -1547,6 +2121,7 @@ mod test { MemorySourceConfig::try_new(&[vec![batch]], schema, None).unwrap(), ))); let exec = Arc::new(InstantManipulateExec { + offset: 0, start: 1_000, end: 1_500, lookback_delta: 1_000, @@ -1554,6 +2129,8 @@ mod test { time_index_column: TIME_INDEX_COLUMN.to_string(), field_column: Some("value".to_string()), reuse_tsid_column: false, + output_schema: input.schema(), + properties: input.properties().clone(), input, metric: ExecutionPlanMetricsSet::new(), }); diff --git a/src/promql/src/extension_plan/normalize.rs b/src/promql/src/extension_plan/normalize.rs index e3410be926..ea2ea875c7 100644 --- a/src/promql/src/extension_plan/normalize.rs +++ b/src/promql/src/extension_plan/normalize.rs @@ -20,7 +20,7 @@ use std::task::{Context, Poll}; use common_query::native_histogram::{START_TIMESTAMP_FIELD, native_histogram_arrow_type}; use datafusion::arrow::array::{Array, BooleanArray, StructArray}; use datafusion::arrow::compute; -use datafusion::common::{DFSchema, DFSchemaRef, Result as DataFusionResult, Statistics}; +use datafusion::common::{Column, DFSchema, DFSchemaRef, Result as DataFusionResult, Statistics}; use datafusion::error::DataFusionError; use datafusion::execution::context::TaskContext; use datafusion::logical_expr::{EmptyRelation, Expr, LogicalPlan, UserDefinedLogicalNodeCore}; @@ -52,7 +52,7 @@ use crate::metrics::PROMQL_SERIES_COUNT; /// the input batch only contains sample points from one time series. /// /// Roughly speaking, this method does these things: -/// - bias sample and native histogram start timestamps by offset +/// - retain raw native sample timestamps while biasing native histogram start timestamps by offset /// - sort the record batch based on timestamp column /// - remove Prometheus stale markers (optional) #[derive(Debug, PartialEq, Eq, Hash, PartialOrd)] @@ -179,6 +179,14 @@ impl UserDefinedLogicalNodeCore for SeriesNormalize { } impl SeriesNormalize { + pub(crate) fn offset_for_time_index(&self, time_index: &Column) -> Option { + let index = self.input.schema().maybe_index_of_column(time_index)?; + let (qualifier, field) = self.input.schema().qualified_field(index); + (field.name() == &self.time_index_column_name + && time_index == &Column::new(qualifier.cloned(), field.name().clone())) + .then_some(self.offset) + } + pub fn new>( offset: Millisecond, time_index_column_name: N, @@ -334,13 +342,8 @@ impl ExecutionPlan for SeriesNormalizeExec { let input = self.input.execute(partition, context)?; let schema = input.schema(); - let time_index = schema - .column_with_name(&self.time_index_column_name) - .expect("time index column not found") - .0; Ok(Box::pin(SeriesNormalizeStream { offset: self.offset, - time_index, filter_stale_markers: self.filter_stale_markers, schema, input, @@ -380,8 +383,6 @@ impl DisplayAs for SeriesNormalizeExec { pub struct SeriesNormalizeStream { offset: Millisecond, - // Column index of TIME INDEX column's position in schema - time_index: usize, filter_stale_markers: bool, schema: SchemaRef, @@ -393,33 +394,12 @@ pub struct SeriesNormalizeStream { impl SeriesNormalizeStream { pub fn normalize(&self, input: RecordBatch) -> DataFusionResult { - let ts_column = input - .column(self.time_index) - .as_any() - .downcast_ref::() - .ok_or_else(|| { - DataFusionError::Execution( - "Time index Column downcast to TimestampMillisecondArray failed".into(), - ) - })?; - - let bias_timestamp = |timestamp: i64| { - timestamp.checked_add(self.offset).ok_or_else(|| { - DataFusionError::Execution("SeriesNormalize: timestamp offset overflow".into()) - }) - }; - - // bias the timestamp column by offset - let ts_column_biased = if self.offset == 0 { - Arc::new(ts_column.clone()) as _ - } else { - Arc::new(ts_column.try_unary::<_, TimestampMillisecondType, _>(&bias_timestamp)?) - }; + // Native sample timestamps remain raw. Manipulators apply the selector offset + // in wide nanosecond arithmetic, avoiding overflow in native Arrow storage. let mut columns = input.columns().to_vec(); - columns[self.time_index] = ts_column_biased; - // Offset selectors move samples into the evaluation timeline. Keep native histogram - // start timestamps on the same timeline for rate and reset calculations. + // Offset selectors move native histogram start timestamps onto the evaluation + // timeline for rate and reset calculations. These payloads are milliseconds. if self.offset != 0 { let native_histogram_type = native_histogram_arrow_type(); for column in &mut columns { @@ -444,7 +424,11 @@ impl SeriesNormalizeStream { if timestamp == 0 { Ok(0) } else { - bias_timestamp(timestamp) + timestamp.checked_add(self.offset).ok_or_else(|| { + DataFusionError::Execution( + "SeriesNormalize: histogram timestamp offset overflow".into(), + ) + }) } })?; // Replace only the start timestamp child to preserve the histogram payload and @@ -516,10 +500,12 @@ impl Stream for SeriesNormalizeStream { mod test { use common_query::native_histogram::{build_histogram_array, read_histogram}; use common_query::prometheus::PROMETHEUS_STALE_NAN_BITS; - use datafusion::arrow::array::Float64Array; + use datafusion::arrow::array::{ + DictionaryArray, Float64Array, TimestampMicrosecondArray, TimestampNanosecondArray, + }; use datafusion::arrow::buffer::NullBuffer; use datafusion::arrow::datatypes::{ - ArrowPrimitiveType, DataType, Field, Schema, TimestampMillisecondType, + ArrowPrimitiveType, DataType, Field, Int64Type, Schema, TimeUnit, TimestampMillisecondType, }; use datafusion::common::ToDFSchema; use datafusion::datasource::memory::MemorySourceConfig; @@ -530,7 +516,9 @@ mod test { use datatypes::arrow_array::StringArray; use super::*; + use crate::extension_plan::RangeManipulate; use crate::extension_plan::test_util::native_histogram; + use crate::range_array::RangeArray; const TIME_INDEX_COLUMN: &str = "timestamp"; @@ -630,11 +618,11 @@ mod test { "+---------------------+--------+------+\ \n| timestamp | value | path |\ \n+---------------------+--------+------+\ - \n| 1970-01-01T00:01:01 | 0.0 | foo |\ - \n| 1970-01-01T00:02:01 | 1.0 | foo |\ - \n| 1970-01-01T00:00:01 | 10.0 | foo |\ - \n| 1970-01-01T00:00:31 | 100.0 | foo |\ - \n| 1970-01-01T00:01:31 | 1000.0 | foo |\ + \n| 1970-01-01T00:01:00 | 0.0 | foo |\ + \n| 1970-01-01T00:02:00 | 1.0 | foo |\ + \n| 1970-01-01T00:00:00 | 10.0 | foo |\ + \n| 1970-01-01T00:00:30 | 100.0 | foo |\ + \n| 1970-01-01T00:01:30 | 1000.0 | foo |\ \n+---------------------+--------+------+", ); @@ -720,12 +708,104 @@ mod test { regular.start_timestamp = Some(500); let mut ordinary_nan = native_histogram(f64::NAN); ordinary_nan.start_timestamp = Some(0); + let mut unknown_start = native_histogram(7.0); + unknown_start.start_timestamp = None; let histograms = build_histogram_array(&[ Some(regular), Some(native_histogram(f64::from_bits(PROMETHEUS_STALE_NAN_BITS))), Some(ordinary_nan), + Some(unknown_start), None, ]); + for (unit, ticks_per_ms) in [ + (TimeUnit::Millisecond, 1_i64), + (TimeUnit::Microsecond, 1_000), + (TimeUnit::Nanosecond, 1_000_000), + ] { + let timestamp_array = |values: Vec| -> Arc { + match unit { + TimeUnit::Millisecond => Arc::new(TimestampMillisecondArray::from(values)), + TimeUnit::Microsecond => Arc::new(TimestampMicrosecondArray::from(values)), + TimeUnit::Nanosecond => Arc::new(TimestampNanosecondArray::from(values)), + TimeUnit::Second => unreachable!(), + } + }; + for offset in [-1_i64, 1] { + let timestamps = timestamp_array( + [1_000, 2_000, 3_000, 4_000, 5_000] + .into_iter() + .map(|timestamp| timestamp * ticks_per_ms) + .collect(), + ); + let schema = Arc::new(Schema::new(vec![ + Field::new(TIME_INDEX_COLUMN, timestamps.data_type().clone(), false), + Field::new("value", histograms.data_type().clone(), true), + ])); + let batch = + RecordBatch::try_new(schema.clone(), vec![timestamps, histograms.clone()]) + .unwrap(); + let input = Arc::new(DataSourceExec::new(Arc::new( + MemorySourceConfig::try_new(&[vec![batch]], schema, None).unwrap(), + ))); + let exec = Arc::new(SeriesNormalizeExec { + offset, + time_index_column_name: TIME_INDEX_COLUMN.to_string(), + filter_stale_markers: true, + tag_columns: Vec::new(), + input, + metric: ExecutionPlanMetricsSet::new(), + }); + let context = SessionContext::default(); + let batches = datafusion::physical_plan::collect(exec, context.task_ctx()) + .await + .unwrap(); + assert_eq!( + batches.iter().map(RecordBatch::num_rows).sum::(), + 4, + "unit={unit:?}, offset={offset}" + ); + let batch = batches.iter().find(|batch| batch.num_rows() == 4).unwrap(); + let expected_timestamps = timestamp_array( + [1_000, 3_000, 4_000, 5_000] + .into_iter() + .map(|timestamp| timestamp * ticks_per_ms) + .collect(), + ); + assert_eq!( + batch.column(0).to_data(), + expected_timestamps.to_data(), + "unit={unit:?}, offset={offset}" + ); + let values = batch + .column(1) + .as_any() + .downcast_ref::() + .unwrap(); + let regular = read_histogram(values, 0).unwrap().unwrap(); + assert_eq!( + (regular.sum, regular.start_timestamp), + (42.0, Some(500 + offset)), + "unit={unit:?}, offset={offset}" + ); + let ordinary_nan = read_histogram(values, 1).unwrap().unwrap(); + assert!(ordinary_nan.sum.is_nan()); + assert_eq!(ordinary_nan.start_timestamp, Some(0)); + let unknown_start = read_histogram(values, 2).unwrap().unwrap(); + assert_eq!( + (unknown_start.sum, unknown_start.start_timestamp), + (7.0, None) + ); + assert!(read_histogram(values, 3).unwrap().is_none()); + } + } + + let mut known_start = native_histogram(42.0); + known_start.start_timestamp = Some(500); + let mut sentinel_start = native_histogram(8.0); + sentinel_start.start_timestamp = Some(0); + let unknown_start = native_histogram(7.0); + let histograms = + build_histogram_array(&[Some(known_start), Some(sentinel_start), Some(unknown_start)]); let schema = Arc::new(Schema::new(vec![ Field::new( TIME_INDEX_COLUMN, @@ -737,47 +817,84 @@ mod test { let batch = RecordBatch::try_new( schema.clone(), vec![ - Arc::new(TimestampMillisecondArray::from(vec![ - 1_000, 2_000, 3_000, 4_000, - ])), + Arc::new(TimestampMillisecondArray::from(vec![1_000; 3])), histograms, ], ) .unwrap(); + let logical_input = LogicalPlan::EmptyRelation(EmptyRelation { + produce_one_row: false, + schema: schema.clone().to_dfschema_ref().unwrap(), + }); + let normalized = + SeriesNormalize::new(1_000, TIME_INDEX_COLUMN, false, Vec::new(), logical_input); + let range = RangeManipulate::new( + 2_000, + 2_000, + 1, + 1, + TIME_INDEX_COLUMN.to_string(), + vec!["value".to_string()], + LogicalPlan::Extension(datafusion::logical_expr::Extension { + node: Arc::new(normalized), + }), + ) + .unwrap(); let input = Arc::new(DataSourceExec::new(Arc::new( MemorySourceConfig::try_new(&[vec![batch]], schema, None).unwrap(), ))); - let exec = Arc::new(SeriesNormalizeExec { + let normalized_input = Arc::new(SeriesNormalizeExec { offset: 1_000, time_index_column_name: TIME_INDEX_COLUMN.to_string(), - filter_stale_markers: true, + filter_stale_markers: false, tag_columns: Vec::new(), input, metric: ExecutionPlanMetricsSet::new(), }); - - let context = SessionContext::default(); - let batches = datafusion::physical_plan::collect(exec, context.task_ctx()) - .await - .unwrap(); - let batch = batches.iter().find(|batch| batch.num_rows() == 3).unwrap(); - let values = batch - .column(1) + let output = datafusion::physical_plan::collect( + range.to_execution_plan(normalized_input), + SessionContext::default().task_ctx(), + ) + .await + .unwrap(); + let values = RangeArray::try_new( + output[0] + .column(1) + .as_any() + .downcast_ref::>() + .unwrap() + .clone(), + ) + .unwrap(); + let values = values.get(0).unwrap(); + let values = values .as_any() .downcast_ref::() .unwrap(); - - let timestamps = batch - .column(0) - .as_any() - .downcast_ref::() - .unwrap(); - assert_eq!(timestamps.values(), &[2_000, 4_000, 5_000]); - let regular = read_histogram(values, 0).unwrap().unwrap(); - assert_eq!((regular.sum, regular.start_timestamp), (42.0, Some(1_500))); - let ordinary_nan = read_histogram(values, 1).unwrap().unwrap(); - assert!(ordinary_nan.sum.is_nan()); - assert_eq!(ordinary_nan.start_timestamp, Some(0)); - assert!(read_histogram(values, 2).unwrap().is_none()); + assert_eq!( + read_histogram(values, 0).unwrap().unwrap().start_timestamp, + Some(1_500) + ); + assert_eq!( + read_histogram(values, 1).unwrap().unwrap().start_timestamp, + Some(0) + ); + assert_eq!( + read_histogram(values, 2).unwrap().unwrap().start_timestamp, + None + ); + let timestamps = RangeArray::try_new( + output[0] + .column(2) + .as_any() + .downcast_ref::>() + .unwrap() + .clone(), + ) + .unwrap(); + assert_eq!( + timestamps.get(0).unwrap().to_data(), + TimestampMillisecondArray::from(vec![2_000; 3]).to_data() + ); } } diff --git a/src/promql/src/extension_plan/range_manipulate.rs b/src/promql/src/extension_plan/range_manipulate.rs index 44c5f49094..99d39e6584 100644 --- a/src/promql/src/extension_plan/range_manipulate.rs +++ b/src/promql/src/extension_plan/range_manipulate.rs @@ -21,7 +21,7 @@ use std::task::{Context, Poll}; use common_telemetry::{debug, warn}; use datafusion::arrow::array::{Array, ArrayRef, Int64Array, TimestampMillisecondArray}; use datafusion::arrow::compute; -use datafusion::arrow::datatypes::{Field, SchemaRef}; +use datafusion::arrow::datatypes::{DataType, Field, SchemaRef, TimeUnit}; use datafusion::arrow::error::ArrowError; use datafusion::arrow::record_batch::RecordBatch; use datafusion::common::stats::Precision; @@ -46,7 +46,8 @@ use snafu::ResultExt; use crate::error::{DeserializeSnafu, Result}; use crate::extension_plan::{ - METRIC_NUM_SERIES, Millisecond, resolve_column_name, serialize_column_index, + METRIC_NUM_SERIES, Millisecond, local_offset, nanoseconds_per_native_tick, + native_timestamp_values, resolve_column_name, serialize_column_index, timestamp_unit, }; use crate::metrics::PROMQL_SERIES_COUNT; use crate::range_array::RangeArray; @@ -142,9 +143,21 @@ impl RangeManipulate { )); }; let ts_col_field = &columns[ts_col_index]; + let output_time_field = Arc::new( + ts_col_field + .as_ref() + .clone() + .with_data_type(DataType::Timestamp(TimeUnit::Millisecond, None)), + ); + new_columns[ts_col_index] = ( + input_schema.qualified_field(ts_col_index).0.cloned(), + output_time_field.clone(), + ); let timestamp_range_field = Field::new( Self::build_timestamp_range_name(time_index), - RangeArray::convert_field(ts_col_field).data_type().clone(), + RangeArray::convert_field(output_time_field.as_ref()) + .data_type() + .clone(), ts_col_field.is_nullable(), ); new_columns.push((None, Arc::new(timestamp_range_field))); @@ -177,6 +190,7 @@ impl RangeManipulate { properties.boundedness, )); Arc::new(RangeManipulateExec { + offset: local_offset(&self.input, &self.time_index), start: self.start, end: self.end, interval: self.interval, @@ -410,6 +424,7 @@ impl UserDefinedLogicalNodeCore for RangeManipulate { #[derive(Debug)] pub struct RangeManipulateExec { + offset: Millisecond, start: Millisecond, end: Millisecond, interval: Millisecond, @@ -470,6 +485,7 @@ impl ExecutionPlan for RangeManipulateExec { properties.boundedness, )); Ok(Arc::new(Self { + offset: self.offset, start: self.start, end: self.end, interval: self.interval, @@ -515,14 +531,17 @@ impl ExecutionPlan for RangeManipulateExec { .0 }) .collect(); + let time_unit = timestamp_unit(schema.field(time_index).data_type())?; let aligned_ts_array = RangeManipulateStream::build_aligned_ts_array(self.start, self.end, self.interval); Ok(Box::pin(RangeManipulateStream { + offset: self.offset, start: self.start, end: self.end, interval: self.interval, range: self.range, time_index, + time_unit, field_columns, aligned_ts_array, output_schema: self.output_schema.clone(), @@ -579,11 +598,13 @@ impl DisplayAs for RangeManipulateExec { } pub struct RangeManipulateStream { + offset: Millisecond, start: Millisecond, end: Millisecond, interval: Millisecond, range: Millisecond, time_index: usize, + time_unit: TimeUnit, field_columns: Vec, aligned_ts_array: ArrayRef, @@ -655,11 +676,31 @@ impl RangeManipulateStream { new_columns[*index] = new_column; } - // push timestamp range column - let ts_range_column = - RangeArray::from_ranges(input.column(self.time_index).clone(), ranges.clone()) - .map_err(|e| ArrowError::InvalidArgumentError(e.to_string()))? - .into_dict(); + // The timestamp range payload is always millisecond ABI. Shift in wide + // native precision before truncating toward zero, preserving null validity. + let scale = nanoseconds_per_native_tick(self.time_unit); + let timestamps = native_timestamp_values(input.column(self.time_index).as_ref())?; + let timestamp_values = timestamps + .iter() + .enumerate() + .map(|(index, timestamp)| { + if !input.column(self.time_index).is_valid(index) { + return Ok(None); + } + let shifted_ns = (*timestamp as i128) * scale + (self.offset as i128) * 1_000_000; + i64::try_from(shifted_ns / 1_000_000) + .map(Some) + .map_err(|_| { + ArrowError::ComputeError( + "RangeManipulate timestamp payload overflow".into(), + ) + }) + }) + .collect::, _>>()?; + let timestamp_values = TimestampMillisecondArray::from(timestamp_values); + let ts_range_column = RangeArray::from_ranges(Arc::new(timestamp_values), ranges.clone()) + .map_err(|e| ArrowError::InvalidArgumentError(e.to_string()))? + .into_dict(); new_columns.push(Arc::new(ts_range_column)); // truncate other columns @@ -694,52 +735,59 @@ impl RangeManipulateStream { &self, input: &RecordBatch, ) -> DataFusionResult<(Vec<(u32, u32)>, (i64, i64))> { - let ts_column = input - .column(self.time_index) - .as_any() - .downcast_ref::() - .ok_or_else(|| { - DataFusionError::Execution( - "Time index Column downcast to TimestampMillisecondArray failed".into(), - ) - })?; - - let len = ts_column.len(); + let ts_column = input.column(self.time_index); + let scale = nanoseconds_per_native_tick(self.time_unit); + let timestamps = native_timestamp_values(ts_column.as_ref())?; + let timestamp = + |index| (timestamps[index] as i128) * scale + (self.offset as i128) * 1_000_000; + let len = timestamps.len(); if len == 0 { return Ok((vec![], (self.start, self.end))); } - // shorten the range to calculate - let first_ts = ts_column.value(0); - // Preserve the query's alignment pattern when optimizing start time - let remainder = (first_ts - self.start).rem_euclid(self.interval); - let first_ts_aligned = if remainder == 0 { - first_ts - } else { - first_ts + (self.interval - remainder) - }; - let last_ts = ts_column.value(ts_column.len() - 1); - let last_ts_with_range = last_ts + self.range; - let remainder = (last_ts_with_range - self.start).rem_euclid(self.interval); + // Shorten the range using wide arithmetic so timestamps near the native + // type limits retain every query-aligned evaluation point. + let query_start = self.start as i128; + let query_end = self.end as i128; + let interval = self.interval as i128; + let first_ts = timestamp(0).div_euclid(1_000_000); + // Preserve the query's alignment pattern when optimizing start time. + let remainder = (first_ts - query_start).rem_euclid(interval); + let first_ts_aligned = first_ts + (interval - remainder).rem_euclid(interval); + let last_ts_with_range = + (timestamp(len - 1) + (self.range as i128) * 1_000_000).div_euclid(1_000_000); + let remainder = (last_ts_with_range - query_start).rem_euclid(interval); let last_ts_aligned = last_ts_with_range - remainder; - let start = self.start.max(first_ts_aligned); - let end = self.end.min(last_ts_aligned); + let start = query_start.max(first_ts_aligned); + let end = query_end.min(last_ts_aligned); if start > end { - return Ok((vec![], (start, end))); + let bounds = if start >= i64::MIN as i128 + && start <= i64::MAX as i128 + && end >= i64::MIN as i128 + && end <= i64::MAX as i128 + { + (start as i64, end as i64) + } else { + (self.start, self.end) + }; + return Ok((vec![], bounds)); } - let mut ranges = Vec::with_capacity(((self.end - self.start) / self.interval + 1) as usize); + // The intersection is within the declared i64 query bounds. + let start = start as i64; + let end = end as i64; + let mut ranges = Vec::new(); // calculate for every aligned timestamp (`curr_ts`), assume the ts column is ordered. let mut left = 0usize; let mut right = 0usize; for curr_ts in (start..=end).step_by(self.interval as _) { - let start_ts = curr_ts - self.range; + let start_ts = (curr_ts as i128) * 1_000_000 - (self.range as i128) * 1_000_000; - while left < len && ts_column.value(left) <= start_ts { + while left < len && timestamp(left) <= start_ts { left += 1; } right = right.max(left); - while right < len && ts_column.value(right) <= curr_ts { + while right < len && timestamp(right) <= (curr_ts as i128) * 1_000_000 { right += 1; } @@ -756,14 +804,20 @@ impl RangeManipulateStream { #[cfg(test)] mod test { - use datafusion::arrow::array::{ArrayRef, DictionaryArray, Float64Array, StringArray}; + use datafusion::arrow::array::{ + ArrayRef, DictionaryArray, Float64Array, StringArray, TimestampMicrosecondArray, + TimestampNanosecondArray, TimestampSecondArray, + }; + use datafusion::arrow::buffer::NullBuffer; use datafusion::arrow::datatypes::{ ArrowPrimitiveType, DataType, Field, Int64Type, Schema, TimestampMillisecondType, }; use datafusion::common::ToDFSchema; use datafusion::datasource::memory::MemorySourceConfig; use datafusion::datasource::source::DataSourceExec; - use datafusion::logical_expr::{EmptyRelation, LogicalPlan}; + use datafusion::logical_expr::{ + EmptyRelation, Extension, LogicalPlan, UserDefinedLogicalNodeCore, + }; use datafusion::physical_expr::Partitioning; use datafusion::physical_plan::execution_plan::{Boundedness, EmissionType}; use datafusion::physical_plan::memory::MemoryStream; @@ -845,6 +899,7 @@ mod test { Boundedness::Bounded, )); let normalize_exec = Arc::new(RangeManipulateExec { + offset: 0, start, end, interval, @@ -888,6 +943,365 @@ mod test { assert_eq!(result_literal, expected); } + #[tokio::test] + async fn native_timestamps_preserve_range_membership_and_ms_payload() { + for (unit, ticks_per_ms) in [ + (TimeUnit::Microsecond, 1_000_i64), + (TimeUnit::Nanosecond, 1_000_000_i64), + ] { + let lower = 1_000 * ticks_per_ms; + let upper = 1_001 * ticks_per_ms; + // Exclude the lower boundary and future sample; retain both native + // samples in the same millisecond bucket and the exact upper sample. + let timestamps = vec![lower, lower + 1, lower + 2, upper, upper + 1]; + let time: ArrayRef = match unit { + TimeUnit::Microsecond => Arc::new(TimestampMicrosecondArray::from(timestamps)), + TimeUnit::Nanosecond => Arc::new(TimestampNanosecondArray::from(timestamps)), + _ => unreachable!(), + }; + let schema = Arc::new(Schema::new(vec![ + Field::new(TIME_INDEX_COLUMN, DataType::Timestamp(unit, None), false), + Field::new("value", DataType::Float64, true), + ])); + let batch = RecordBatch::try_new( + schema.clone(), + vec![ + time, + Arc::new(Float64Array::from(vec![10.0, 20.0, 30.0, 40.0, 50.0])), + ], + ) + .unwrap(); + let logical_input = LogicalPlan::EmptyRelation(EmptyRelation { + produce_one_row: false, + schema: schema.clone().to_dfschema_ref().unwrap(), + }); + let plan = RangeManipulate::new( + 1_001, + 1_001, + 1, + 1, + TIME_INDEX_COLUMN.to_string(), + vec!["value".to_string()], + logical_input.clone(), + ) + .unwrap(); + let output_time = Field::new( + TIME_INDEX_COLUMN, + DataType::Timestamp(TimeUnit::Millisecond, None), + false, + ); + let output_schema = Arc::new(Schema::new(vec![ + output_time.clone(), + RangeArray::convert_field(&Field::new("value", DataType::Float64, true)), + Field::new( + RangeManipulate::build_timestamp_range_name(TIME_INDEX_COLUMN), + RangeArray::convert_field(&output_time).data_type().clone(), + false, + ), + ])); + assert_eq!(plan.schema().as_arrow(), output_schema.as_ref()); + + let rebuilt = RangeManipulate::deserialize(&plan.serialize()) + .unwrap() + .with_exprs_and_inputs(vec![], vec![logical_input]) + .unwrap(); + assert_eq!(rebuilt.schema(), plan.schema()); + assert_eq!(rebuilt.input.schema().as_arrow(), schema.as_ref()); + + let input = Arc::new(DataSourceExec::new(Arc::new( + MemorySourceConfig::try_new(&[vec![batch]], schema.clone(), None).unwrap(), + ))); + let exec = rebuilt.to_execution_plan(input); + assert_eq!(exec.schema(), output_schema); + assert_eq!(exec.children()[0].schema(), schema); + + let batches = + datafusion::physical_plan::collect(exec, SessionContext::default().task_ctx()) + .await + .unwrap(); + assert_eq!(batches.len(), 1, "{unit:?}"); + let output = &batches[0]; + assert_eq!(output.schema(), output_schema); + assert_eq!(output.num_rows(), 1); + assert_eq!( + output + .column(0) + .as_any() + .downcast_ref::() + .unwrap() + .values() + .as_ref(), + &[1_001] + ); + + // RangeArray packs offset/length into dictionary keys; Arrow dictionary + // equality treats those packed keys as indices and cannot compare them. + let values = RangeArray::try_new( + output + .column(1) + .as_any() + .downcast_ref::>() + .unwrap() + .clone(), + ) + .unwrap(); + assert_eq!(values.get_offset_length(0), Some((1, 3))); + assert_eq!( + values.get(0).unwrap().to_data(), + Float64Array::from(vec![20.0, 30.0, 40.0]).to_data() + ); + let timestamps = RangeArray::try_new( + output + .column(2) + .as_any() + .downcast_ref::>() + .unwrap() + .clone(), + ) + .unwrap(); + assert_eq!(timestamps.get_offset_length(0), Some((1, 3))); + assert_eq!( + timestamps.get(0).unwrap().to_data(), + TimestampMillisecondArray::from(vec![1_000, 1_000, 1_001]).to_data() + ); + } + } + + #[tokio::test] + async fn logical_normalize_offset_survives_rebuild_and_executes() { + for (name, time_unit, raw, offset, start, range, expected_payload) in [ + ( + "millisecond offset", + TimeUnit::Millisecond, + 0, + 1_000, + 1_000, + 1_000, + 1_000, + ), + ( + "negative native lower limit with positive window", + TimeUnit::Nanosecond, + -9_223_112_837_000_000_000, + -259_200_000, + -9_223_372_037_000, + 300_000, + -9_223_372_037_000, + ), + ( + "second timestamp with negative fractional offset", + TimeUnit::Second, + 1, + -500, + 1_000, + 1_000, + 500, + ), + ] { + let schema = Arc::new(Schema::new(vec![ + Field::new( + TIME_INDEX_COLUMN, + DataType::Timestamp(time_unit, None), + false, + ), + Field::new("value", DataType::Float64, true), + ])); + let input = LogicalPlan::EmptyRelation(EmptyRelation { + produce_one_row: false, + schema: schema.clone().to_dfschema_ref().unwrap(), + }); + let normalize = crate::extension_plan::SeriesNormalize::new( + offset, + TIME_INDEX_COLUMN, + false, + Vec::new(), + input.clone(), + ); + let normalize = + crate::extension_plan::SeriesNormalize::deserialize(&normalize.serialize()) + .unwrap() + .with_exprs_and_inputs(vec![], vec![input]) + .unwrap(); + let normalized = LogicalPlan::Extension(Extension { + node: Arc::new(normalize), + }); + let plan = RangeManipulate::new( + start, + start, + 1, + range, + TIME_INDEX_COLUMN.to_string(), + vec!["value".to_string()], + normalized.clone(), + ) + .unwrap(); + let rebuilt = RangeManipulate::deserialize(&plan.serialize()) + .unwrap() + .with_exprs_and_inputs(vec![], vec![normalized]) + .unwrap(); + let timestamp: ArrayRef = match time_unit { + TimeUnit::Millisecond => Arc::new(TimestampMillisecondArray::from(vec![raw])), + TimeUnit::Nanosecond => Arc::new(TimestampNanosecondArray::from(vec![raw])), + TimeUnit::Second => Arc::new(TimestampSecondArray::from(vec![raw])), + _ => unreachable!(), + }; + let batch = RecordBatch::try_new( + schema.clone(), + vec![timestamp, Arc::new(Float64Array::from(vec![7.0]))], + ) + .unwrap(); + let exec_input = Arc::new(DataSourceExec::new(Arc::new( + MemorySourceConfig::try_new(&[vec![batch]], schema, None).unwrap(), + ))); + let output = datafusion::physical_plan::collect( + rebuilt.to_execution_plan(exec_input), + SessionContext::default().task_ctx(), + ) + .await + .unwrap(); + assert_eq!(output.len(), 1, "{name}"); + let output = &output[0]; + assert_eq!(output.num_rows(), 1, "{name}"); + assert_eq!( + output + .column(0) + .as_any() + .downcast_ref::() + .unwrap() + .value(0), + start, + "{name}" + ); + let values = RangeArray::try_new( + output + .column(1) + .as_any() + .downcast_ref::>() + .unwrap() + .clone(), + ) + .unwrap(); + assert_eq!(values.get_offset_length(0), Some((0, 1)), "{name}"); + assert_eq!( + values.get(0).unwrap().to_data(), + Float64Array::from(vec![7.0]).to_data(), + "{name}" + ); + let timestamps = RangeArray::try_new( + output + .column(2) + .as_any() + .downcast_ref::>() + .unwrap() + .clone(), + ) + .unwrap(); + assert_eq!(timestamps.get_offset_length(0), Some((0, 1)), "{name}"); + assert_eq!( + timestamps.get(0).unwrap().to_data(), + TimestampMillisecondArray::from(vec![expected_payload]).to_data(), + "{name}" + ); + } + } + + #[tokio::test] + async fn range_payload_preserves_null_timestamp_and_rejects_offset_overflow() { + let schema = Arc::new(Schema::new(vec![ + Field::new(TIME_INDEX_COLUMN, TimestampMillisecondType::DATA_TYPE, true), + Field::new("value", DataType::Float64, true), + ])); + let null_timestamp = + TimestampMillisecondArray::new(vec![1_000].into(), Some(NullBuffer::from(vec![false]))); + let batch = RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(null_timestamp), + Arc::new(Float64Array::from(vec![7.0])), + ], + ) + .unwrap(); + let input = Arc::new(DataSourceExec::new(Arc::new( + MemorySourceConfig::try_new(&[vec![batch]], schema.clone(), None).unwrap(), + ))); + let plan = RangeManipulate::new( + 1_000, + 1_000, + 1, + 1, + TIME_INDEX_COLUMN.to_string(), + vec!["value".to_string()], + LogicalPlan::EmptyRelation(EmptyRelation { + produce_one_row: false, + schema: schema.clone().to_dfschema_ref().unwrap(), + }), + ) + .unwrap(); + let output = datafusion::physical_plan::collect( + plan.to_execution_plan(input), + SessionContext::default().task_ctx(), + ) + .await + .unwrap(); + let timestamps = RangeArray::try_new( + output[0] + .column(2) + .as_any() + .downcast_ref::>() + .unwrap() + .clone(), + ) + .unwrap(); + let payload = timestamps.get(0).unwrap(); + let payload = payload + .as_any() + .downcast_ref::() + .unwrap(); + assert_eq!(payload.len(), 1); + assert!(!payload.is_valid(0)); + + let batch = RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(TimestampMillisecondArray::from(vec![0, i64::MAX])), + Arc::new(Float64Array::from(vec![7.0, 8.0])), + ], + ) + .unwrap(); + let input = Arc::new(DataSourceExec::new(Arc::new( + MemorySourceConfig::try_new(&[vec![batch]], schema.clone(), None).unwrap(), + ))); + let normalized = crate::extension_plan::SeriesNormalize::new( + 1, + TIME_INDEX_COLUMN, + false, + Vec::new(), + LogicalPlan::EmptyRelation(EmptyRelation { + produce_one_row: false, + schema: schema.to_dfschema_ref().unwrap(), + }), + ); + let plan = RangeManipulate::new( + 1, + 1, + 1, + 1, + TIME_INDEX_COLUMN.to_string(), + vec!["value".to_string()], + LogicalPlan::Extension(Extension { + node: Arc::new(normalized), + }), + ) + .unwrap(); + let error = datafusion::physical_plan::collect( + plan.to_execution_plan(input), + SessionContext::default().task_ctx(), + ) + .await + .unwrap_err(); + assert!(error.to_string().contains("timestamp payload overflow")); + } + #[tokio::test] async fn pruning_should_keep_time_and_value_columns_for_exec() { let schema = Arc::new(Schema::new(vec![ @@ -1042,11 +1456,13 @@ mod test { let empty_stream = MemoryStream::try_new(vec![], schema.clone(), None).unwrap(); let stream = RangeManipulateStream { + offset: 0, start: 1758093274000, // ends in 4000 end: 1758093334000, // ends in 4000 interval: 30000, // 30s step range: 60000, // 60s lookback time_index: 0, + time_unit: TimeUnit::Millisecond, field_columns: vec![], aligned_ts_array: Arc::new(TimestampMillisecondArray::from(vec![0i64; 0])), output_schema: schema.clone(), @@ -1106,11 +1522,13 @@ mod test { )])); let empty_stream = MemoryStream::try_new(vec![], schema.clone(), None).unwrap(); let stream = RangeManipulateStream { + offset: 0, start: query_start, end: query_end, interval, range, time_index: 0, + time_unit: TimeUnit::Millisecond, field_columns: vec![], aligned_ts_array: Arc::new(TimestampMillisecondArray::from(vec![0i64; 0])), output_schema: schema.clone(), @@ -1227,6 +1645,15 @@ mod test { } } + #[test] + fn calculate_range_keeps_extreme_range_tail() { + let (ranges, bounds) = + calculate_range_for_test(i64::MAX - 1, i64::MAX, 1, i64::MAX, &[i64::MAX]); + + assert_eq!(bounds, (i64::MAX, i64::MAX)); + assert_eq!(ranges, vec![(0, 1)]); + } + #[test] fn calculate_range_matches_bruteforce_oracle_for_deterministic_cases() { let cases = vec![ diff --git a/src/query/src/promql/planner.rs b/src/query/src/promql/planner.rs index 7f1d5d35fb..34d0e20840 100644 --- a/src/query/src/promql/planner.rs +++ b/src/query/src/promql/planner.rs @@ -1943,6 +1943,11 @@ impl PromPlanner { if let Some(empty_plan) = self.setup_context().await? { return Ok(empty_plan); } + let offset_ms = match offset { + Some(Offset::Pos(duration)) => duration.as_millis() as Millisecond, + Some(Offset::Neg(duration)) => -(duration.as_millis() as Millisecond), + None => 0, + }; let normalize = self .selector_to_series_normalize_plan(offset, matchers, false) .await?; @@ -1974,8 +1979,48 @@ impl PromPlanner { DfExpr::Column(Column::new(qualifier.cloned(), field.name().clone())) }) .collect::>(); - project_exprs - .push(build_special_time_expr(&time_index_column).alias(×tamp_value_column)); + // `timestamp()` preserves the shifted selector timeline even though + // SeriesNormalize now retains raw native timestamp storage. Decimal + // arithmetic shifts before truncating to milliseconds. + let unit_factor = match col(&time_index_column) + .get_type(normalize.schema()) + .context(DataFusionPlanningSnafu)? + { + ArrowDataType::Timestamp(ArrowTimeUnit::Second, _) => (1_000_i128, 4, 0), + ArrowDataType::Timestamp(ArrowTimeUnit::Millisecond, _) => (1, 1, 0), + ArrowDataType::Timestamp(ArrowTimeUnit::Microsecond, _) => (1, 4, 3), + ArrowDataType::Timestamp(ArrowTimeUnit::Nanosecond, _) => (1, 7, 6), + _ => unreachable!("time index is a timestamp"), + }; + let sample_time = col(&time_index_column) + .cast_to(&ArrowDataType::Int64, normalize.schema()) + .context(DataFusionPlanningSnafu)? + .cast_to(&ArrowDataType::Decimal128(19, 0), normalize.schema()) + .context(DataFusionPlanningSnafu)?; + let sample_time = DfExpr::BinaryExpr(BinaryExpr { + left: Box::new(sample_time), + op: Operator::Multiply, + right: Box::new(lit(ScalarValue::Decimal128( + Some(unit_factor.0), + unit_factor.1, + unit_factor.2, + ))), + }); + let sample_time = DfExpr::BinaryExpr(BinaryExpr { + left: Box::new(sample_time), + op: Operator::Plus, + right: Box::new(lit(ScalarValue::Decimal128(Some(offset_ms as i128), 19, 0))), + }) + .cast_to(&ArrowDataType::Int64, normalize.schema()) + .context(DataFusionPlanningSnafu)? + .cast_to(&ArrowDataType::Float64, normalize.schema()) + .context(DataFusionPlanningSnafu)?; + let sample_time = DfExpr::BinaryExpr(BinaryExpr { + left: Box::new(sample_time), + op: Operator::Divide, + right: Box::new(lit(1000.0)), + }); + project_exprs.push(sample_time.alias(×tamp_value_column)); let normalize = LogicalPlanBuilder::from(normalize) .project(project_exprs) .context(DataFusionPlanningSnafu)? @@ -2339,14 +2384,18 @@ impl PromPlanner { None => 0, }; let mut scan_filters = Self::matchers_to_expr(label_matchers.clone(), table_schema)?; - if let Some(time_index_filter) = self.build_time_index_filter(offset_duration)? { + if let Some(time_index_filter) = + self.build_time_index_filter(offset_duration, table_schema)? + { scan_filters.push(time_index_filter); } - table_scan = LogicalPlanBuilder::from(table_scan) - .filter(conjunction(scan_filters).unwrap()) // Safety: `scan_filters` is not empty. - .context(DataFusionPlanningSnafu)? - .build() - .context(DataFusionPlanningSnafu)?; + if let Some(filter) = conjunction(scan_filters) { + table_scan = LogicalPlanBuilder::from(table_scan) + .filter(filter) + .context(DataFusionPlanningSnafu)? + .build() + .context(DataFusionPlanningSnafu)?; + } // make a projection plan if there is any `__field__` matcher if let Some(field_matchers) = &self.ctx.field_column_matcher { @@ -2718,72 +2767,83 @@ impl PromPlanner { Ok(table_ref) } - fn build_time_index_filter(&self, offset_duration: i64) -> Result> { + fn build_time_index_filter( + &self, + offset_duration: i64, + schema: &DFSchemaRef, + ) -> Result> { let start = self.ctx.start; let end = self.ctx.end; if end < start { return InvalidTimeRangeSnafu { start, end }.fail(); } - let lookback_delta = self.ctx.lookback_delta; - let range = self.ctx.range.unwrap_or_default(); - let interval = self.ctx.interval; let time_index_expr = self.create_time_index_column_expr()?; - let num_points = (end - start) / interval; - - // Prometheus semantics: - // - Instant selector lookback: (eval_ts - lookback_delta, eval_ts] - // - Range selector: (eval_ts - range, eval_ts] - // - // So samples positioned exactly at the lower boundary must be excluded. We align the scan - // lower bound with Prometheus by shifting it forward by 1ms (millisecond granularity), - // while still using a `>=` filter. - let selector_window = if range == 0 { lookback_delta } else { range }; - let lower_exclusive_adjustment = if selector_window > 0 { 1 } else { 0 }; - - // Scan a continuous time range - if (end - start) / interval > MAX_SCATTER_POINTS || interval <= INTERVAL_1H { - let single_time_range = time_index_expr - .clone() - .gt_eq(DfExpr::Literal( - ScalarValue::TimestampMillisecond( - Some( - self.ctx.start - offset_duration - selector_window - + lower_exclusive_adjustment, - ), - None, - ), - None, - )) - .and(time_index_expr.lt_eq(DfExpr::Literal( - ScalarValue::TimestampMillisecond(Some(self.ctx.end - offset_duration), None), - None, - ))); - return Ok(Some(single_time_range)); - } - - // Otherwise scan scatter ranges separately - let mut filters = Vec::with_capacity(num_points as usize + 1); - for timestamp in (start..=end).step_by(interval as usize) { - filters.push( + let time_index_name = self.ctx.time_index_column.as_ref().unwrap(); + let unit = schema + .index_of_column_by_name(None, time_index_name) + .and_then(|index| match schema.field(index).data_type() { + ArrowDataType::Timestamp(unit, _) => Some(*unit), + _ => None, + }) + .unwrap_or(ArrowTimeUnit::Millisecond); + let scalar = |milliseconds: i64| -> Option { + let value = match unit { + ArrowTimeUnit::Second => milliseconds.div_euclid(1_000), + ArrowTimeUnit::Millisecond => milliseconds, + ArrowTimeUnit::Microsecond => milliseconds.checked_mul(1_000)?, + ArrowTimeUnit::Nanosecond => milliseconds.checked_mul(1_000_000)?, + }; + Some(match unit { + ArrowTimeUnit::Second => ScalarValue::TimestampSecond(Some(value), None), + ArrowTimeUnit::Millisecond => ScalarValue::TimestampMillisecond(Some(value), None), + ArrowTimeUnit::Microsecond => ScalarValue::TimestampMicrosecond(Some(value), None), + ArrowTimeUnit::Nanosecond => ScalarValue::TimestampNanosecond(Some(value), None), + }) + }; + let window = self.ctx.range.unwrap_or(self.ctx.lookback_delta); + let filter = |lower_ms: i64, upper_ms: i64| -> Option { + let lower = DfExpr::Literal(scalar(lower_ms)?, None); + let lower_filter = if window == 0 { + time_index_expr.clone().gt_eq(lower) + } else if unit == ArrowTimeUnit::Millisecond + && let Some(inclusive_lower) = lower_ms.checked_add(1) + { time_index_expr .clone() - .gt_eq(DfExpr::Literal( - ScalarValue::TimestampMillisecond( - Some( - timestamp - offset_duration - selector_window - + lower_exclusive_adjustment, - ), - None, - ), - None, - )) - .and(time_index_expr.clone().lt_eq(DfExpr::Literal( - ScalarValue::TimestampMillisecond(Some(timestamp - offset_duration), None), - None, - ))), + .gt_eq(DfExpr::Literal(scalar(inclusive_lower)?, None)) + } else { + time_index_expr.clone().gt(lower) + }; + Some( + lower_filter.and( + time_index_expr + .clone() + .lt_eq(DfExpr::Literal(scalar(upper_ms)?, None)), + ), ) + }; + let bounds = |timestamp: i64| { + timestamp + .checked_sub(offset_duration) + .and_then(|upper| upper.checked_sub(window).map(|lower| (lower, upper))) + }; + let num_points = (end as i128 - start as i128) / self.ctx.interval as i128; + if num_points > MAX_SCATTER_POINTS as i128 || self.ctx.interval <= INTERVAL_1H { + return Ok(bounds(start) + .zip(bounds(end)) + .and_then(|((lower, _), (_, upper))| filter(lower, upper))); + } + let mut filters = Vec::new(); + for timestamp in (start..=end).step_by(self.ctx.interval as usize) { + let Some((lower, upper)) = bounds(timestamp) else { + // An unrepresentable envelope must not discard samples. + return Ok(None); + }; + let Some(filter) = filter(lower, upper) else { + return Ok(None); + }; + filters.push(filter); } - Ok(filters.into_iter().reduce(DfExpr::or)) } @@ -2884,14 +2944,16 @@ impl PromPlanner { self.ctx.tag_columns.clone() }; - let is_time_index_ms = scan_table + let time_index_data_type = scan_table .schema() .timestamp_column() .with_context(|| TimeIndexNotFoundSnafu { table: maybe_phy_table_ref.to_quoted_string(), })? .data_type - == ConcreteDataType::timestamp_millisecond_datatype(); + .clone(); + let is_time_index_second = + time_index_data_type == ConcreteDataType::timestamp_second_datatype(); let scan_projection = if table_id_filter.is_some() { let mut required_columns = HashSet::new(); @@ -2944,8 +3006,8 @@ impl PromPlanner { .context(DataFusionPlanningSnafu)?; } - if !is_time_index_ms { - // cast to ms if time_index not in Millisecond precision + if is_time_index_second { + // Promote seconds so millisecond offsets remain exact; retain finer precision. let expr: Vec<_> = self .create_field_column_exprs()? .into_iter() @@ -2980,8 +3042,12 @@ impl PromPlanner { .context(DataFusionPlanningSnafu)? .build() .context(DataFusionPlanningSnafu)?; - } else if table_id_filter.is_some() { - // Drop the internal `__table_id` column after filtering. + } else if table_id_filter.is_some() + || time_index_data_type == ConcreteDataType::timestamp_microsecond_datatype() + || time_index_data_type == ConcreteDataType::timestamp_nanosecond_datatype() + { + // Drop the internal `__table_id` column after filtering and preserve PromQL's + // field/tag/timestamp column order for native microsecond/nanosecond timestamps. let project_exprs = self .create_field_column_exprs()? .into_iter() @@ -8832,7 +8898,7 @@ mod test { \n Projection: some_metric.timestamp, value AS value, some_metric.tag_0 [timestamp:Timestamp(ms), value:Float64, tag_0:Utf8]\ \n Projection: some_metric.timestamp, __promql_timestamp_value_ AS value, some_metric.tag_0 [timestamp:Timestamp(ms), value:Float64, tag_0:Utf8]\ \n PromInstantManipulate: range=[0..100000000], lookback=[1000], interval=[5000], time index=[timestamp] [tag_0:Utf8, timestamp:Timestamp(ms), field_0:Float64;N, __promql_timestamp_value_:Float64]\ - \n Projection: some_metric.tag_0, some_metric.timestamp, some_metric.field_0, CAST(CAST(some_metric.timestamp AS Int64) AS Float64) / Float64(1000) AS __promql_timestamp_value_ [tag_0:Utf8, timestamp:Timestamp(ms), field_0:Float64;N, __promql_timestamp_value_:Float64]\ + \n Projection: some_metric.tag_0, some_metric.timestamp, some_metric.field_0, CAST(CAST(CAST(CAST(some_metric.timestamp AS Int64) AS Decimal128(19, 0)) * Decimal128(Some(1),1,0) + Decimal128(Some(0),19,0) AS Int64) AS Float64) / Float64(1000) AS __promql_timestamp_value_ [tag_0:Utf8, timestamp:Timestamp(ms), field_0:Float64;N, __promql_timestamp_value_:Float64]\ \n PromSeriesDivide: tags=[\"tag_0\"] [tag_0:Utf8, timestamp:Timestamp(ms), field_0:Float64;N]\ \n Sort: some_metric.tag_0 ASC NULLS FIRST, some_metric.timestamp ASC NULLS FIRST [tag_0:Utf8, timestamp:Timestamp(ms), field_0:Float64;N]\ \n Filter: some_metric.tag_0 != Utf8(\"bar\") AND some_metric.timestamp >= TimestampMillisecond(-999, None) AND some_metric.timestamp <= TimestampMillisecond(100000000, None) [tag_0:Utf8, timestamp:Timestamp(ms), field_0:Float64;N]\ @@ -9051,7 +9117,17 @@ mod test { let manipulate = find_instant_manipulate(&plan).unwrap(); let exec = manipulate.to_execution_plan(Arc::new(DataSourceExec::new(Arc::new( - MemorySourceConfig::try_new(&[], Arc::new(ArrowSchema::empty()), None).unwrap(), + MemorySourceConfig::try_new( + &[], + Arc::new( + datafusion_expr::UserDefinedLogicalNodeCore::inputs(manipulate)[0] + .schema() + .as_arrow() + .clone(), + ), + None, + ) + .unwrap(), )))); assert!(format!("{exec:?}").contains("reuse_tsid_column: true")); } @@ -12135,6 +12211,57 @@ mod test { } } + #[tokio::test] + async fn native_scan_bounds_preserve_zero_lookback_and_overflow() { + let table_provider = build_test_table_provider( + &[(DEFAULT_SCHEMA_NAME.to_string(), "some_metric".to_string())], + 1, + 1, + ) + .await; + let mut planner = PromPlanner { + table_provider, + ctx: PromPlannerContext::from_eval_stmt(&build_eval_stmt("some_metric")), + promql_annotations: None, + }; + planner.ctx.time_index_column = Some("timestamp".to_string()); + planner.ctx.start = 1_000; + planner.ctx.lookback_delta = 0; + let schema = Arc::new( + DFSchema::try_from(ArrowSchema::new(vec![Field::new( + "timestamp", + ArrowDataType::Timestamp(ArrowTimeUnit::Nanosecond, None), + false, + )])) + .unwrap(), + ); + for (end, interval, windows) in [ + (1_000, 1_000, 1), + (2_000, 1_000, 1), + (7_201_000, 7_200_000, 2), + ] { + planner.ctx.end = end; + planner.ctx.interval = interval; + let filter = planner + .build_time_index_filter(0, &schema) + .unwrap() + .unwrap() + .to_string(); + assert_eq!(filter.matches(">=").count(), windows, "{filter}"); + assert!( + filter.contains("TimestampNanosecond(1000000000, None)"), + "{filter}" + ); + } + planner.ctx.end = i64::MAX; + assert!( + planner + .build_time_index_filter(0, &schema) + .unwrap() + .is_none() + ); + } + #[tokio::test] async fn test_non_ms_precision() { let catalog_list = MemoryCatalogManager::with_default_setup(); @@ -12205,12 +12332,7 @@ mod test { .unwrap(); assert_eq!( plan.display_indent_schema().to_string(), - "PromInstantManipulate: range=[0..100000000], lookback=[1000], interval=[5000], time index=[timestamp] [field:Float64;N, tag:Utf8, timestamp:Timestamp(ms)]\ - \n PromSeriesDivide: tags=[\"tag\"] [field:Float64;N, tag:Utf8, timestamp:Timestamp(ms)]\ - \n Sort: metrics.tag ASC NULLS FIRST, metrics.timestamp ASC NULLS FIRST [field:Float64;N, tag:Utf8, timestamp:Timestamp(ms)]\ - \n Filter: metrics.tag = Utf8(\"1\") AND metrics.timestamp >= TimestampMillisecond(-999, None) AND metrics.timestamp <= TimestampMillisecond(100000000, None) [field:Float64;N, tag:Utf8, timestamp:Timestamp(ms)]\ - \n Projection: metrics.field, metrics.tag, CAST(metrics.timestamp AS Timestamp(ms)) AS timestamp [field:Float64;N, tag:Utf8, timestamp:Timestamp(ms)]\ - \n TableScan: metrics [tag:Utf8, timestamp:Timestamp(ns), field:Float64;N]" + "PromInstantManipulate: range=[0..100000000], lookback=[1000], interval=[5000], time index=[timestamp] [field:Float64;N, tag:Utf8, timestamp:Timestamp(ms)]\n PromSeriesDivide: tags=[\"tag\"] [field:Float64;N, tag:Utf8, timestamp:Timestamp(ns)]\n Sort: metrics.tag ASC NULLS FIRST, metrics.timestamp ASC NULLS FIRST [field:Float64;N, tag:Utf8, timestamp:Timestamp(ns)]\n Filter: metrics.tag = Utf8(\"1\") AND metrics.timestamp > TimestampNanosecond(-1000000000, None) AND metrics.timestamp <= TimestampNanosecond(100000000000000, None) [field:Float64;N, tag:Utf8, timestamp:Timestamp(ns)]\n Projection: metrics.field, metrics.tag, metrics.timestamp [field:Float64;N, tag:Utf8, timestamp:Timestamp(ns)]\n TableScan: metrics [tag:Utf8, timestamp:Timestamp(ns), field:Float64;N]" ); let plan = PromPlanner::stmt_to_plan( DfTableSourceProvider::new( @@ -12235,15 +12357,7 @@ mod test { .unwrap(); assert_eq!( plan.display_indent_schema().to_string(), - "Filter: prom_avg_over_time(timestamp_range,field) IS NOT NULL [timestamp:Timestamp(ms), prom_avg_over_time(timestamp_range,field):Float64;N, tag:Utf8]\ - \n Projection: metrics.timestamp, prom_avg_over_time(timestamp_range, field) AS prom_avg_over_time(timestamp_range,field), metrics.tag [timestamp:Timestamp(ms), prom_avg_over_time(timestamp_range,field):Float64;N, tag:Utf8]\ - \n PromRangeManipulate: req range=[0..100000000], interval=[5000], eval range=[5000], time index=[timestamp], values=[\"field\"] [field:Dictionary(Int64, Float64);N, tag:Utf8, timestamp:Timestamp(ms), timestamp_range:Dictionary(Int64, Timestamp(ms))]\ - \n PromSeriesNormalize: offset=[0], time index=[timestamp], filter NaN: [true] [field:Float64;N, tag:Utf8, timestamp:Timestamp(ms)]\ - \n PromSeriesDivide: tags=[\"tag\"] [field:Float64;N, tag:Utf8, timestamp:Timestamp(ms)]\ - \n Sort: metrics.tag ASC NULLS FIRST, metrics.timestamp ASC NULLS FIRST [field:Float64;N, tag:Utf8, timestamp:Timestamp(ms)]\ - \n Filter: metrics.tag = Utf8(\"1\") AND metrics.timestamp >= TimestampMillisecond(-4999, None) AND metrics.timestamp <= TimestampMillisecond(100000000, None) [field:Float64;N, tag:Utf8, timestamp:Timestamp(ms)]\ - \n Projection: metrics.field, metrics.tag, CAST(metrics.timestamp AS Timestamp(ms)) AS timestamp [field:Float64;N, tag:Utf8, timestamp:Timestamp(ms)]\ - \n TableScan: metrics [tag:Utf8, timestamp:Timestamp(ns), field:Float64;N]" + "Filter: prom_avg_over_time(timestamp_range,field) IS NOT NULL [timestamp:Timestamp(ms), prom_avg_over_time(timestamp_range,field):Float64;N, tag:Utf8]\n Projection: metrics.timestamp, prom_avg_over_time(timestamp_range, field) AS prom_avg_over_time(timestamp_range,field), metrics.tag [timestamp:Timestamp(ms), prom_avg_over_time(timestamp_range,field):Float64;N, tag:Utf8]\n PromRangeManipulate: req range=[0..100000000], interval=[5000], eval range=[5000], time index=[timestamp], values=[\"field\"] [field:Dictionary(Int64, Float64);N, tag:Utf8, timestamp:Timestamp(ms), timestamp_range:Dictionary(Int64, Timestamp(ms))]\n PromSeriesNormalize: offset=[0], time index=[timestamp], filter NaN: [true] [field:Float64;N, tag:Utf8, timestamp:Timestamp(ns)]\n PromSeriesDivide: tags=[\"tag\"] [field:Float64;N, tag:Utf8, timestamp:Timestamp(ns)]\n Sort: metrics.tag ASC NULLS FIRST, metrics.timestamp ASC NULLS FIRST [field:Float64;N, tag:Utf8, timestamp:Timestamp(ns)]\n Filter: metrics.tag = Utf8(\"1\") AND metrics.timestamp > TimestampNanosecond(-5000000000, None) AND metrics.timestamp <= TimestampNanosecond(100000000000000, None) [field:Float64;N, tag:Utf8, timestamp:Timestamp(ns)]\n Projection: metrics.field, metrics.tag, metrics.timestamp [field:Float64;N, tag:Utf8, timestamp:Timestamp(ns)]\n TableScan: metrics [tag:Utf8, timestamp:Timestamp(ns), field:Float64;N]" ); } diff --git a/tests/cases/standalone/common/promql/native_time_selection.result b/tests/cases/standalone/common/promql/native_time_selection.result new file mode 100644 index 0000000000..1715356f81 --- /dev/null +++ b/tests/cases/standalone/common/promql/native_time_selection.result @@ -0,0 +1,361 @@ +-- Regression coverage for instant and range selection on native microsecond and +-- nanosecond time indexes. +CREATE TABLE native_time_us ( + ts TIMESTAMP(6) TIME INDEX, + series STRING PRIMARY KEY, + val DOUBLE, +); + +Affected Rows: 0 + +INSERT INTO native_time_us VALUES + (1000001, 'future', 101), + (1000000, 'exact', 201), + (1000001, 'exact', 202), + (-299000000, 'lowerbound', 301), + (-298999999, 'lowerplus', 302), + (1000000, 'positive_lowerbound', 701), + (1000001, 'positive_lowerplus', 702), + (1000001, 'multi', 401), + (-1000000, 'offset', 501), + (0, 'offset', 502), + (1000000, 'offset', 503), + (999999, 'past', 602), + (999001, 'past', 601), + (0, 'window', 1), + (1, 'window', 2), + (2, 'window', 5), + (1000000, 'window', 3), + (1000001, 'window', 4); + +Affected Rows: 18 + +-- Future-only selection is empty before flushing, exercising the memtable path. +TQL EVAL (1, 1, '1s', '300s') native_time_us{series="future"}; + +++ +++ + +ADMIN FLUSH_TABLE('native_time_us'); + ++-------------------------------------+ +| ADMIN FLUSH_TABLE('native_time_us') | ++-------------------------------------+ +| 0 | ++-------------------------------------+ + +-- At 1s, selection keeps an exact native timestamp. +TQL EVAL (1, 1, '1s', '300s') native_time_us{series="exact"}; + ++-------+--------+---------------------+ +| val | series | ts | ++-------+--------+---------------------+ +| 201.0 | exact | 1970-01-01T00:00:01 | ++-------+--------+---------------------+ + +TQL EVAL (1, 1, '1s', '300s') timestamp(native_time_us{series="future"}); + +++ +++ + +TQL EVAL (1, 1, '1s', '300s') timestamp(native_time_us{series="exact"}); + ++---------------------+-------+--------+ +| ts | value | series | ++---------------------+-------+--------+ +| 1970-01-01T00:00:01 | 1.0 | exact | ++---------------------+-------+--------+ + +-- Instant lookback bounds are exclusive: these return only 302 and 702. +TQL EVAL (1, 1, '1s', '300s') native_time_us{series=~"lower.*"}; + ++-------+-----------+---------------------+ +| val | series | ts | ++-------+-----------+---------------------+ +| 302.0 | lowerplus | 1970-01-01T00:00:01 | ++-------+-----------+---------------------+ + +TQL EVAL (301, 301, '1s', '300s') native_time_us{series=~"positive_lower.*"}; + ++-------+--------------------+---------------------+ +| val | series | ts | ++-------+--------------------+---------------------+ +| 702.0 | positive_lowerplus | 1970-01-01T00:05:01 | ++-------+--------------------+---------------------+ + +-- The sub-millisecond point belongs only to the 2s evaluation step. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (1, 2, '1s', '300s') native_time_us{series="multi"}; + ++-------+--------+---------------------+ +| val | series | ts | ++-------+--------+---------------------+ +| 401.0 | multi | 1970-01-01T00:00:02 | ++-------+--------+---------------------+ + +-- The latest native timestamp below 1s is retained even when inserts are unordered. +TQL EVAL (1, 1, '1s', '300s') native_time_us{series="past"}; + ++-------+--------+---------------------+ +| val | series | ts | ++-------+--------+---------------------+ +| 602.0 | past | 1970-01-01T00:00:01 | ++-------+--------+---------------------+ + +-- Offsets select native timestamps, including stored negative time. +TQL EVAL (0, 0, '1s', '300s') native_time_us{series="offset"}; + ++-------+--------+---------------------+ +| val | series | ts | ++-------+--------+---------------------+ +| 502.0 | offset | 1970-01-01T00:00:00 | ++-------+--------+---------------------+ + +TQL EVAL (0, 0, '1s', '300s') native_time_us{series="offset"} offset 1s; + ++-------+--------+---------------------+ +| val | series | ts | ++-------+--------+---------------------+ +| 501.0 | offset | 1970-01-01T00:00:00 | ++-------+--------+---------------------+ + +TQL EVAL (0, 0, '1s', '300s') native_time_us{series="offset"} offset -1s; + ++-------+--------+---------------------+ +| val | series | ts | ++-------+--------+---------------------+ +| 503.0 | offset | 1970-01-01T00:00:00 | ++-------+--------+---------------------+ + +-- [1s] at 1s excludes 0 and 1s+tick, retaining 0+tick, 0+2ticks, and 1s. +TQL EVAL (1, 1, '1s', '300s') count_over_time(native_time_us{series="window"}[1s]); + ++---------------------+------------------------------------+--------+ +| ts | prom_count_over_time(ts_range,val) | series | ++---------------------+------------------------------------+--------+ +| 1970-01-01T00:00:01 | 3.0 | window | ++---------------------+------------------------------------+--------+ + +TQL EVAL (1, 1, '1s', '300s') sum_over_time(native_time_us{series="window"}[1s]); + ++---------------------+----------------------------------+--------+ +| ts | prom_sum_over_time(ts_range,val) | series | ++---------------------+----------------------------------+--------+ +| 1970-01-01T00:00:01 | 10.0 | window | ++---------------------+----------------------------------+--------+ + +TQL EVAL (1, 1, '1s', '300s') last_over_time(native_time_us{series="window"}[1s]); + ++---------------------+-----------------------------------+--------+ +| ts | prom_last_over_time(ts_range,val) | series | ++---------------------+-----------------------------------+--------+ +| 1970-01-01T00:00:01 | 3.0 | window | ++---------------------+-----------------------------------+--------+ + +-- The inner selector consumes native time; the subquery consumes ms evaluations. +TQL EVAL (1, 1, '1s') last_over_time((native_time_us{series="exact"})[1s:1s]); + ++---------------------+-----------------------------------+--------+ +| ts | prom_last_over_time(ts_range,val) | series | ++---------------------+-----------------------------------+--------+ +| 1970-01-01T00:00:01 | 201.0 | exact | ++---------------------+-----------------------------------+--------+ + +DROP TABLE native_time_us; + +Affected Rows: 0 + +CREATE TABLE native_time_ns ( + ts TIMESTAMP(9) TIME INDEX, + series STRING PRIMARY KEY, + val DOUBLE, +); + +Affected Rows: 0 + +INSERT INTO native_time_ns VALUES + (1000000001, 'future', 101), + (1000000000, 'exact', 201), + (1000000001, 'exact', 202), + (-299000000000, 'lowerbound', 301), + (-298999999999, 'lowerplus', 302), + (1000000000, 'positive_lowerbound', 701), + (1000000001, 'positive_lowerplus', 702), + (1000000001, 'multi', 401), + (-1000000000, 'offset', 501), + (0, 'offset', 502), + (1000000000, 'offset', 503), + (999999000, 'past', 602), + (999001000, 'past', 601), + (0, 'window', 1), + (1, 'window', 2), + (2, 'window', 5), + (1000000000, 'window', 3), + (1000000001, 'window', 4); + +Affected Rows: 18 + +-- Future-only selection is empty before flushing, exercising the memtable path. +TQL EVAL (1, 1, '1s', '300s') native_time_ns{series="future"}; + +++ +++ + +ADMIN FLUSH_TABLE('native_time_ns'); + ++-------------------------------------+ +| ADMIN FLUSH_TABLE('native_time_ns') | ++-------------------------------------+ +| 0 | ++-------------------------------------+ + +-- At 1s, selection keeps an exact native timestamp. +TQL EVAL (1, 1, '1s', '300s') native_time_ns{series="exact"}; + ++-------+--------+---------------------+ +| val | series | ts | ++-------+--------+---------------------+ +| 201.0 | exact | 1970-01-01T00:00:01 | ++-------+--------+---------------------+ + +TQL EVAL (1, 1, '1s', '300s') timestamp(native_time_ns{series="future"}); + +++ +++ + +TQL EVAL (1, 1, '1s', '300s') timestamp(native_time_ns{series="exact"}); + ++---------------------+-------+--------+ +| ts | value | series | ++---------------------+-------+--------+ +| 1970-01-01T00:00:01 | 1.0 | exact | ++---------------------+-------+--------+ + +-- Instant lookback bounds are exclusive: these return only 302 and 702. +TQL EVAL (1, 1, '1s', '300s') native_time_ns{series=~"lower.*"}; + ++-------+-----------+---------------------+ +| val | series | ts | ++-------+-----------+---------------------+ +| 302.0 | lowerplus | 1970-01-01T00:00:01 | ++-------+-----------+---------------------+ + +TQL EVAL (301, 301, '1s', '300s') native_time_ns{series=~"positive_lower.*"}; + ++-------+--------------------+---------------------+ +| val | series | ts | ++-------+--------------------+---------------------+ +| 702.0 | positive_lowerplus | 1970-01-01T00:05:01 | ++-------+--------------------+---------------------+ + +-- The sub-millisecond point belongs only to the 2s evaluation step. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (1, 2, '1s', '300s') native_time_ns{series="multi"}; + ++-------+--------+---------------------+ +| val | series | ts | ++-------+--------+---------------------+ +| 401.0 | multi | 1970-01-01T00:00:02 | ++-------+--------+---------------------+ + +-- The latest native timestamp below 1s is retained even when inserts are unordered. +TQL EVAL (1, 1, '1s', '300s') native_time_ns{series="past"}; + ++-------+--------+---------------------+ +| val | series | ts | ++-------+--------+---------------------+ +| 602.0 | past | 1970-01-01T00:00:01 | ++-------+--------+---------------------+ + +-- Offsets select native timestamps, including stored negative time. +TQL EVAL (0, 0, '1s', '300s') native_time_ns{series="offset"}; + ++-------+--------+---------------------+ +| val | series | ts | ++-------+--------+---------------------+ +| 502.0 | offset | 1970-01-01T00:00:00 | ++-------+--------+---------------------+ + +TQL EVAL (0, 0, '1s', '300s') native_time_ns{series="offset"} offset 1s; + ++-------+--------+---------------------+ +| val | series | ts | ++-------+--------+---------------------+ +| 501.0 | offset | 1970-01-01T00:00:00 | ++-------+--------+---------------------+ + +TQL EVAL (0, 0, '1s', '300s') native_time_ns{series="offset"} offset -1s; + ++-------+--------+---------------------+ +| val | series | ts | ++-------+--------+---------------------+ +| 503.0 | offset | 1970-01-01T00:00:00 | ++-------+--------+---------------------+ + +-- [1s] at 1s excludes 0 and 1s+tick, retaining 0+tick, 0+2ticks, and 1s. +TQL EVAL (1, 1, '1s', '300s') count_over_time(native_time_ns{series="window"}[1s]); + ++---------------------+------------------------------------+--------+ +| ts | prom_count_over_time(ts_range,val) | series | ++---------------------+------------------------------------+--------+ +| 1970-01-01T00:00:01 | 3.0 | window | ++---------------------+------------------------------------+--------+ + +TQL EVAL (1, 1, '1s', '300s') sum_over_time(native_time_ns{series="window"}[1s]); + ++---------------------+----------------------------------+--------+ +| ts | prom_sum_over_time(ts_range,val) | series | ++---------------------+----------------------------------+--------+ +| 1970-01-01T00:00:01 | 10.0 | window | ++---------------------+----------------------------------+--------+ + +TQL EVAL (1, 1, '1s', '300s') last_over_time(native_time_ns{series="window"}[1s]); + ++---------------------+-----------------------------------+--------+ +| ts | prom_last_over_time(ts_range,val) | series | ++---------------------+-----------------------------------+--------+ +| 1970-01-01T00:00:01 | 3.0 | window | ++---------------------+-----------------------------------+--------+ + +-- The inner selector consumes native time; the subquery consumes ms evaluations. +TQL EVAL (1, 1, '1s') last_over_time((native_time_ns{series="exact"})[1s:1s]); + ++---------------------+-----------------------------------+--------+ +| ts | prom_last_over_time(ts_range,val) | series | ++---------------------+-----------------------------------+--------+ +| 1970-01-01T00:00:01 | 201.0 | exact | ++---------------------+-----------------------------------+--------+ + +DROP TABLE native_time_ns; + +Affected Rows: 0 + +-- Second precision is promoted before applying fractional-second offsets. +CREATE TABLE native_time_sec (ts TIMESTAMP(0) TIME INDEX, val DOUBLE); + +Affected Rows: 0 + +INSERT INTO native_time_sec VALUES (0, 10), (1, 11), (2, 12); + +Affected Rows: 3 + +TQL EVAL (1, 1, '1s', '1s') native_time_sec offset 500ms; + ++------+---------------------+ +| val | ts | ++------+---------------------+ +| 10.0 | 1970-01-01T00:00:01 | ++------+---------------------+ + +TQL EVAL (1, 1, '1s', '1s') native_time_sec offset -500ms; + ++------+---------------------+ +| val | ts | ++------+---------------------+ +| 11.0 | 1970-01-01T00:00:01 | ++------+---------------------+ + +DROP TABLE native_time_sec; + +Affected Rows: 0 + diff --git a/tests/cases/standalone/common/promql/native_time_selection.sql b/tests/cases/standalone/common/promql/native_time_selection.sql new file mode 100644 index 0000000000..d6ec7d3d13 --- /dev/null +++ b/tests/cases/standalone/common/promql/native_time_selection.sql @@ -0,0 +1,133 @@ +-- Regression coverage for instant and range selection on native microsecond and +-- nanosecond time indexes. + +CREATE TABLE native_time_us ( + ts TIMESTAMP(6) TIME INDEX, + series STRING PRIMARY KEY, + val DOUBLE, +); + +INSERT INTO native_time_us VALUES + (1000001, 'future', 101), + (1000000, 'exact', 201), + (1000001, 'exact', 202), + (-299000000, 'lowerbound', 301), + (-298999999, 'lowerplus', 302), + (1000000, 'positive_lowerbound', 701), + (1000001, 'positive_lowerplus', 702), + (1000001, 'multi', 401), + (-1000000, 'offset', 501), + (0, 'offset', 502), + (1000000, 'offset', 503), + (999999, 'past', 602), + (999001, 'past', 601), + (0, 'window', 1), + (1, 'window', 2), + (2, 'window', 5), + (1000000, 'window', 3), + (1000001, 'window', 4); + +-- Future-only selection is empty before flushing, exercising the memtable path. +TQL EVAL (1, 1, '1s', '300s') native_time_us{series="future"}; + +ADMIN FLUSH_TABLE('native_time_us'); + +-- At 1s, selection keeps an exact native timestamp. +TQL EVAL (1, 1, '1s', '300s') native_time_us{series="exact"}; +TQL EVAL (1, 1, '1s', '300s') timestamp(native_time_us{series="future"}); +TQL EVAL (1, 1, '1s', '300s') timestamp(native_time_us{series="exact"}); + +-- Instant lookback bounds are exclusive: these return only 302 and 702. +TQL EVAL (1, 1, '1s', '300s') native_time_us{series=~"lower.*"}; +TQL EVAL (301, 301, '1s', '300s') native_time_us{series=~"positive_lower.*"}; + +-- The sub-millisecond point belongs only to the 2s evaluation step. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (1, 2, '1s', '300s') native_time_us{series="multi"}; + +-- The latest native timestamp below 1s is retained even when inserts are unordered. +TQL EVAL (1, 1, '1s', '300s') native_time_us{series="past"}; + +-- Offsets select native timestamps, including stored negative time. +TQL EVAL (0, 0, '1s', '300s') native_time_us{series="offset"}; +TQL EVAL (0, 0, '1s', '300s') native_time_us{series="offset"} offset 1s; +TQL EVAL (0, 0, '1s', '300s') native_time_us{series="offset"} offset -1s; + +-- [1s] at 1s excludes 0 and 1s+tick, retaining 0+tick, 0+2ticks, and 1s. +TQL EVAL (1, 1, '1s', '300s') count_over_time(native_time_us{series="window"}[1s]); +TQL EVAL (1, 1, '1s', '300s') sum_over_time(native_time_us{series="window"}[1s]); +TQL EVAL (1, 1, '1s', '300s') last_over_time(native_time_us{series="window"}[1s]); + +-- The inner selector consumes native time; the subquery consumes ms evaluations. +TQL EVAL (1, 1, '1s') last_over_time((native_time_us{series="exact"})[1s:1s]); + +DROP TABLE native_time_us; + +CREATE TABLE native_time_ns ( + ts TIMESTAMP(9) TIME INDEX, + series STRING PRIMARY KEY, + val DOUBLE, +); + +INSERT INTO native_time_ns VALUES + (1000000001, 'future', 101), + (1000000000, 'exact', 201), + (1000000001, 'exact', 202), + (-299000000000, 'lowerbound', 301), + (-298999999999, 'lowerplus', 302), + (1000000000, 'positive_lowerbound', 701), + (1000000001, 'positive_lowerplus', 702), + (1000000001, 'multi', 401), + (-1000000000, 'offset', 501), + (0, 'offset', 502), + (1000000000, 'offset', 503), + (999999000, 'past', 602), + (999001000, 'past', 601), + (0, 'window', 1), + (1, 'window', 2), + (2, 'window', 5), + (1000000000, 'window', 3), + (1000000001, 'window', 4); + +-- Future-only selection is empty before flushing, exercising the memtable path. +TQL EVAL (1, 1, '1s', '300s') native_time_ns{series="future"}; + +ADMIN FLUSH_TABLE('native_time_ns'); + +-- At 1s, selection keeps an exact native timestamp. +TQL EVAL (1, 1, '1s', '300s') native_time_ns{series="exact"}; +TQL EVAL (1, 1, '1s', '300s') timestamp(native_time_ns{series="future"}); +TQL EVAL (1, 1, '1s', '300s') timestamp(native_time_ns{series="exact"}); + +-- Instant lookback bounds are exclusive: these return only 302 and 702. +TQL EVAL (1, 1, '1s', '300s') native_time_ns{series=~"lower.*"}; +TQL EVAL (301, 301, '1s', '300s') native_time_ns{series=~"positive_lower.*"}; + +-- The sub-millisecond point belongs only to the 2s evaluation step. +-- SQLNESS SORT_RESULT 3 1 +TQL EVAL (1, 2, '1s', '300s') native_time_ns{series="multi"}; + +-- The latest native timestamp below 1s is retained even when inserts are unordered. +TQL EVAL (1, 1, '1s', '300s') native_time_ns{series="past"}; + +-- Offsets select native timestamps, including stored negative time. +TQL EVAL (0, 0, '1s', '300s') native_time_ns{series="offset"}; +TQL EVAL (0, 0, '1s', '300s') native_time_ns{series="offset"} offset 1s; +TQL EVAL (0, 0, '1s', '300s') native_time_ns{series="offset"} offset -1s; + +-- [1s] at 1s excludes 0 and 1s+tick, retaining 0+tick, 0+2ticks, and 1s. +TQL EVAL (1, 1, '1s', '300s') count_over_time(native_time_ns{series="window"}[1s]); +TQL EVAL (1, 1, '1s', '300s') sum_over_time(native_time_ns{series="window"}[1s]); +TQL EVAL (1, 1, '1s', '300s') last_over_time(native_time_ns{series="window"}[1s]); + +-- The inner selector consumes native time; the subquery consumes ms evaluations. +TQL EVAL (1, 1, '1s') last_over_time((native_time_ns{series="exact"})[1s:1s]); + +DROP TABLE native_time_ns; + +-- Second precision is promoted before applying fractional-second offsets. +CREATE TABLE native_time_sec (ts TIMESTAMP(0) TIME INDEX, val DOUBLE); +INSERT INTO native_time_sec VALUES (0, 10), (1, 11), (2, 12); +TQL EVAL (1, 1, '1s', '1s') native_time_sec offset 500ms; +TQL EVAL (1, 1, '1s', '1s') native_time_sec offset -500ms; +DROP TABLE native_time_sec; diff --git a/tests/cases/standalone/common/promql/precisions.result b/tests/cases/standalone/common/promql/precisions.result index e57f5b04ee..4026a35832 100644 --- a/tests/cases/standalone/common/promql/precisions.result +++ b/tests/cases/standalone/common/promql/precisions.result @@ -133,10 +133,8 @@ TQL EVAL (0, 15, '5s') avg_over_time(host_sec{host="host1"}[5s]) + avg_over_time -- Verify that PromQL time predicates on non-millisecond time indexes are -- pushed into the scan as native timestamp range filters. --- Original instant selector filter is built on the millisecond alias: --- host = "host1" AND ts_ms >= -299999ms AND ts_ms <= 10000ms --- After pushing through `CAST(raw_ts AS Timestamp(ms)) AS ts` and applying --- DataFusion cast preimage, it becomes a native half-open range on raw_ts. +-- Instant selection compares raw timestamps before millisecond output conversion: +-- host = "host1" AND ts_us > -300000000us AND ts_us <= 10000000us. -- SQLNESS REPLACE (RoundRobinBatch.*) REDACTED -- SQLNESS REPLACE (peers.*) REDACTED -- SQLNESS REPLACE (Hash.*) REDACTED @@ -152,17 +150,17 @@ TQL EXPLAIN (0, 10, '5s') host_micro{host="host1"}; | | PromInstantManipulate: range=[0..10000], lookback=[300000], interval=[5000], time index=[ts] | | | PromSeriesDivide: tags=["host"] | | | Sort: host_micro.host ASC NULLS FIRST, host_micro.ts ASC NULLS FIRST | -| | Projection: host_micro.val, host_micro.host, CAST(host_micro.ts AS Timestamp(ms)) AS ts | -| | Filter: host_micro.host = Utf8("host1") AND host_micro.ts >= TimestampMicrosecond(-299999999, None) AND host_micro.ts < TimestampMicrosecond(10001000, None) | -| | TableScan: host_micro, partial_filters=[host_micro.host = Utf8("host1"), host_micro.ts >= TimestampMicrosecond(-299999999, None), host_micro.ts < TimestampMicrosecond(10001000, None)] | +| | Projection: host_micro.val, host_micro.host, host_micro.ts | +| | Filter: host_micro.host = Utf8("host1") AND host_micro.ts > TimestampMicrosecond(-300000000, None) AND host_micro.ts <= TimestampMicrosecond(10000000, None) | +| | TableScan: host_micro, partial_filters=[host_micro.host = Utf8("host1"), host_micro.ts > TimestampMicrosecond(-300000000, None), host_micro.ts <= TimestampMicrosecond(10000000, None)] | | | ]] | | physical_plan | CooperativeExec | | | MergeScanExec: REDACTED | | | +---------------+---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+ --- The same instant-selector cast-preimage path should work for nanosecond indexes. --- Expected native bounds: ts_ns >= -299999999999ns AND ts_ns < 10001000000ns. +-- The same exclusive-lower, inclusive-upper window applies to nanosecond indexes. +-- Expected native bounds: ts_ns > -300000000000ns AND ts_ns <= 10000000000ns. -- SQLNESS REPLACE (RoundRobinBatch.*) REDACTED -- SQLNESS REPLACE (peers.*) REDACTED -- SQLNESS REPLACE (Hash.*) REDACTED @@ -177,9 +175,9 @@ TQL EXPLAIN (0, 10, '5s') host_nano{host="host1"}; | | PromInstantManipulate: range=[0..10000], lookback=[300000], interval=[5000], time index=[ts] | | | PromSeriesDivide: tags=["host"] | | | Sort: host_nano.host ASC NULLS FIRST, host_nano.ts ASC NULLS FIRST | -| | Projection: host_nano.val, host_nano.host, CAST(host_nano.ts AS Timestamp(ms)) AS ts | -| | Filter: host_nano.host = Utf8("host1") AND host_nano.ts >= TimestampNanosecond(-299999999999, None) AND host_nano.ts < TimestampNanosecond(10001000000, None) | -| | TableScan: host_nano, partial_filters=[host_nano.host = Utf8("host1"), host_nano.ts >= TimestampNanosecond(-299999999999, None), host_nano.ts < TimestampNanosecond(10001000000, None)] | +| | Projection: host_nano.val, host_nano.host, host_nano.ts | +| | Filter: host_nano.host = Utf8("host1") AND host_nano.ts > TimestampNanosecond(-300000000000, None) AND host_nano.ts <= TimestampNanosecond(10000000000, None) | +| | TableScan: host_nano, partial_filters=[host_nano.host = Utf8("host1"), host_nano.ts > TimestampNanosecond(-300000000000, None), host_nano.ts <= TimestampNanosecond(10000000000, None)] | | | ]] | | physical_plan | CooperativeExec | | | MergeScanExec: REDACTED @@ -187,8 +185,8 @@ TQL EXPLAIN (0, 10, '5s') host_nano{host="host1"}; +---------------+---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+ -- Range selectors use their range window instead of the default lookback. --- Original range selector filter for [5s]: --- host = "host1" AND ts_ms >= -4999ms AND ts_ms <= 10000ms +-- Native range selector filter for [5s]: +-- host = "host1" AND ts_us > -5000000us AND ts_us <= 10000000us -- SQLNESS REPLACE (RoundRobinBatch.*) REDACTED -- SQLNESS REPLACE (peers.*) REDACTED -- SQLNESS REPLACE (Hash.*) REDACTED @@ -206,17 +204,17 @@ TQL EXPLAIN (0, 10, '5s') avg_over_time(host_micro{host="host1"}[5s]); | | PromSeriesNormalize: offset=[0], time index=[ts], filter NaN: [true] | | | PromSeriesDivide: tags=["host"] | | | Sort: host_micro.host ASC NULLS FIRST, host_micro.ts ASC NULLS FIRST | -| | Projection: host_micro.val, host_micro.host, CAST(host_micro.ts AS Timestamp(ms)) AS ts | -| | Filter: host_micro.host = Utf8("host1") AND host_micro.ts >= TimestampMicrosecond(-4999999, None) AND host_micro.ts < TimestampMicrosecond(10001000, None) | -| | TableScan: host_micro, partial_filters=[host_micro.host = Utf8("host1"), host_micro.ts >= TimestampMicrosecond(-4999999, None), host_micro.ts < TimestampMicrosecond(10001000, None)] | +| | Projection: host_micro.val, host_micro.host, host_micro.ts | +| | Filter: host_micro.host = Utf8("host1") AND host_micro.ts > TimestampMicrosecond(-5000000, None) AND host_micro.ts <= TimestampMicrosecond(10000000, None) | +| | TableScan: host_micro, partial_filters=[host_micro.host = Utf8("host1"), host_micro.ts > TimestampMicrosecond(-5000000, None), host_micro.ts <= TimestampMicrosecond(10000000, None)] | | | ]] | | physical_plan | CooperativeExec | | | MergeScanExec: REDACTED | | | +---------------+-------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+ --- The same range-selector cast-preimage path should work for nanosecond indexes. --- Expected native bounds: ts_ns >= -4999999999ns AND ts_ns < 10001000000ns. +-- Range membership also retains nanosecond precision. +-- Expected native bounds: ts_ns > -5000000000ns AND ts_ns <= 10000000000ns. -- SQLNESS REPLACE (RoundRobinBatch.*) REDACTED -- SQLNESS REPLACE (peers.*) REDACTED -- SQLNESS REPLACE (Hash.*) REDACTED @@ -234,9 +232,9 @@ TQL EXPLAIN (0, 10, '5s') avg_over_time(host_nano{host="host1"}[5s]); | | PromSeriesNormalize: offset=[0], time index=[ts], filter NaN: [true] | | | PromSeriesDivide: tags=["host"] | | | Sort: host_nano.host ASC NULLS FIRST, host_nano.ts ASC NULLS FIRST | -| | Projection: host_nano.val, host_nano.host, CAST(host_nano.ts AS Timestamp(ms)) AS ts | -| | Filter: host_nano.host = Utf8("host1") AND host_nano.ts >= TimestampNanosecond(-4999999999, None) AND host_nano.ts < TimestampNanosecond(10001000000, None) | -| | TableScan: host_nano, partial_filters=[host_nano.host = Utf8("host1"), host_nano.ts >= TimestampNanosecond(-4999999999, None), host_nano.ts < TimestampNanosecond(10001000000, None)] | +| | Projection: host_nano.val, host_nano.host, host_nano.ts | +| | Filter: host_nano.host = Utf8("host1") AND host_nano.ts > TimestampNanosecond(-5000000000, None) AND host_nano.ts <= TimestampNanosecond(10000000000, None) | +| | TableScan: host_nano, partial_filters=[host_nano.host = Utf8("host1"), host_nano.ts > TimestampNanosecond(-5000000000, None), host_nano.ts <= TimestampNanosecond(10000000000, None)] | | | ]] | | physical_plan | CooperativeExec | | | MergeScanExec: REDACTED diff --git a/tests/cases/standalone/common/promql/precisions.sql b/tests/cases/standalone/common/promql/precisions.sql index 01e6cae9fe..8f7fec37a8 100644 --- a/tests/cases/standalone/common/promql/precisions.sql +++ b/tests/cases/standalone/common/promql/precisions.sql @@ -68,10 +68,8 @@ TQL EVAL (0, 15, '5s') avg_over_time(host_sec{host="host1"}[5s]) + avg_over_time -- Verify that PromQL time predicates on non-millisecond time indexes are -- pushed into the scan as native timestamp range filters. --- Original instant selector filter is built on the millisecond alias: --- host = "host1" AND ts_ms >= -299999ms AND ts_ms <= 10000ms --- After pushing through `CAST(raw_ts AS Timestamp(ms)) AS ts` and applying --- DataFusion cast preimage, it becomes a native half-open range on raw_ts. +-- Instant selection compares raw timestamps before millisecond output conversion: +-- host = "host1" AND ts_us > -300000000us AND ts_us <= 10000000us. -- SQLNESS REPLACE (RoundRobinBatch.*) REDACTED -- SQLNESS REPLACE (peers.*) REDACTED -- SQLNESS REPLACE (Hash.*) REDACTED @@ -80,8 +78,8 @@ TQL EVAL (0, 15, '5s') avg_over_time(host_sec{host="host1"}[5s]) + avg_over_time -- SQLNESS REPLACE host_nano.__table_id\s*=\s*UInt32\(\d+\) host_nano.__table_id=UInt32(REDACTED) TQL EXPLAIN (0, 10, '5s') host_micro{host="host1"}; --- The same instant-selector cast-preimage path should work for nanosecond indexes. --- Expected native bounds: ts_ns >= -299999999999ns AND ts_ns < 10001000000ns. +-- The same exclusive-lower, inclusive-upper window applies to nanosecond indexes. +-- Expected native bounds: ts_ns > -300000000000ns AND ts_ns <= 10000000000ns. -- SQLNESS REPLACE (RoundRobinBatch.*) REDACTED -- SQLNESS REPLACE (peers.*) REDACTED -- SQLNESS REPLACE (Hash.*) REDACTED @@ -90,8 +88,8 @@ TQL EXPLAIN (0, 10, '5s') host_micro{host="host1"}; TQL EXPLAIN (0, 10, '5s') host_nano{host="host1"}; -- Range selectors use their range window instead of the default lookback. --- Original range selector filter for [5s]: --- host = "host1" AND ts_ms >= -4999ms AND ts_ms <= 10000ms +-- Native range selector filter for [5s]: +-- host = "host1" AND ts_us > -5000000us AND ts_us <= 10000000us -- SQLNESS REPLACE (RoundRobinBatch.*) REDACTED -- SQLNESS REPLACE (peers.*) REDACTED -- SQLNESS REPLACE (Hash.*) REDACTED @@ -99,8 +97,8 @@ TQL EXPLAIN (0, 10, '5s') host_nano{host="host1"}; -- SQLNESS REPLACE host_micro.__table_id\s*=\s*UInt32\(\d+\) host_micro.__table_id=UInt32(REDACTED) TQL EXPLAIN (0, 10, '5s') avg_over_time(host_micro{host="host1"}[5s]); --- The same range-selector cast-preimage path should work for nanosecond indexes. --- Expected native bounds: ts_ns >= -4999999999ns AND ts_ns < 10001000000ns. +-- Range membership also retains nanosecond precision. +-- Expected native bounds: ts_ns > -5000000000ns AND ts_ns <= 10000000000ns. -- SQLNESS REPLACE (RoundRobinBatch.*) REDACTED -- SQLNESS REPLACE (peers.*) REDACTED -- SQLNESS REPLACE (Hash.*) REDACTED diff --git a/tests/cases/standalone/common/tql-explain-analyze/explain.result b/tests/cases/standalone/common/tql-explain-analyze/explain.result index 2d6eab49da..6b03513a11 100644 --- a/tests/cases/standalone/common/tql-explain-analyze/explain.result +++ b/tests/cases/standalone/common/tql-explain-analyze/explain.result @@ -391,9 +391,9 @@ TQL EXPLAIN (0, 10, '5s') test_nano; | | PromInstantManipulate: range=[0..10000], lookback=[300000], interval=[5000], time index=[j] | | | PromSeriesDivide: tags=["k"] | | | Sort: test_nano.k ASC NULLS FIRST, test_nano.j ASC NULLS FIRST | -| | Projection: test_nano.i, test_nano.k, CAST(test_nano.j AS Timestamp(ms)) AS j | -| | Filter: test_nano.j >= TimestampNanosecond(-299999999999, None) AND test_nano.j < TimestampNanosecond(10001000000, None) | -| | TableScan: test_nano, partial_filters=[test_nano.j >= TimestampNanosecond(-299999999999, None), test_nano.j < TimestampNanosecond(10001000000, None)] | +| | Projection: test_nano.i, test_nano.k, test_nano.j | +| | Filter: test_nano.j > TimestampNanosecond(-300000000000, None) AND test_nano.j <= TimestampNanosecond(10000000000, None) | +| | TableScan: test_nano, partial_filters=[test_nano.j > TimestampNanosecond(-300000000000, None), test_nano.j <= TimestampNanosecond(10000000000, None)] | | | ]] | | physical_plan | CooperativeExec | | | MergeScanExec: REDACTED @@ -415,8 +415,8 @@ TQL EXPLAIN VERBOSE (0, 10, '5s') test_nano; | initial_logical_plan_| PromInstantManipulate: range=[0..10000], lookback=[300000], interval=[5000], time index=[j]_| |_|_PromSeriesDivide: tags=["k"]_| |_|_Sort: test_nano.k ASC NULLS FIRST, test_nano.j ASC NULLS FIRST_| -|_|_Filter: test_nano.j >= TimestampMillisecond(-299999, None) AND test_nano.j <= TimestampMillisecond(10000, None)_| -|_|_Projection: test_nano.i, test_nano.k, CAST(test_nano.j AS Timestamp(ms)) AS j_| +|_|_Filter: test_nano.j > TimestampNanosecond(-300000000000, None) AND test_nano.j <= TimestampNanosecond(10000000000, None)_| +|_|_Projection: test_nano.i, test_nano.k, test_nano.j_| |_|_TableScan: test_nano_| | logical_plan after apply_function_rewrites_| SAME TEXT AS ABOVE_| | logical_plan after count_wildcard_to_time_index_rule_| SAME TEXT AS ABOVE_| @@ -430,9 +430,9 @@ TQL EXPLAIN VERBOSE (0, 10, '5s') test_nano; |_| PromInstantManipulate: range=[0..10000], lookback=[300000], interval=[5000], time index=[j]_| |_|_PromSeriesDivide: tags=["k"]_| |_|_Sort: test_nano.k ASC NULLS FIRST, test_nano.j ASC NULLS FIRST_| -|_|_Projection: test_nano.i, test_nano.k, CAST(test_nano.j AS Timestamp(ms)) AS j_| -|_|_Filter: test_nano.j >= TimestampNanosecond(-299999999999, None) AND test_nano.j < TimestampNanosecond(10001000000, None)_| -|_|_TableScan: test_nano, partial_filters=[test_nano.j >= TimestampNanosecond(-299999999999, None), test_nano.j < TimestampNanosecond(10001000000, None)] | +|_|_Projection: test_nano.i, test_nano.k, test_nano.j_| +|_|_Filter: test_nano.j > TimestampNanosecond(-300000000000, None) AND test_nano.j <= TimestampNanosecond(10000000000, None)_| +|_|_TableScan: test_nano, partial_filters=[test_nano.j > TimestampNanosecond(-300000000000, None), test_nano.j <= TimestampNanosecond(10000000000, None)] | |_| ]]_| | logical_plan after JsonSchemaConcretizeRule_| SAME TEXT AS ABOVE_| | logical_plan after FixStateUdafOrderingAnalyzer_| SAME TEXT AS ABOVE_| @@ -464,9 +464,9 @@ TQL EXPLAIN VERBOSE (0, 10, '5s') test_nano; |_| PromInstantManipulate: range=[0..10000], lookback=[300000], interval=[5000], time index=[j]_| |_|_PromSeriesDivide: tags=["k"]_| |_|_Sort: test_nano.k ASC NULLS FIRST, test_nano.j ASC NULLS FIRST_| -|_|_Projection: test_nano.i, test_nano.k, CAST(test_nano.j AS Timestamp(ms)) AS j_| -|_|_Filter: test_nano.j >= TimestampNanosecond(-299999999999, None) AND test_nano.j < TimestampNanosecond(10001000000, None)_| -|_|_TableScan: test_nano, partial_filters=[test_nano.j >= TimestampNanosecond(-299999999999, None), test_nano.j < TimestampNanosecond(10001000000, None)] | +|_|_Projection: test_nano.i, test_nano.k, test_nano.j_| +|_|_Filter: test_nano.j > TimestampNanosecond(-300000000000, None) AND test_nano.j <= TimestampNanosecond(10000000000, None)_| +|_|_TableScan: test_nano, partial_filters=[test_nano.j > TimestampNanosecond(-300000000000, None), test_nano.j <= TimestampNanosecond(10000000000, None)] | |_| ]]_| | logical_plan after ScanHintRule_| SAME TEXT AS ABOVE_| | logical_plan after JsonTypeConcretizeRule_| SAME TEXT AS ABOVE_| @@ -500,9 +500,9 @@ TQL EXPLAIN VERBOSE (0, 10, '5s') test_nano; |_| PromInstantManipulate: range=[0..10000], lookback=[300000], interval=[5000], time index=[j]_| |_|_PromSeriesDivide: tags=["k"]_| |_|_Sort: test_nano.k ASC NULLS FIRST, test_nano.j ASC NULLS FIRST_| -|_|_Projection: test_nano.i, test_nano.k, CAST(test_nano.j AS Timestamp(ms)) AS j_| -|_|_Filter: test_nano.j >= TimestampNanosecond(-299999999999, None) AND test_nano.j < TimestampNanosecond(10001000000, None)_| -|_|_TableScan: test_nano, partial_filters=[test_nano.j >= TimestampNanosecond(-299999999999, None), test_nano.j < TimestampNanosecond(10001000000, None)] | +|_|_Projection: test_nano.i, test_nano.k, test_nano.j_| +|_|_Filter: test_nano.j > TimestampNanosecond(-300000000000, None) AND test_nano.j <= TimestampNanosecond(10000000000, None)_| +|_|_TableScan: test_nano, partial_filters=[test_nano.j > TimestampNanosecond(-300000000000, None), test_nano.j <= TimestampNanosecond(10000000000, None)] | |_| ]]_| | initial_physical_plan_| MergeScanExec: REDACTED |_|_| diff --git a/tests/cases/standalone/common/tql/general_table.result b/tests/cases/standalone/common/tql/general_table.result index 6172f40aaa..61d414c333 100644 --- a/tests/cases/standalone/common/tql/general_table.result +++ b/tests/cases/standalone/common/tql/general_table.result @@ -42,7 +42,7 @@ TQL analyze (0, 10, '1s') sum by(job) (irate(cpu_usage{job="fire"}[5s])) / 1e9; |_|_|_PromRangeManipulateExec: req range=[0..10000], interval=[1000], eval range=[5000], time index=[ts] REDACTED |_|_|_PromSeriesNormalizeExec: offset=[0], time index=[ts], filter NaN: [true] REDACTED |_|_|_PromSeriesDivideExec: tags=["job"] REDACTED -|_|_|_ProjectionExec: expr=[value@1 as value, job@0 as job, CAST(ts@2 AS Timestamp(ms)) as ts] REDACTED +|_|_|_ProjectionExec: expr=[value@1 as value, job@0 as job, ts@2 as ts] REDACTED |_|_|_ScanExec: REDACTED |_|_|_| |_|_| Total rows: 0_|