Merge branch 'fix/promql-native-time-selection' into work/9070-native-lastrow

This commit is contained in:
discord9
2026-09-10 16:33:50 +08:00
11 changed files with 2263 additions and 353 deletions
+188 -3
View File
@@ -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 = <TimestampMillisecondType as ArrowPrimitiveType>::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::<TimestampSecondArray>()
.map(|a| a.values().as_ref()),
DataType::Timestamp(TimeUnit::Millisecond, _) => array
.as_any()
.downcast_ref::<TimestampMillisecondArray>()
.map(|a| a.values().as_ref()),
DataType::Timestamp(TimeUnit::Microsecond, _) => array
.as_any()
.downcast_ref::<TimestampMicrosecondArray>()
.map(|a| a.values().as_ref()),
DataType::Timestamp(TimeUnit::Nanosecond, _) => array
.as_any()
.downcast_ref::<TimestampNanosecondArray>()
.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<TimeUnit> {
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::<SeriesNormalize>()
.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(&divide, "timestamp"));
}
}
File diff suppressed because it is too large Load Diff
+185 -68
View File
@@ -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<Millisecond> {
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<N: AsRef<str>>(
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<RecordBatch> {
let ts_column = input
.column(self.time_index)
.as_any()
.downcast_ref::<TimestampMillisecondArray>()
.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<i64>| -> Arc<dyn Array> {
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::<usize>(),
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::<datafusion::arrow::array::StructArray>()
.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::<DictionaryArray<Int64Type>>()
.unwrap()
.clone(),
)
.unwrap();
let values = values.get(0).unwrap();
let values = values
.as_any()
.downcast_ref::<datafusion::arrow::array::StructArray>()
.unwrap();
let timestamps = batch
.column(0)
.as_any()
.downcast_ref::<TimestampMillisecondArray>()
.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::<DictionaryArray<Int64Type>>()
.unwrap()
.clone(),
)
.unwrap();
assert_eq!(
timestamps.get(0).unwrap().to_data(),
TimestampMillisecondArray::from(vec![2_000; 3]).to_data()
);
}
}
+467 -40
View File
@@ -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<usize>,
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::<std::result::Result<Vec<_>, _>>()?;
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::<TimestampMillisecondArray>()
.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::<TimestampMillisecondArray>()
.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::<DictionaryArray<Int64Type>>()
.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::<DictionaryArray<Int64Type>>()
.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::<TimestampMillisecondArray>()
.unwrap()
.value(0),
start,
"{name}"
);
let values = RangeArray::try_new(
output
.column(1)
.as_any()
.downcast_ref::<DictionaryArray<Int64Type>>()
.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::<DictionaryArray<Int64Type>>()
.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::<DictionaryArray<Int64Type>>()
.unwrap()
.clone(),
)
.unwrap();
let payload = timestamps.get(0).unwrap();
let payload = payload
.as_any()
.downcast_ref::<TimestampMillisecondArray>()
.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![
+201 -87
View File
@@ -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::<Vec<_>>();
project_exprs
.push(build_special_time_expr(&time_index_column).alias(&timestamp_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(&timestamp_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<Option<DfExpr>> {
fn build_time_index_filter(
&self,
offset_duration: i64,
schema: &DFSchemaRef,
) -> Result<Option<DfExpr>> {
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<ScalarValue> {
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<DfExpr> {
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]"
);
}
@@ -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
@@ -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;
@@ -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
@@ -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
@@ -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
|_|_|
@@ -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_|