From a598690bf7805714917bd824dd5bcdf64ebffd70 Mon Sep 17 00:00:00 2001 From: discord9 Date: Thu, 16 Jul 2026 10:29:43 +0800 Subject: [PATCH] fix(promql): handle missing labels in or matching (#8504) * fix(promql): handle missing labels in or matching Signed-off-by: discord9 * test(promql): streamline or matching coverage Signed-off-by: discord9 * fix(promql): preserve unmatched rhs series in or Signed-off-by: discord9 * perf(promql): stream union distinct inputs Signed-off-by: discord9 --------- Signed-off-by: discord9 --- .../src/extension_plan/union_distinct_on.rs | 1649 +++++++++++++---- src/promql/src/lib.rs | 2 - src/query/src/promql/planner.rs | 846 ++++++++- .../common/promql/set_operation.result | 24 +- .../promql/tsid_binary_join_regression.result | 24 +- 5 files changed, 2048 insertions(+), 497 deletions(-) diff --git a/src/promql/src/extension_plan/union_distinct_on.rs b/src/promql/src/extension_plan/union_distinct_on.rs index 149475bfe4..a7d8fd4859 100644 --- a/src/promql/src/extension_plan/union_distinct_on.rs +++ b/src/promql/src/extension_plan/union_distinct_on.rs @@ -17,7 +17,7 @@ use std::pin::Pin; use std::sync::Arc; use std::task::{Context, Poll}; -use ahash::{HashMap, RandomState}; +use ahash::{HashSet, RandomState}; use datafusion::arrow::array::UInt64Array; use datafusion::arrow::datatypes::SchemaRef; use datafusion::arrow::record_batch::RecordBatch; @@ -34,14 +34,12 @@ use datafusion::physical_plan::{ }; use datafusion_expr::col; use datatypes::arrow::compute; -use futures::future::BoxFuture; -use futures::{Stream, StreamExt, TryStreamExt, ready}; +use futures::{Stream, StreamExt, ready}; use greptime_proto::substrait_extension as pb; use prost::Message; use snafu::ResultExt; -use crate::error::{DeserializeSnafu, Result}; -use crate::extension_plan::{resolve_column_name, serialize_column_index}; +use crate::error::{DataFusionPlanningSnafu, DeserializeSnafu, Result}; /// A special kind of `UNION`(`OR` in PromQL) operator, for PromQL specific use case. /// @@ -49,35 +47,23 @@ use crate::extension_plan::{resolve_column_name, serialize_column_index}; /// most different part is that it treat left child and right child differently: /// - All columns from left child will be outputted. /// - Only check collisions (when not distinct) on the columns specified by `compare_keys`. -/// - When there is a collision: -/// - If the collision is from right child itself, only the first observed row will be -/// preserved. All others are discarded. -/// - If the collision is from left child, the row in right child will be discarded. -/// - The output order is not maintained. This plan will output left child first, then right child. -/// - The output schema contains all columns from left or right child plans. +/// - Rows from the right child with the same comparison signature are all preserved. +/// - If a signature occurs in the left child, all matching rows from the right child are discarded. +/// - Output preserves each input's batch and row order, with all left rows before right rows. +/// - The output schema is based on the left child schema, with nullability widened from both inputs. /// -/// From the implementation perspective, this operator is similar to `HashJoin`, but the -/// probe side is the right child, and the build side is the left child. Another difference -/// is that the probe is opting-out. -/// -/// This plan will exhaust the right child first to build probe hash table, then streaming -/// on left side, and use the left side to "mask" the hash table. +/// Execution streams the left child first while retaining its comparison signatures. The +/// right child is polled only after the left child completes, then rows whose signatures were +/// observed on the left are omitted. #[derive(Debug, PartialEq, Eq, Hash)] pub struct UnionDistinctOn { left: LogicalPlan, right: LogicalPlan, /// The columns to compare for equality. /// TIME INDEX is included. - compare_keys: Vec, - ts_col: String, + compare_key_indices: Vec, + ts_col_idx: usize, output_schema: DFSchemaRef, - unfix: Option, -} - -#[derive(Debug, PartialEq, Eq, Hash, PartialOrd)] -struct UnfixIndices { - pub compare_key_indices: Vec, - pub ts_col_idx: u64, } impl UnionDistinctOn { @@ -85,21 +71,92 @@ impl UnionDistinctOn { "UnionDistinctOn" } - pub fn new( + pub fn try_new( left: LogicalPlan, right: LogicalPlan, - compare_keys: Vec, - ts_col: String, - output_schema: DFSchemaRef, - ) -> Self { - Self { + compare_key_indices: Vec, + ts_col_idx: usize, + ) -> DataFusionResult { + let output_schema = + Self::validate_children(&left, &right, &compare_key_indices, ts_col_idx)?; + Ok(Self { left, right, - compare_keys, - ts_col, + compare_key_indices, + ts_col_idx, output_schema, - unfix: None, + }) + } + + fn validate_children( + left: &LogicalPlan, + right: &LogicalPlan, + compare_key_indices: &[usize], + ts_col_idx: usize, + ) -> DataFusionResult { + let left_schema = left.schema(); + let right_schema = right.schema(); + let left_fields = left_schema.fields(); + let right_fields = right_schema.fields(); + + if left_fields.len() != right_fields.len() { + return Err(DataFusionError::Plan(format!( + "UnionDistinctOn inputs have different field counts: left={}, right={}", + left_fields.len(), + right_fields.len() + ))); } + + for (column_type, index) in compare_key_indices + .iter() + .map(|index| ("compare key", *index)) + .chain(std::iter::once(("timestamp", ts_col_idx))) + { + if index >= left_fields.len() || index >= right_fields.len() { + return Err(DataFusionError::Plan(format!( + "UnionDistinctOn {column_type} index {index} is out of bounds for inputs with {} fields", + left_fields.len() + ))); + } + } + + for (index, (left_field, right_field)) in left_fields.iter().zip(right_fields).enumerate() { + if left_field.data_type() != right_field.data_type() { + return Err(DataFusionError::Plan(format!( + "UnionDistinctOn input field at index {index} has incompatible data types: left={:?}, right={:?}", + left_field.data_type(), + right_field.data_type() + ))); + } + } + + let output_fields = left_fields + .iter() + .zip(right_fields) + .enumerate() + .map(|(index, (left_field, right_field))| { + let (qualifier, _) = left_schema.qualified_field(index); + ( + qualifier.cloned(), + Arc::new( + left_field + .as_ref() + .clone() + .with_nullable(left_field.is_nullable() || right_field.is_nullable()), + ), + ) + }) + .collect(); + let output_schema = + DFSchema::new_with_metadata(output_fields, left_schema.metadata().clone()).map_err( + |error| { + DataFusionError::Plan(format!( + "Failed to construct UnionDistinctOn output schema: {error}" + )) + }, + )?; + + Ok(Arc::new(output_schema)) } pub fn to_execution_plan( @@ -117,8 +174,8 @@ impl UnionDistinctOn { Arc::new(UnionDistinctOnExec { left: left_exec, right: right_exec, - compare_keys: self.compare_keys.clone(), - ts_col: self.ts_col.clone(), + compare_key_indices: self.compare_key_indices.clone(), + ts_col_idx: self.ts_col_idx, output_schema, metric: ExecutionPlanMetricsSet::new(), properties, @@ -128,12 +185,11 @@ impl UnionDistinctOn { pub fn serialize(&self) -> Vec { let compare_key_indices = self - .compare_keys + .compare_key_indices .iter() - .map(|name| serialize_column_index(&self.output_schema, name)) - .collect::>(); - - let ts_col_idx = serialize_column_index(&self.output_schema, &self.ts_col); + .map(|index| u64::try_from(*index).expect("usize always fits in u64")) + .collect(); + let ts_col_idx = u64::try_from(self.ts_col_idx).expect("usize always fits in u64"); pb::UnionDistinctOn { compare_key_indices, @@ -149,18 +205,33 @@ impl UnionDistinctOn { schema: Arc::new(DFSchema::empty()), }); - let unfix = UnfixIndices { - compare_key_indices: pb_union.compare_key_indices.clone(), - ts_col_idx: pb_union.ts_col_idx, - }; + let compare_key_indices = pb_union + .compare_key_indices + .into_iter() + .map(|index| { + usize::try_from(index).map_err(|_| { + DataFusionError::Plan(format!( + "UnionDistinctOn compare key index {index} does not fit in usize" + )) + }) + }) + .collect::>>() + .context(DataFusionPlanningSnafu)?; + let ts_col_idx = usize::try_from(pb_union.ts_col_idx) + .map_err(|_| { + DataFusionError::Plan(format!( + "UnionDistinctOn timestamp index {} does not fit in usize", + pb_union.ts_col_idx + )) + }) + .context(DataFusionPlanningSnafu)?; Ok(Self { left: placeholder_plan.clone(), right: placeholder_plan, - compare_keys: Vec::new(), - ts_col: String::new(), + compare_key_indices, + ts_col_idx, output_schema: Arc::new(DFSchema::empty()), - unfix: Some(unfix), }) } } @@ -176,11 +247,14 @@ impl PartialOrd for UnionDistinctOn { Some(core::cmp::Ordering::Equal) => {} ord => return ord, } - match self.compare_keys.partial_cmp(&other.compare_keys) { + match self + .compare_key_indices + .partial_cmp(&other.compare_key_indices) + { Some(core::cmp::Ordering::Equal) => {} ord => return ord, } - self.ts_col.partial_cmp(&other.ts_col) + self.ts_col_idx.partial_cmp(&other.ts_col_idx) } } @@ -198,22 +272,21 @@ impl UserDefinedLogicalNodeCore for UnionDistinctOn { } fn expressions(&self) -> Vec { - if self.unfix.is_some() { - return vec![]; - } - - let mut exprs: Vec = self.compare_keys.iter().map(col).collect(); - if !self.compare_keys.iter().any(|key| key == &self.ts_col) { - exprs.push(col(&self.ts_col)); + let fields = self.left.schema().fields(); + let mut exprs = self + .compare_key_indices + .iter() + .filter_map(|index| fields.get(*index).map(|field| col(field.name()))) + .collect::>(); + if !self.compare_key_indices.contains(&self.ts_col_idx) + && let Some(field) = fields.get(self.ts_col_idx) + { + exprs.push(col(field.name())); } exprs } fn necessary_children_exprs(&self, _output_columns: &[usize]) -> Option>> { - if self.unfix.is_some() { - return None; - } - let left_len = self.left.schema().fields().len(); let right_len = self.right.schema().fields().len(); Some(vec![ @@ -223,10 +296,20 @@ impl UserDefinedLogicalNodeCore for UnionDistinctOn { } fn fmt_for_explain(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + let fields = self.left.schema().fields(); + let display_column = |index: usize| match fields.get(index) { + Some(field) => format!("{}@{index}", field.name()), + None => format!("@{index}"), + }; + let compare_keys = self + .compare_key_indices + .iter() + .map(|index| display_column(*index)) + .collect::>(); write!( f, - "UnionDistinctOn: on col=[{:?}], ts_col=[{}]", - self.compare_keys, self.ts_col + "UnionDistinctOn: on col={compare_keys:?}, ts_col={}", + display_column(self.ts_col_idx) ) } @@ -245,38 +328,12 @@ impl UserDefinedLogicalNodeCore for UnionDistinctOn { let left = inputs.next().unwrap(); let right = inputs.next().unwrap(); - if let Some(unfix) = &self.unfix { - let output_schema = left.schema().clone(); - - let compare_keys = unfix - .compare_key_indices - .iter() - .map(|idx| { - resolve_column_name(*idx, &output_schema, "UnionDistinctOn", "compare key") - }) - .collect::>>()?; - - let ts_col = - resolve_column_name(unfix.ts_col_idx, &output_schema, "UnionDistinctOn", "ts")?; - - Ok(Self { - left, - right, - compare_keys, - ts_col, - output_schema, - unfix: None, - }) - } else { - Ok(Self { - left, - right, - compare_keys: self.compare_keys.clone(), - ts_col: self.ts_col.clone(), - output_schema: self.output_schema.clone(), - unfix: None, - }) - } + Self::try_new( + left, + right, + self.compare_key_indices.clone(), + self.ts_col_idx, + ) } } @@ -284,8 +341,8 @@ impl UserDefinedLogicalNodeCore for UnionDistinctOn { pub struct UnionDistinctOnExec { left: Arc, right: Arc, - compare_keys: Vec, - ts_col: String, + compare_key_indices: Vec, + ts_col_idx: usize, output_schema: SchemaRef, metric: ExecutionPlanMetricsSet, properties: Arc, @@ -326,8 +383,8 @@ impl ExecutionPlan for UnionDistinctOnExec { Ok(Arc::new(UnionDistinctOnExec { left, right, - compare_keys: self.compare_keys.clone(), - ts_col: self.ts_col.clone(), + compare_key_indices: self.compare_key_indices.clone(), + ts_col_idx: self.ts_col_idx, output_schema: self.output_schema.clone(), metric: self.metric.clone(), properties: self.properties.clone(), @@ -341,41 +398,23 @@ impl ExecutionPlan for UnionDistinctOnExec { context: Arc, ) -> DataFusionResult { let left_stream = self.left.execute(partition, context.clone())?; - let right_stream = self.right.execute(partition, context.clone())?; - // Convert column name to column index. Add one for the time column. - let mut key_indices = Vec::with_capacity(self.compare_keys.len() + 1); - for key in &self.compare_keys { - let index = self - .output_schema - .column_with_name(key) - .map(|(i, _)| i) - .ok_or_else(|| DataFusionError::Internal(format!("Column {} not found", key)))?; - key_indices.push(index); - } - let ts_index = self - .output_schema - .column_with_name(&self.ts_col) - .map(|(i, _)| i) - .ok_or_else(|| { - DataFusionError::Internal(format!("Column {} not found", self.ts_col)) - })?; - key_indices.push(ts_index); + let mut key_indices = self.compare_key_indices.clone(); + key_indices.push(self.ts_col_idx); - // Build right hash table future. - let hashed_data_future = HashedDataFut::Pending(Box::pin(HashedData::new( - right_stream, - self.random_state.clone(), - key_indices.clone(), - ))); - - let baseline_metric = BaselineMetrics::new(&self.metric, partition); Ok(Box::pin(UnionDistinctOnStream { left: left_stream, - right: hashed_data_future, + right_plan: self.right.clone(), + right_partition: partition, + right_context: context, + right: None, compare_keys: key_indices, output_schema: self.output_schema.clone(), - metric: baseline_metric, + random_state: self.random_state.clone(), + lhs_signatures: HashSet::default(), + hashes: Vec::new(), + phase: StreamPhase::Left, + metric: BaselineMetrics::new(&self.metric, partition), })) } @@ -396,60 +435,139 @@ impl DisplayAs for UnionDistinctOnExec { | DisplayFormatType::TreeRender => { write!( f, - "UnionDistinctOnExec: on col=[{:?}], ts_col=[{}]", - self.compare_keys, self.ts_col + "UnionDistinctOnExec: on col={:?}, ts_col={}", + self.compare_key_indices, self.ts_col_idx ) } } } } -// TODO(ruihang): some unused fields are for metrics, which will be implemented later. -#[allow(dead_code)] pub struct UnionDistinctOnStream { left: SendableRecordBatchStream, - right: HashedDataFut, + right_plan: Arc, + right_partition: usize, + right_context: Arc, + right: Option, /// Include time index compare_keys: Vec, output_schema: SchemaRef, + random_state: RandomState, + lhs_signatures: HashSet, + hashes: Vec, + phase: StreamPhase, metric: BaselineMetrics, } -impl UnionDistinctOnStream { - fn poll_impl(&mut self, cx: &mut Context<'_>) -> Poll::Item>> { - // resolve the right stream - let right = match self.right { - HashedDataFut::Pending(ref mut fut) => { - let right = ready!(fut.as_mut().poll(cx))?; - self.right = HashedDataFut::Ready(right); - let HashedDataFut::Ready(right_ref) = &mut self.right else { - unreachable!() - }; - right_ref - } - HashedDataFut::Ready(ref mut right) => right, - HashedDataFut::Empty => return Poll::Ready(None), - }; +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum StreamPhase { + Left, + Right, + Done, +} - // poll left and probe with right - let next_left = ready!(self.left.poll_next_unpin(cx)); - match next_left { - Some(Ok(left)) => { - // observe left batch and return it - right.update_map(&left)?; - Poll::Ready(Some(Ok(left))) +impl UnionDistinctOnStream { + fn hash_batch(&mut self, batch: &RecordBatch) -> DataFusionResult<()> { + let arrays = self + .compare_keys + .iter() + .map(|index| batch.column(*index).clone()) + .collect::>(); + self.hashes.clear(); + self.hashes.resize(batch.num_rows(), 0); + hash_utils::create_hashes(&arrays, &self.random_state, &mut self.hashes)?; + Ok(()) + } + + fn filter_rhs_batch(&mut self, batch: RecordBatch) -> DataFusionResult> { + self.hash_batch(&batch)?; + let mut survivor_indices: Option> = None; + for (index, hash) in self.hashes.iter().enumerate() { + if self.lhs_signatures.contains(hash) { + survivor_indices.get_or_insert_with(|| (0..index).collect()); + } else if let Some(indices) = &mut survivor_indices { + indices.push(index); } - Some(Err(e)) => Poll::Ready(Some(Err(e))), + } + + match survivor_indices { + None => Ok(Some(with_schema(batch, self.output_schema.clone())?)), + Some(indices) if indices.is_empty() => Ok(None), + Some(indices) => Ok(Some(with_schema( + take_batch(&batch, &indices)?, + self.output_schema.clone(), + )?)), + } + } + + fn terminal_error(&mut self, error: DataFusionError) -> Poll::Item>> { + self.phase = StreamPhase::Done; + Poll::Ready(Some(Err(error))) + } + + fn poll_right(&mut self, cx: &mut Context<'_>) -> Poll::Item>> { + if self.right.is_none() { + let right = match self + .right_plan + .execute(self.right_partition, self.right_context.clone()) + { + Ok(right) => right, + Err(error) => return self.terminal_error(error), + }; + self.right = Some(right); + } + + let right = self.right.as_mut().expect("right stream is initialized"); + match ready!(right.poll_next_unpin(cx)) { + Some(Ok(batch)) if self.lhs_signatures.is_empty() => { + match with_schema(batch, self.output_schema.clone()) { + Ok(batch) => Poll::Ready(Some(Ok(batch))), + Err(error) => self.terminal_error(error), + } + } + Some(Ok(batch)) => match self.filter_rhs_batch(batch) { + Ok(Some(batch)) => Poll::Ready(Some(Ok(batch))), + Ok(None) => { + // One fully filtered input batch has been consumed. Yield so a sequence + // of filtered batches cannot monopolize a single downstream poll. + cx.waker().wake_by_ref(); + Poll::Pending + } + Err(error) => self.terminal_error(error), + }, + Some(Err(error)) => self.terminal_error(error), None => { - // left stream is exhausted, so we can send the right part - let right = std::mem::replace(&mut self.right, HashedDataFut::Empty); - let HashedDataFut::Ready(data) = right else { - unreachable!() - }; - Poll::Ready(Some(data.finish())) + self.phase = StreamPhase::Done; + Poll::Ready(None) } } } + + fn poll_impl(&mut self, cx: &mut Context<'_>) -> Poll::Item>> { + match self.phase { + StreamPhase::Left => match ready!(self.left.poll_next_unpin(cx)) { + Some(Ok(batch)) => { + if batch.num_rows() > 0 { + if let Err(error) = self.hash_batch(&batch) { + return self.terminal_error(error); + } + self.lhs_signatures.extend(self.hashes.iter().copied()); + } + match with_schema(batch, self.output_schema.clone()) { + Ok(batch) => Poll::Ready(Some(Ok(batch))), + Err(error) => self.terminal_error(error), + } + } + Some(Err(error)) => self.terminal_error(error), + None => { + self.phase = StreamPhase::Right; + self.poll_right(cx) + } + }, + StreamPhase::Right => self.poll_right(cx), + StreamPhase::Done => Poll::Ready(None), + } + } } impl RecordBatchStream for UnionDistinctOnStream { @@ -462,185 +580,49 @@ impl Stream for UnionDistinctOnStream { type Item = DataFusionResult; fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { - self.poll_impl(cx) + let poll = self.poll_impl(cx); + self.metric.record_poll(poll) } } -/// Simple future state for [HashedData] -enum HashedDataFut { - /// The result is not ready - Pending(BoxFuture<'static, DataFusionResult>), - /// The result is ready - Ready(HashedData), - /// The result is taken - Empty, -} - -/// ALL input batches and its hash table -struct HashedData { - // TODO(ruihang): use `JoinHashMap` instead after upgrading to DF 34.0 - /// Hash table for all input batches. The key is hash value, and the value - /// is the index of `bathc`. - hash_map: HashMap, - /// Output batch. - batch: RecordBatch, - /// The indices of the columns to be hashed. - hash_key_indices: Vec, - random_state: RandomState, -} - -impl HashedData { - pub async fn new( - input: SendableRecordBatchStream, - random_state: RandomState, - hash_key_indices: Vec, - ) -> DataFusionResult { - // Collect all batches from the input stream - let initial = (Vec::new(), 0); - let schema = input.schema(); - let (batches, _num_rows) = input - .try_fold(initial, |mut acc, batch| async { - // Update rowcount - acc.1 += batch.num_rows(); - // Push batch to output - acc.0.push(batch); - Ok(acc) - }) - .await?; - - // Create hash for each batch - let mut hash_map = HashMap::default(); - let mut hashes_buffer = Vec::new(); - let mut interleave_indices = Vec::new(); - for (batch_number, batch) in batches.iter().enumerate() { - hashes_buffer.resize(batch.num_rows(), 0); - // get columns for hashing - let arrays = hash_key_indices - .iter() - .map(|i| batch.column(*i).clone()) - .collect::>(); - - // compute hash - let hash_values = - hash_utils::create_hashes(&arrays, &random_state, &mut hashes_buffer)?; - for (row_number, hash_value) in hash_values.iter().enumerate() { - // Only keeps the first observed row for each hash value - if hash_map - .try_insert(*hash_value, interleave_indices.len()) - .is_ok() - { - interleave_indices.push((batch_number, row_number)); - } - } - } - - // Finalize the hash map - let batch = interleave_batches(schema, batches, interleave_indices)?; - - Ok(Self { - hash_map, - batch, - hash_key_indices, - random_state, - }) - } - - /// Remove rows that hash value present in the input - /// record batch from the hash map. - pub fn update_map(&mut self, input: &RecordBatch) -> DataFusionResult<()> { - // get columns for hashing - let mut hashes_buffer = Vec::new(); - let arrays = self - .hash_key_indices - .iter() - .map(|i| input.column(*i).clone()) - .collect::>(); - - // compute hash - hashes_buffer.resize(input.num_rows(), 0); - let hash_values = - hash_utils::create_hashes(&arrays, &self.random_state, &mut hashes_buffer)?; - - // remove those hashes - for hash in hash_values { - self.hash_map.remove(hash); - } - - Ok(()) - } - - pub fn finish(self) -> DataFusionResult { - let valid_indices = self.hash_map.values().copied().collect::>(); - let result = take_batch(&self.batch, &valid_indices)?; - Ok(result) - } -} - -/// Utility function to interleave batches. Based on [interleave](datafusion::arrow::compute::interleave) -fn interleave_batches( - schema: SchemaRef, - batches: Vec, - indices: Vec<(usize, usize)>, -) -> DataFusionResult { - if batches.is_empty() { - if indices.is_empty() { - return Ok(RecordBatch::new_empty(schema)); - } else { - return Err(DataFusionError::Internal( - "Cannot interleave empty batches with non-empty indices".to_string(), - )); - } - } - - // transform batches into arrays - let mut arrays = vec![vec![]; schema.fields().len()]; - for batch in &batches { - for (i, array) in batch.columns().iter().enumerate() { - arrays[i].push(array.as_ref()); - } - } - - // interleave arrays - let interleaved_arrays: Vec<_> = arrays - .into_iter() - .map(|array| compute::interleave(&array, &indices)) - .collect::>()?; - - // assemble new record batch - RecordBatch::try_new(schema, interleaved_arrays) +fn with_schema(batch: RecordBatch, schema: SchemaRef) -> DataFusionResult { + RecordBatch::try_new(schema, batch.columns().to_vec()) .map_err(|e| DataFusionError::ArrowError(Box::new(e), None)) } -/// Utility function to take rows from a record batch. Based on [take](datafusion::arrow::compute::take) +/// Utility function to take ordered rows from a record batch. fn take_batch(batch: &RecordBatch, indices: &[usize]) -> DataFusionResult { - // fast path if batch.num_rows() == indices.len() { return Ok(batch.clone()); } - let schema = batch.schema(); - - let indices_array = UInt64Array::from_iter(indices.iter().map(|i| *i as u64)); + let indices_array = UInt64Array::from_iter(indices.iter().map(|index| *index as u64)); let arrays = batch .columns() .iter() .map(|array| compute::take(array, &indices_array, None)) .collect::, _>>() - .map_err(|e| DataFusionError::ArrowError(Box::new(e), None))?; + .map_err(|error| DataFusionError::ArrowError(Box::new(error), None))?; - let result = RecordBatch::try_new(schema, arrays) - .map_err(|e| DataFusionError::ArrowError(Box::new(e), None))?; - Ok(result) + RecordBatch::try_new(batch.schema(), arrays) + .map_err(|error| DataFusionError::ArrowError(Box::new(error), None)) } #[cfg(test)] mod test { + use std::collections::VecDeque; use std::sync::Arc; + use std::sync::atomic::{AtomicUsize, Ordering}; + use std::task::{Wake, Waker}; - use datafusion::arrow::array::Int32Array; - use datafusion::arrow::datatypes::{DataType, Field, Schema}; + use datafusion::arrow::array::{Array, Float64Array, Int32Array, Int64Array, StringArray}; + use datafusion::arrow::datatypes::{DataType, Field, Schema, SchemaRef}; use datafusion::common::ToDFSchema; + use datafusion::datasource::memory::MemorySourceConfig; + use datafusion::datasource::source::DataSourceExec; use datafusion::logical_expr::{EmptyRelation, LogicalPlan}; + use datafusion::prelude::SessionContext; + use futures::StreamExt; use super::*; @@ -660,13 +642,7 @@ mod test { produce_one_row: false, schema: df_schema.clone(), }); - let plan = UnionDistinctOn::new( - left, - right, - vec!["k".to_string()], - "ts".to_string(), - df_schema, - ); + let plan = UnionDistinctOn::try_new(left, right, vec![1], 0).unwrap(); // Simulate a parent projection requesting only one output column. let output_columns = [2usize]; @@ -676,56 +652,6 @@ mod test { assert_eq!(required[1].as_slice(), &[0, 1, 2]); } - #[test] - fn test_interleave_batches() { - let schema = Schema::new(vec![ - Field::new("a", DataType::Int32, false), - Field::new("b", DataType::Int32, false), - ]); - - let batch1 = RecordBatch::try_new( - Arc::new(schema.clone()), - vec![ - Arc::new(Int32Array::from(vec![1, 2, 3])), - Arc::new(Int32Array::from(vec![4, 5, 6])), - ], - ) - .unwrap(); - - let batch2 = RecordBatch::try_new( - Arc::new(schema.clone()), - vec![ - Arc::new(Int32Array::from(vec![7, 8, 9])), - Arc::new(Int32Array::from(vec![10, 11, 12])), - ], - ) - .unwrap(); - - let batch3 = RecordBatch::try_new( - Arc::new(schema.clone()), - vec![ - Arc::new(Int32Array::from(vec![13, 14, 15])), - Arc::new(Int32Array::from(vec![16, 17, 18])), - ], - ) - .unwrap(); - - let batches = vec![batch1, batch2, batch3]; - let indices = vec![(0, 0), (1, 0), (2, 0), (0, 1), (1, 1), (2, 1)]; - let result = interleave_batches(Arc::new(schema.clone()), batches, indices).unwrap(); - - let expected = RecordBatch::try_new( - Arc::new(schema), - vec![ - Arc::new(Int32Array::from(vec![1, 7, 13, 2, 8, 14])), - Arc::new(Int32Array::from(vec![4, 10, 16, 5, 11, 17])), - ], - ) - .unwrap(); - - assert_eq!(result, expected); - } - #[test] fn test_take_batch() { let schema = Schema::new(vec![ @@ -757,39 +683,958 @@ mod test { assert_eq!(result, expected); } + fn empty_plan(schema: datafusion::common::DFSchemaRef) -> LogicalPlan { + LogicalPlan::EmptyRelation(EmptyRelation { + produce_one_row: false, + schema, + }) + } + + fn schemas(left: SchemaRef, right: SchemaRef) -> (LogicalPlan, LogicalPlan) { + ( + empty_plan(left.to_dfschema_ref().unwrap()), + empty_plan(right.to_dfschema_ref().unwrap()), + ) + } + + fn schema(prefix: &str, key: &str, nullable: bool) -> SchemaRef { + Arc::new(Schema::new(vec![ + Field::new(format!("{prefix}_ts"), DataType::Int64, false), + Field::new(format!("{prefix}_{key}"), DataType::Utf8, nullable), + Field::new(format!("{prefix}_value"), DataType::Float64, nullable), + ])) + } + + fn batch(schema: SchemaRef, ts: i64, key: Option<&str>, value: Option) -> RecordBatch { + RecordBatch::try_new( + schema, + vec![ + Arc::new(Int64Array::from(vec![ts])), + Arc::new(StringArray::from(vec![key])), + Arc::new(Float64Array::from(vec![value])), + ], + ) + .unwrap() + } + + fn source_exec(batch: RecordBatch) -> Arc { + source_exec_batches(batch.schema(), vec![batch]) + } + + async fn execute( + plan: &UnionDistinctOn, + left: RecordBatch, + right: RecordBatch, + ) -> Vec { + datafusion::physical_plan::collect( + plan.to_execution_plan(source_exec(left), source_exec(right)), + SessionContext::default().task_ctx(), + ) + .await + .unwrap() + } + + fn source_exec_batches(schema: SchemaRef, batches: Vec) -> Arc { + Arc::new(DataSourceExec::new(Arc::new( + MemorySourceConfig::try_new(&[batches], schema, None).unwrap(), + ))) + } + + async fn execute_batches( + plan: &UnionDistinctOn, + left_schema: SchemaRef, + left: Vec, + right_schema: SchemaRef, + right: Vec, + ) -> Vec { + datafusion::physical_plan::collect( + plan.to_execution_plan( + source_exec_batches(left_schema, left), + source_exec_batches(right_schema, right), + ), + SessionContext::default().task_ctx(), + ) + .await + .unwrap() + } + + fn simple_schema(prefix: &str, nullable: bool) -> SchemaRef { + schema(prefix, "label", nullable) + } + + fn simple_batch(schema: SchemaRef, rows: &[(i64, &str, f64)]) -> RecordBatch { + RecordBatch::try_new( + schema, + vec![ + Arc::new(Int64Array::from_iter_values(rows.iter().map(|row| row.0))), + Arc::new(StringArray::from_iter_values(rows.iter().map(|row| row.1))), + Arc::new(Float64Array::from_iter_values(rows.iter().map(|row| row.2))), + ], + ) + .unwrap() + } + + fn simple_rows(batches: &[RecordBatch]) -> Vec<(i64, String, f64)> { + batches + .iter() + .flat_map(|batch| { + let timestamps = batch + .column(0) + .as_any() + .downcast_ref::() + .unwrap(); + let labels = batch + .column(1) + .as_any() + .downcast_ref::() + .unwrap(); + let values = batch + .column(2) + .as_any() + .downcast_ref::() + .unwrap(); + (0..batch.num_rows()).map(move |row| { + ( + timestamps.value(row), + labels.value(row).to_string(), + values.value(row), + ) + }) + }) + .collect() + } + + fn simple_plan(left_schema: SchemaRef, right_schema: SchemaRef) -> UnionDistinctOn { + let (left, right) = schemas(left_schema, right_schema); + UnionDistinctOn::try_new(left, right, vec![1], 0).unwrap() + } + + #[derive(Clone, Debug)] + enum TestEvent { + Batch(RecordBatch), + Error(&'static str), + } + + struct TestStream { + schema: SchemaRef, + events: VecDeque, + polls: Arc, + } + + impl Stream for TestStream { + type Item = DataFusionResult; + + fn poll_next(mut self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll> { + self.polls.fetch_add(1, Ordering::SeqCst); + match self.events.pop_front() { + Some(TestEvent::Batch(batch)) => Poll::Ready(Some(Ok(batch))), + Some(TestEvent::Error(message)) => { + Poll::Ready(Some(Err(DataFusionError::Internal(message.to_string())))) + } + None => Poll::Ready(None), + } + } + } + + impl RecordBatchStream for TestStream { + fn schema(&self) -> SchemaRef { + self.schema.clone() + } + } + + struct CountingWake(AtomicUsize); + + impl Wake for CountingWake { + fn wake(self: Arc) { + self.0.fetch_add(1, Ordering::SeqCst); + } + + fn wake_by_ref(self: &Arc) { + self.0.fetch_add(1, Ordering::SeqCst); + } + } + + fn test_stream( + schema: SchemaRef, + events: Vec, + polls: Arc, + ) -> SendableRecordBatchStream { + Box::pin(TestStream { + schema, + events: events.into(), + polls, + }) + } + + #[derive(Debug)] + struct TestExec { + schema: SchemaRef, + events: Vec, + polls: Arc, + executions: Arc, + execute_error: Option<&'static str>, + properties: Arc, + } + + impl TestExec { + fn new( + schema: SchemaRef, + events: Vec, + polls: Arc, + executions: Arc, + ) -> Self { + Self { + properties: Arc::new(PlanProperties::new( + EquivalenceProperties::new(schema.clone()), + Partitioning::UnknownPartitioning(1), + EmissionType::Incremental, + Boundedness::Bounded, + )), + schema, + events, + polls, + executions, + execute_error: None, + } + } + + fn with_execute_error(mut self, error: &'static str) -> Self { + self.execute_error = Some(error); + self + } + } + + impl DisplayAs for TestExec { + fn fmt_as(&self, _t: DisplayFormatType, _f: &mut std::fmt::Formatter) -> std::fmt::Result { + Ok(()) + } + } + + impl ExecutionPlan for TestExec { + fn name(&self) -> &str { + "TestExec" + } + + fn as_any(&self) -> &dyn Any { + self + } + + fn properties(&self) -> &Arc { + &self.properties + } + + fn children(&self) -> Vec<&Arc> { + vec![] + } + + fn with_new_children( + self: Arc, + _children: Vec>, + ) -> DataFusionResult> { + Ok(self) + } + + fn execute( + &self, + _partition: usize, + _context: Arc, + ) -> DataFusionResult { + self.executions.fetch_add(1, Ordering::SeqCst); + if let Some(error) = self.execute_error { + return Err(DataFusionError::Internal(error.to_string())); + } + Ok(test_stream( + self.schema.clone(), + self.events.clone(), + self.polls.clone(), + )) + } + } + + fn test_union_stream( + left: SendableRecordBatchStream, + right: Arc, + output_schema: SchemaRef, + ) -> UnionDistinctOnStream { + let metrics = ExecutionPlanMetricsSet::new(); + UnionDistinctOnStream { + left, + right_plan: right, + right_partition: 0, + right_context: SessionContext::default().task_ctx(), + right: None, + compare_keys: vec![1, 0], + output_schema, + random_state: RandomState::new(), + lhs_signatures: HashSet::default(), + hashes: Vec::new(), + phase: StreamPhase::Left, + metric: BaselineMetrics::new(&metrics, 0), + } + } + + fn series_rows(batches: &[RecordBatch]) -> Vec<(i64, &str, f64, &str)> { + batches + .iter() + .flat_map(|batch| { + let timestamps = batch + .column(0) + .as_any() + .downcast_ref::() + .unwrap(); + let series = batch + .column(1) + .as_any() + .downcast_ref::() + .unwrap(); + let values = batch + .column(2) + .as_any() + .downcast_ref::() + .unwrap(); + let keys = batch + .column(3) + .as_any() + .downcast_ref::() + .unwrap(); + (0..batch.num_rows()).map(move |row| { + ( + timestamps.value(row), + series.value(row), + values.value(row), + keys.value(row), + ) + }) + }) + .collect() + } + #[tokio::test] - async fn encode_decode_union_distinct_on() { + async fn serialize_deserialize_and_execute_with_different_input_names() { + let (left_schema, right_schema) = + (schema("left", "job", false), schema("right", "job", false)); + let (left_plan, right_plan) = schemas(left_schema.clone(), right_schema.clone()); + let decoded = UnionDistinctOn::deserialize( + &UnionDistinctOn::try_new(left_plan.clone(), right_plan.clone(), vec![1], 0) + .unwrap() + .serialize(), + ) + .unwrap() + .with_exprs_and_inputs(vec![], vec![left_plan, right_plan]) + .unwrap(); + assert_eq!( + (decoded.compare_key_indices.as_slice(), decoded.ts_col_idx), + (&[1usize][..], 0) + ); + assert_eq!( + decoded.output_schema, + left_schema.clone().to_dfschema_ref().unwrap() + ); + let result = execute( + &decoded, + batch(left_schema.clone(), 1, Some("left"), Some(10.0)), + batch(right_schema, 2, Some("right"), Some(20.0)), + ) + .await; + assert!(result.iter().all(|batch| batch.schema() == left_schema)); + } + + #[tokio::test] + async fn execute_widens_nullable_fields_and_emits_rhs_nulls() { + let (left_schema, right_schema) = ( + schema("left", "label", false), + schema("right", "label", true), + ); + let (left_plan, right_plan) = schemas(left_schema.clone(), right_schema.clone()); + let plan = UnionDistinctOn::try_new(left_plan, right_plan, vec![1, 2], 0).unwrap(); + let declared_schema = plan.output_schema.inner().clone(); + assert!( + declared_schema.fields()[1..] + .iter() + .all(|field| field.is_nullable()) + ); + let result = execute( + &plan, + batch(left_schema, 1, Some("present"), Some(10.0)), + batch(right_schema, 2, None, None), + ) + .await; + assert_eq!(result.len(), 2); + assert!(result.iter().all(|batch| batch.schema() == declared_schema)); + assert!(result[1].column(1).is_null(0)); + assert!(result[1].column(2).is_null(0)); + } + + #[tokio::test] + async fn empty_lhs_preserves_distinct_rhs_series_with_same_normalized_key() { + let left_schema = Arc::new(Schema::new(vec![ + Field::new("ts", DataType::Int64, false), + Field::new("series", DataType::Utf8, false), + Field::new("value", DataType::Float64, false), + Field::new("__normalized_absent_label", DataType::Utf8, false), + ])); + let right_schema = Arc::new(Schema::new(vec![ + Field::new("rhs_ts", DataType::Int64, false), + Field::new("rhs_series", DataType::Utf8, false), + Field::new("rhs_value", DataType::Float64, false), + Field::new("rhs_normalized_absent_label", DataType::Utf8, false), + ])); + let (left, right) = schemas(left_schema.clone(), right_schema.clone()); + let plan = UnionDistinctOn::try_new(left, right, vec![3], 0).unwrap(); + let declared_schema = plan.output_schema.inner().clone(); + let rhs = RecordBatch::try_new( + right_schema, + vec![ + Arc::new(Int64Array::from(vec![1_000, 1_000])), + Arc::new(StringArray::from(vec!["series_a", "series_b"])), + Arc::new(Float64Array::from(vec![10.0, 20.0])), + Arc::new(StringArray::from(vec!["", ""])), + ], + ) + .unwrap(); + + let result = execute(&plan, RecordBatch::new_empty(left_schema), rhs).await; + assert!(result.iter().all(|batch| batch.schema() == declared_schema)); + assert_eq!(result.iter().map(RecordBatch::num_rows).sum::(), 2); + + let mut rows = series_rows(&result); + rows.sort_by_key(|row| row.1); + assert_eq!( + rows, + vec![(1_000, "series_a", 10.0, ""), (1_000, "series_b", 20.0, "")] + ); + } + + #[tokio::test] + async fn lhs_signature_suppresses_all_matching_rhs_rows_and_preserves_lhs() { + let left_schema = Arc::new(Schema::new(vec![ + Field::new("ts", DataType::Int64, false), + Field::new("series", DataType::Utf8, false), + Field::new("value", DataType::Float64, false), + Field::new("__normalized_absent_label", DataType::Utf8, false), + ])); + let right_schema = Arc::new(Schema::new(vec![ + Field::new("rhs_ts", DataType::Int64, false), + Field::new("rhs_series", DataType::Utf8, false), + Field::new("rhs_value", DataType::Float64, false), + Field::new("rhs_normalized_absent_label", DataType::Utf8, false), + ])); + let (left, right) = schemas(left_schema.clone(), right_schema.clone()); + let plan = UnionDistinctOn::try_new(left, right, vec![3], 0).unwrap(); + let declared_schema = plan.output_schema.inner().clone(); + let lhs = RecordBatch::try_new( + left_schema, + vec![ + Arc::new(Int64Array::from(vec![1_000, 2_000])), + Arc::new(StringArray::from(vec!["lhs_match", "lhs_only"])), + Arc::new(Float64Array::from(vec![1.0, 2.0])), + Arc::new(StringArray::from(vec!["", "left-only"])), + ], + ) + .unwrap(); + let rhs = RecordBatch::try_new( + right_schema, + vec![ + Arc::new(Int64Array::from(vec![1_000, 1_000, 3_000])), + Arc::new(StringArray::from(vec![ + "rhs_match_a", + "rhs_match_b", + "rhs_unmatched", + ])), + Arc::new(Float64Array::from(vec![10.0, 20.0, 30.0])), + Arc::new(StringArray::from(vec!["", "", ""])), + ], + ) + .unwrap(); + + let result = execute(&plan, lhs, rhs).await; + assert_eq!(result.iter().map(RecordBatch::num_rows).sum::(), 3); + assert!(result.iter().all(|batch| batch.schema() == declared_schema)); + let mut rows = series_rows(&result); + rows.sort_by_key(|row| row.1); + assert_eq!( + rows, + vec![ + (1_000, "lhs_match", 1.0, ""), + (2_000, "lhs_only", 2.0, "left-only"), + (3_000, "rhs_unmatched", 30.0, ""), + ] + ); + } + + #[tokio::test] + async fn empty_lhs_preserves_duplicate_rhs_rows_across_batches() { + let left_schema = simple_schema("left", false); + let right_schema = simple_schema("right", false); + let plan = simple_plan(left_schema.clone(), right_schema.clone()); + let result = execute_batches( + &plan, + left_schema, + vec![], + right_schema.clone(), + vec![ + simple_batch( + right_schema.clone(), + &[(1, "same", 10.0), (1, "same", 11.0)], + ), + simple_batch(right_schema, &[(1, "same", 12.0)]), + ], + ) + .await; + assert_eq!( + simple_rows(&result), + vec![ + (1, "same".to_string(), 10.0), + (1, "same".to_string(), 11.0), + (1, "same".to_string(), 12.0), + ] + ); + } + + #[tokio::test] + async fn lhs_signature_suppresses_rhs_duplicates_across_batches() { + let left_schema = simple_schema("left", false); + let right_schema = simple_schema("right", false); + let plan = simple_plan(left_schema.clone(), right_schema.clone()); + let result = execute_batches( + &plan, + left_schema.clone(), + vec![simple_batch(left_schema, &[(1, "match", 1.0)])], + right_schema.clone(), + vec![ + simple_batch(right_schema.clone(), &[(1, "match", 10.0)]), + simple_batch(right_schema, &[(1, "match", 11.0)]), + ], + ) + .await; + assert_eq!(simple_rows(&result), vec![(1, "match".to_string(), 1.0)]); + } + + #[tokio::test] + async fn hash_scratch_reset_suppresses_null_label_after_same_sized_lhs_batches() { + let left_schema = simple_schema("left", true); + let right_schema = simple_schema("right", true); + let plan = simple_plan(left_schema.clone(), right_schema.clone()); + let lhs_poison = RecordBatch::try_new( + left_schema.clone(), + vec![ + Arc::new(Int64Array::from(vec![1])), + Arc::new(StringArray::from(vec![Some("poison")])), + Arc::new(Float64Array::from(vec![1.0])), + ], + ) + .unwrap(); + let lhs_null = RecordBatch::try_new( + left_schema.clone(), + vec![ + Arc::new(Int64Array::from(vec![42])), + Arc::new(StringArray::from(vec![None::<&str>])), + Arc::new(Float64Array::from(vec![2.0])), + ], + ) + .unwrap(); + let rhs_null = RecordBatch::try_new( + right_schema.clone(), + vec![ + Arc::new(Int64Array::from(vec![42])), + Arc::new(StringArray::from(vec![None::<&str>])), + Arc::new(Float64Array::from(vec![999.0])), + ], + ) + .unwrap(); + + let result = execute_batches( + &plan, + left_schema, + vec![lhs_poison, lhs_null], + right_schema, + vec![rhs_null], + ) + .await; + assert_eq!(result.iter().map(RecordBatch::num_rows).sum::(), 2); + let values = result + .iter() + .flat_map(|batch| { + batch + .column(2) + .as_any() + .downcast_ref::() + .unwrap() + .values() + .iter() + .copied() + }) + .collect::>(); + assert_eq!(values, vec![1.0, 2.0]); + } + + #[tokio::test] + async fn mixed_rhs_duplicates_retain_unmatched_rows_in_order() { + let left_schema = simple_schema("left", false); + let right_schema = simple_schema("right", false); + let plan = simple_plan(left_schema.clone(), right_schema.clone()); + let result = execute_batches( + &plan, + left_schema.clone(), + vec![simple_batch(left_schema, &[(1, "match", 1.0)])], + right_schema.clone(), + vec![simple_batch( + right_schema, + &[ + (1, "match", 10.0), + (1, "keep", 11.0), + (1, "keep", 12.0), + (1, "match", 13.0), + (1, "later", 14.0), + ], + )], + ) + .await; + assert_eq!( + simple_rows(&result), + vec![ + (1, "match".to_string(), 1.0), + (1, "keep".to_string(), 11.0), + (1, "keep".to_string(), 12.0), + (1, "later".to_string(), 14.0), + ] + ); + } + + #[tokio::test] + async fn same_label_at_different_timestamps_is_unmatched() { + let left_schema = simple_schema("left", false); + let right_schema = simple_schema("right", false); + let plan = simple_plan(left_schema.clone(), right_schema.clone()); + let result = execute_batches( + &plan, + left_schema.clone(), + vec![simple_batch(left_schema, &[(1, "label", 1.0)])], + right_schema.clone(), + vec![simple_batch(right_schema, &[(2, "label", 2.0)])], + ) + .await; + assert_eq!( + simple_rows(&result), + vec![(1, "label".to_string(), 1.0), (2, "label".to_string(), 2.0)] + ); + } + + #[tokio::test] + async fn empty_input_combinations_complete_cleanly() { + let left_schema = simple_schema("left", false); + let right_schema = simple_schema("right", false); + let plan = simple_plan(left_schema.clone(), right_schema.clone()); + + let empty_left = execute_batches( + &plan, + left_schema.clone(), + vec![], + right_schema.clone(), + vec![simple_batch(right_schema.clone(), &[(1, "rhs", 1.0)])], + ) + .await; + assert_eq!(simple_rows(&empty_left), vec![(1, "rhs".to_string(), 1.0)]); + + let empty_right = execute_batches( + &plan, + left_schema.clone(), + vec![simple_batch(left_schema.clone(), &[(1, "lhs", 1.0)])], + right_schema.clone(), + vec![], + ) + .await; + assert_eq!(simple_rows(&empty_right), vec![(1, "lhs".to_string(), 1.0)]); + + let both_empty = execute_batches(&plan, left_schema, vec![], right_schema, vec![]).await; + assert!(both_empty.is_empty()); + } + + #[tokio::test] + async fn zero_row_batches_on_either_side_are_handled() { + let left_schema = simple_schema("left", false); + let right_schema = simple_schema("right", false); + let plan = simple_plan(left_schema.clone(), right_schema.clone()); + + let zero_row_left = execute_batches( + &plan, + left_schema.clone(), + vec![simple_batch(left_schema.clone(), &[])], + right_schema.clone(), + vec![simple_batch(right_schema.clone(), &[(1, "rhs", 1.0)])], + ) + .await; + assert_eq!( + simple_rows(&zero_row_left), + vec![(1, "rhs".to_string(), 1.0)] + ); + + let zero_row_right = execute_batches( + &plan, + left_schema.clone(), + vec![simple_batch(left_schema, &[(1, "lhs", 1.0)])], + right_schema.clone(), + vec![simple_batch(right_schema, &[])], + ) + .await; + assert_eq!( + simple_rows(&zero_row_right), + vec![(1, "lhs".to_string(), 1.0)] + ); + } + + #[tokio::test] + async fn metrics_record_output_rows_and_completion() { + let left_schema = simple_schema("left", false); + let right_schema = simple_schema("right", false); + let plan = simple_plan(left_schema.clone(), right_schema.clone()); + let exec = plan.to_execution_plan( + source_exec(simple_batch(left_schema, &[(1, "lhs", 1.0)])), + source_exec(simple_batch(right_schema, &[(2, "rhs", 2.0)])), + ); + + let output = + datafusion::physical_plan::collect(exec.clone(), SessionContext::default().task_ctx()) + .await + .unwrap(); + assert_eq!(simple_rows(&output).len(), 2); + + let metrics = exec.metrics().unwrap(); + assert_eq!(metrics.output_rows(), Some(2)); + assert!(metrics.iter().any(|metric| { + metric.value().name() == "end_timestamp" && metric.value().as_usize() > 0 + })); + } + + #[tokio::test] + async fn output_keeps_lhs_before_rhs_and_input_order() { + let left_schema = simple_schema("left", false); + let right_schema = simple_schema("right", false); + let plan = simple_plan(left_schema.clone(), right_schema.clone()); + let result = execute_batches( + &plan, + left_schema.clone(), + vec![ + simple_batch(left_schema.clone(), &[(1, "lhs-a", 1.0)]), + simple_batch(left_schema, &[(2, "lhs-b", 2.0)]), + ], + right_schema.clone(), + vec![ + simple_batch(right_schema.clone(), &[(3, "rhs-a", 3.0)]), + simple_batch(right_schema, &[(4, "rhs-b", 4.0)]), + ], + ) + .await; + assert_eq!( + simple_rows(&result), + vec![ + (1, "lhs-a".to_string(), 1.0), + (2, "lhs-b".to_string(), 2.0), + (3, "rhs-a".to_string(), 3.0), + (4, "rhs-b".to_string(), 4.0), + ] + ); + } + + #[test] + fn fully_filtered_rhs_batch_wakes_before_later_unmatched_batch() { + let schema = simple_schema("stream", false); + let left_polls = Arc::new(AtomicUsize::new(0)); + let right_polls = Arc::new(AtomicUsize::new(0)); + let right_executions = Arc::new(AtomicUsize::new(0)); + let left = test_stream( + schema.clone(), + vec![TestEvent::Batch(simple_batch( + schema.clone(), + &[(1, "match", 1.0)], + ))], + left_polls, + ); + let right = Arc::new(TestExec::new( + schema.clone(), + vec![ + TestEvent::Batch(simple_batch(schema.clone(), &[(1, "match", 2.0)])), + TestEvent::Batch(simple_batch(schema.clone(), &[(2, "keep", 3.0)])), + ], + right_polls.clone(), + right_executions.clone(), + )); + let mut stream = Box::pin(test_union_stream(left, right, schema)); + let wake = Arc::new(CountingWake(AtomicUsize::new(0))); + let waker = Waker::from(wake.clone()); + let mut context = Context::from_waker(&waker); + + assert!(matches!( + stream.as_mut().poll_next(&mut context), + Poll::Ready(Some(Ok(_))) + )); + assert!(matches!( + stream.as_mut().poll_next(&mut context), + Poll::Pending + )); + assert_eq!(right_executions.load(Ordering::SeqCst), 1); + assert_eq!(right_polls.load(Ordering::SeqCst), 1); + assert_eq!(wake.0.load(Ordering::SeqCst), 1); + assert!(matches!( + stream.as_mut().poll_next(&mut context), + Poll::Ready(Some(Ok(batch))) + if simple_rows(std::slice::from_ref(&batch)) + == vec![(2, "keep".to_string(), 3.0)] + )); + } + + #[tokio::test] + async fn rhs_is_not_polled_before_lhs_eof_and_drop_preserves_backpressure() { + let schema = simple_schema("stream", false); + let left_polls = Arc::new(AtomicUsize::new(0)); + let right_polls = Arc::new(AtomicUsize::new(0)); + let left = test_stream( + schema.clone(), + vec![TestEvent::Batch(simple_batch( + schema.clone(), + &[(1, "lhs", 1.0)], + ))], + left_polls.clone(), + ); + let right_executions = Arc::new(AtomicUsize::new(0)); + let right = Arc::new(TestExec::new( + schema.clone(), + vec![TestEvent::Batch(simple_batch( + schema.clone(), + &[(2, "rhs", 2.0)], + ))], + right_polls.clone(), + right_executions.clone(), + )); + let mut stream = Box::pin(test_union_stream(left, right, schema)); + + let first = stream.next().await.unwrap().unwrap(); + assert_eq!(simple_rows(&[first]), vec![(1, "lhs".to_string(), 1.0)]); + assert_eq!(left_polls.load(Ordering::SeqCst), 1); + assert_eq!(right_polls.load(Ordering::SeqCst), 0); + assert_eq!(right_executions.load(Ordering::SeqCst), 0); + drop(stream); + assert_eq!(right_polls.load(Ordering::SeqCst), 0); + assert_eq!(right_executions.load(Ordering::SeqCst), 0); + } + + #[tokio::test] + async fn lhs_and_delayed_rhs_errors_propagate() { + let schema = simple_schema("stream", false); + let polls = Arc::new(AtomicUsize::new(0)); + let mut left_error = Box::pin(test_union_stream( + test_stream( + schema.clone(), + vec![TestEvent::Error("left failed")], + polls.clone(), + ), + Arc::new(TestExec::new( + schema.clone(), + vec![], + polls.clone(), + Arc::new(AtomicUsize::new(0)), + )), + schema.clone(), + )); + assert!(left_error.next().await.unwrap().is_err()); + assert!(left_error.next().await.is_none()); + + let mut right_error = Box::pin(test_union_stream( + test_stream( + schema.clone(), + vec![TestEvent::Batch(simple_batch( + schema.clone(), + &[(1, "lhs", 1.0)], + ))], + polls.clone(), + ), + Arc::new(TestExec::new( + schema.clone(), + vec![TestEvent::Error("right failed")], + polls, + Arc::new(AtomicUsize::new(0)), + )), + schema.clone(), + )); + assert!(right_error.next().await.unwrap().is_ok()); + assert!(right_error.next().await.unwrap().is_err()); + assert!(right_error.next().await.is_none()); + + let right_execute_error = Arc::new( + TestExec::new( + schema.clone(), + vec![], + Arc::new(AtomicUsize::new(0)), + Arc::new(AtomicUsize::new(0)), + ) + .with_execute_error("right execute failed"), + ); + let mut right_execute_error = Box::pin(test_union_stream( + test_stream(schema.clone(), vec![], Arc::new(AtomicUsize::new(0))), + right_execute_error, + schema, + )); + assert!(right_execute_error.next().await.unwrap().is_err()); + assert!(right_execute_error.next().await.is_none()); + } + + #[tokio::test] + async fn rhs_partial_and_full_batches_rebind_to_declared_left_schema() { + let left_schema = schema("left", "label", false); + let right_schema = schema("right", "label", true); + let (left, right) = schemas(left_schema.clone(), right_schema.clone()); + let plan = UnionDistinctOn::try_new(left, right, vec![1], 0).unwrap(); + let declared_schema = plan.output_schema.inner().clone(); + let partial_rhs = RecordBatch::try_new( + right_schema.clone(), + vec![ + Arc::new(Int64Array::from(vec![1, 2])), + Arc::new(StringArray::from(vec![Some("match"), None])), + Arc::new(Float64Array::from(vec![Some(10.0), None])), + ], + ) + .unwrap(); + let full_rhs = batch(right_schema.clone(), 3, Some("full"), Some(30.0)); + let result = execute_batches( + &plan, + left_schema.clone(), + vec![batch(left_schema, 1, Some("match"), Some(1.0))], + right_schema, + vec![partial_rhs, full_rhs], + ) + .await; + assert!(result.iter().all(|batch| batch.schema() == declared_schema)); + assert_eq!(result.iter().map(RecordBatch::num_rows).sum::(), 3); + assert!(result[1].column(1).is_null(0)); + } + + #[test] + fn malformed_indices_and_incompatible_inputs_fail_before_execution() { let schema = Arc::new(Schema::new(vec![ Field::new("ts", DataType::Int64, false), Field::new("job", DataType::Utf8, false), - Field::new("value", DataType::Float64, false), ])); - let df_schema = schema.clone().to_dfschema_ref().unwrap(); - let left_plan = LogicalPlan::EmptyRelation(EmptyRelation { - produce_one_row: false, - schema: df_schema.clone(), - }); - let right_plan = LogicalPlan::EmptyRelation(EmptyRelation { - produce_one_row: false, - schema: df_schema.clone(), - }); - let plan_node = UnionDistinctOn::new( - left_plan.clone(), - right_plan.clone(), - vec!["job".to_string()], - "ts".to_string(), - df_schema.clone(), - ); - - let bytes = plan_node.serialize(); - - let union_distinct_on = UnionDistinctOn::deserialize(&bytes).unwrap(); - let union_distinct_on = union_distinct_on - .with_exprs_and_inputs(vec![], vec![left_plan, right_plan]) + let invalid = |compare_key_indices, ts_col_idx| { + let decoded = UnionDistinctOn::deserialize( + &pb::UnionDistinctOn { + compare_key_indices, + ts_col_idx, + } + .encode_to_vec(), + ) .unwrap(); - - assert_eq!(union_distinct_on.compare_keys, vec!["job".to_string()]); - assert_eq!(union_distinct_on.ts_col, "ts"); - assert_eq!(union_distinct_on.output_schema, df_schema); + let (left, right) = schemas(schema.clone(), schema.clone()); + decoded + .with_exprs_and_inputs(vec![], vec![left, right]) + .is_err() + }; + assert!(invalid(vec![2], 0)); + assert!(invalid(vec![1], 2)); + let incompatible_schema = Arc::new(Schema::new(vec![ + Field::new("other_ts", DataType::Int64, false), + Field::new("other_job", DataType::Int64, false), + ])); + let (left, right) = schemas(schema, incompatible_schema); + assert!(UnionDistinctOn::try_new(left, right, vec![1], 0).is_err()); } } diff --git a/src/promql/src/lib.rs b/src/promql/src/lib.rs index 2de1c6508e..78bd012c4d 100644 --- a/src/promql/src/lib.rs +++ b/src/promql/src/lib.rs @@ -12,8 +12,6 @@ // See the License for the specific language governing permissions and // limitations under the License. -#![feature(map_try_insert)] - pub mod error; pub mod extension_plan; pub mod functions; diff --git a/src/query/src/promql/planner.rs b/src/query/src/promql/planner.rs index 61913156cf..284701f337 100644 --- a/src/query/src/promql/planner.rs +++ b/src/query/src/promql/planner.rs @@ -50,6 +50,7 @@ use datafusion_expr::utils::conjunction; use datafusion_expr::{ ExprSchemable, Literal, Projection, SortExpr, TableScan, TableSource, col, lit, }; +use datafusion_functions::core::coalesce; use datatypes::arrow::datatypes::{DataType as ArrowDataType, TimeUnit as ArrowTimeUnit}; use datatypes::data_type::ConcreteDataType; use itertools::Itertools; @@ -4423,6 +4424,59 @@ impl PromPlanner { // Take the name of first field column. The length is checked above. let left_field_col = left_context.field_columns.first().unwrap(); let right_field_col = right_context.field_columns.first().unwrap(); + let left_field = left + .schema() + .iter() + .find(|(_, field)| field.name() == left_field_col) + .map(|(qualifier, field)| (qualifier.cloned(), field.data_type().clone())) + .with_context(|| ColumnNotFoundSnafu { + col: left_field_col.clone(), + })?; + let right_field = right + .schema() + .iter() + .find(|(_, field)| field.name() == right_field_col) + .map(|(qualifier, field)| (qualifier.cloned(), field.data_type().clone())) + .with_context(|| ColumnNotFoundSnafu { + col: right_field_col.clone(), + })?; + let target_field_type = if left_field.1 == right_field.1 { + left_field.1.clone() + } else if matches!( + left_field.1, + ArrowDataType::Int8 + | ArrowDataType::Int16 + | ArrowDataType::Int32 + | ArrowDataType::Int64 + | ArrowDataType::UInt8 + | ArrowDataType::UInt16 + | ArrowDataType::UInt32 + | ArrowDataType::UInt64 + | ArrowDataType::Float32 + | ArrowDataType::Float64 + ) && matches!( + right_field.1, + ArrowDataType::Int8 + | ArrowDataType::Int16 + | ArrowDataType::Int32 + | ArrowDataType::Int64 + | ArrowDataType::UInt8 + | ArrowDataType::UInt16 + | ArrowDataType::UInt32 + | ArrowDataType::UInt64 + | ArrowDataType::Float32 + | ArrowDataType::Float64 + ) { + ArrowDataType::Float64 + } else { + return UnexpectedPlanExprSnafu { + desc: format!( + "OR value fields have incompatible types: {:?} and {:?}", + left_field.1, right_field.1 + ), + } + .fail(); + }; let left_has_tsid = left .schema() .fields() @@ -4459,10 +4513,26 @@ impl PromPlanner { all_columns.sort_unstable(); // use left time index column name as the result time index column name all_columns.insert(0, left_time_index_column.clone()); + let mut occupied_column_names = left + .schema() + .fields() + .iter() + .chain(right.schema().fields().iter()) + .map(|field| field.name().clone()) + .collect::>(); // step 1: align schema using project, fill non-exist columns with null let left_proj_exprs = all_columns.iter().map(|col| { - if tags_not_in_left.contains(col) { + if col == left_field_col && left_field.1 != target_field_type { + DfExpr::Cast(Cast { + expr: Box::new(DfExpr::Column(Column::new( + left_field.0.clone(), + left_field_col, + ))), + data_type: target_field_type.clone(), + }) + .alias(left_field_col.clone()) + } else if tags_not_in_left.contains(col) { DfExpr::Literal(ScalarValue::Utf8(None), None).alias(col.clone()) } else { DfExpr::Column(Column::new(None::, col)) @@ -4475,25 +4545,22 @@ impl PromPlanner { .alias(left_time_index_column.clone()); // The field column in right side may not have qualifier (it may be removed by join operation), // so we need to find it from the schema. - let right_qualifier_for_field = right - .schema() - .iter() - .find(|(_, f)| f.name() == right_field_col) - .map(|(q, _)| q) - .with_context(|| ColumnNotFoundSnafu { - col: right_field_col.clone(), - })? - .cloned(); - // `skip(1)` to skip the time index column let right_proj_exprs_without_time_index = all_columns.iter().skip(1).map(|col| { // expr - if col == left_field_col && left_field_col != right_field_col { - // qualify field in right side if necessary to handle different field name - DfExpr::Column(Column::new( - right_qualifier_for_field.clone(), - right_field_col, - )) + if col == left_field_col { + let expr = DfExpr::Column(Column::new(right_field.0.clone(), right_field_col)); + if right_field.1 != target_field_type { + DfExpr::Cast(Cast { + expr: Box::new(expr), + data_type: target_field_type.clone(), + }) + .alias(left_field_col.clone()) + } else if left_field_col != right_field_col { + expr.alias(left_field_col.clone()) + } else { + expr + } } else if tags_not_in_right.contains(col) { DfExpr::Literal(ScalarValue::Utf8(None), None).alias(col.clone()) } else { @@ -4537,24 +4604,168 @@ impl PromPlanner { }; // sort to ensure the generated plan is not volatile match_columns.sort_unstable(); - // step 3: build `UnionDistinctOn` plan - let schema = left_projected.schema().clone(); - let union_distinct_on = UnionDistinctOn::new( - left_projected, - right_projected, - match_columns, - left_time_index_column.clone(), - schema, + match_columns.dedup(); + occupied_column_names.extend( + left_projected + .schema() + .fields() + .iter() + .chain(right_projected.schema().fields().iter()) + .map(|field| field.name().clone()), ); - let result = LogicalPlan::Extension(Extension { + + let visible_schema = left_projected.schema().clone(); + let visible_left_exprs = left_projected + .schema() + .iter() + .map(|(qualifier, field)| { + DfExpr::Column(Column::new(qualifier.cloned(), field.name().clone())) + }) + .collect::>(); + let visible_right_exprs = right_projected + .schema() + .iter() + .map(|(qualifier, field)| { + DfExpr::Column(Column::new(qualifier.cloned(), field.name().clone())) + }) + .collect::>(); + let mut left_match_exprs = Vec::with_capacity(match_columns.len()); + let mut right_match_exprs = Vec::with_capacity(match_columns.len()); + let mut next_internal_column = 0; + + for label in &match_columns { + let left_field = if left_tag_cols_set.contains(label) { + Some( + left_projected + .schema() + .iter() + .find(|(_, field)| field.name() == label) + .map(|(qualifier, field)| (qualifier.cloned(), field.data_type().clone())) + .with_context(|| ColumnNotFoundSnafu { col: label.clone() })?, + ) + } else { + None + }; + let right_field = if right_tag_cols_set.contains(label) { + Some( + right_projected + .schema() + .iter() + .find(|(_, field)| field.name() == label) + .map(|(qualifier, field)| (qualifier.cloned(), field.data_type().clone())) + .with_context(|| ColumnNotFoundSnafu { col: label.clone() })?, + ) + } else { + None + }; + let data_type = match (left_field.as_ref(), right_field.as_ref()) { + (Some((_, left_type)), Some((_, right_type))) if left_type == right_type => { + left_type.clone() + } + (Some((_, left_type)), Some((_, right_type))) => { + return UnexpectedPlanExprSnafu { + desc: format!( + "OR match label {label} has incompatible types: {left_type:?} and {right_type:?}" + ), + } + .fail(); + } + (Some((_, data_type)), None) | (None, Some((_, data_type))) => data_type.clone(), + (None, None) => ArrowDataType::Utf8, + }; + let empty = match data_type { + ArrowDataType::Utf8 => ScalarValue::Utf8(Some(String::new())), + ArrowDataType::LargeUtf8 => ScalarValue::LargeUtf8(Some(String::new())), + _ => { + return UnexpectedPlanExprSnafu { + desc: format!("OR match label {label} must be a string"), + } + .fail(); + } + }; + let internal_name = loop { + let name = format!("__promql_or_match_{next_internal_column}"); + next_internal_column += 1; + if occupied_column_names.insert(name.clone()) { + break name; + } + }; + let normalize = |field: Option<(Option, ArrowDataType)>| { + let expr = if let Some((qualifier, _)) = field { + DfExpr::ScalarFunction(ScalarFunction { + func: coalesce(), + args: vec![ + DfExpr::Column(Column::new(qualifier, label.clone())), + DfExpr::Literal(empty.clone(), None), + ], + }) + } else { + DfExpr::Literal(empty.clone(), None) + }; + expr.alias(internal_name.clone()) + }; + left_match_exprs.push(normalize(left_field)); + right_match_exprs.push(normalize(right_field)); + } + + let left_augmented = LogicalPlanBuilder::from(left_projected) + .project(visible_left_exprs.into_iter().chain(left_match_exprs)) + .context(DataFusionPlanningSnafu)? + .build() + .context(DataFusionPlanningSnafu)?; + let right_augmented = LogicalPlanBuilder::from(right_projected) + .project(visible_right_exprs.into_iter().chain(right_match_exprs)) + .context(DataFusionPlanningSnafu)? + .build() + .context(DataFusionPlanningSnafu)?; + + // step 3: build `UnionDistinctOn` with normalized internal match keys. + let visible_field_count = visible_schema.fields().len(); + let compare_key_indices = + (visible_field_count..visible_field_count + match_columns.len()).collect::>(); + let (time_qualifier, _) = visible_schema + .iter() + .find(|(_, field)| field.name() == &left_time_index_column) + .with_context(|| TimeIndexNotFoundSnafu { + table: left_qualifier_string.clone(), + })?; + let ts_col_idx = left_augmented + .schema() + .iter() + .position(|(qualifier, field)| { + qualifier == time_qualifier && field.name() == &left_time_index_column + }) + .with_context(|| TimeIndexNotFoundSnafu { + table: left_qualifier_string.clone(), + })?; + let union_distinct_on = UnionDistinctOn::try_new( + left_augmented, + right_augmented, + compare_key_indices, + ts_col_idx, + ) + .context(DataFusionPlanningSnafu)?; + let augmented_result = LogicalPlan::Extension(Extension { node: Arc::new(union_distinct_on), }); + let result = LogicalPlanBuilder::from(augmented_result) + .project(visible_schema.iter().map(|(qualifier, field)| { + DfExpr::Column(Column::new(qualifier.cloned(), field.name().clone())) + })) + .context(DataFusionPlanningSnafu)? + .build() + .context(DataFusionPlanningSnafu)?; // step 4: update context - self.ctx.time_index_column = Some(left_time_index_column); - self.ctx.tag_columns = all_tags.into_iter().collect(); - self.ctx.field_columns = vec![left_field_col.clone()]; - self.ctx.use_tsid = left_has_tsid && right_has_tsid; + let output_field_col = left_field_col.clone(); + let mut output_context = left_context; + let mut visible_tags = all_tags.into_iter().collect::>(); + visible_tags.sort_unstable(); + output_context.time_index_column = Some(left_time_index_column); + output_context.tag_columns = visible_tags; + output_context.field_columns = vec![output_field_col]; + output_context.use_tsid = left_has_tsid && right_has_tsid; + self.ctx = output_context; Ok(result) } @@ -4761,15 +4972,23 @@ mod test { use common_catalog::consts::{DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME}; use common_query::prelude::greptime_timestamp; use common_query::test_util::DummyDecoder; - use datafusion::arrow::datatypes::Schema as ArrowSchema; + use datafusion::arrow::array::{ + Array, Float64Array, Int64Array, StringArray, TimestampMillisecondArray, + }; + use datafusion::arrow::datatypes::{Field, Schema as ArrowSchema}; + use datafusion::arrow::record_batch::RecordBatch; + use datafusion::catalog::{CatalogProvider, MemoryCatalogProvider, MemorySchemaProvider}; use datafusion::datasource::memory::MemorySourceConfig; use datafusion::datasource::source::DataSourceExec; + use datafusion::datasource::{MemTable, provider_as_source}; + use datafusion::execution::context::SessionContext; use datafusion::logical_expr::Extension; use datatypes::prelude::ConcreteDataType; use datatypes::schema::{ColumnSchema, Schema}; use promql_parser::label::Labels; use promql_parser::parser; use session::context::QueryContext; + use substrait::{DFLogicalSubstraitConvertor, SubstraitPlan}; use table::metadata::{TableInfoBuilder, TableMetaBuilder}; use table::test_util::EmptyTable; @@ -4777,6 +4996,7 @@ mod test { use crate::QueryEngineContext; use crate::options::QueryOptions; use crate::parser::QueryLanguageParser; + use crate::query_engine::DefaultSerializer; fn find_instant_manipulate(plan: &LogicalPlan) -> Option<&InstantManipulate> { if let LogicalPlan::Extension(Extension { node }) = plan @@ -4872,6 +5092,262 @@ mod test { } } + enum DirectOrValue { + Float64(f64), + Int64(i64), + Utf8(&'static str), + } + + impl DirectOrValue { + fn data_type(&self) -> ArrowDataType { + match self { + Self::Float64(_) => ArrowDataType::Float64, + Self::Int64(_) => ArrowDataType::Int64, + Self::Utf8(_) => ArrowDataType::Utf8, + } + } + fn array(&self) -> Arc { + match self { + Self::Float64(v) => Arc::new(Float64Array::from(vec![*v])), + Self::Int64(v) => Arc::new(Int64Array::from(vec![*v])), + Self::Utf8(v) => Arc::new(StringArray::from(vec![*v])), + } + } + } + + struct DirectOrSource { + name: &'static str, + empty: bool, + timestamp: i64, + tags: Vec<(&'static str, Option<&'static str>)>, + value: DirectOrValue, + } + + fn source( + name: &'static str, + empty: bool, + timestamp: i64, + tags: Vec<(&'static str, Option<&'static str>)>, + value: DirectOrValue, + ) -> DirectOrSource { + DirectOrSource { + name, + empty, + timestamp, + tags, + value, + } + } + + fn tagged_source( + name: &'static str, + empty: bool, + tag: (&'static str, Option<&'static str>), + value: DirectOrValue, + ) -> DirectOrSource { + source(name, empty, 1, vec![("job", Some("job")), tag], value) + } + + fn job_source(name: &'static str, value: DirectOrValue) -> DirectOrSource { + source(name, true, 1, vec![("job", Some("job"))], value) + } + + fn table(source: &DirectOrSource) -> Arc { + let mut fields = vec![Field::new( + "ts", + ArrowDataType::Timestamp(ArrowTimeUnit::Millisecond, None), + false, + )]; + fields.extend( + source + .tags + .iter() + .map(|(name, _)| Field::new(*name, ArrowDataType::Utf8, true)), + ); + fields.push(Field::new("v", source.value.data_type(), true)); + let schema = Arc::new(ArrowSchema::new(fields)); + let partitions = if source.empty { + vec![vec![]] + } else { + let mut columns: Vec> = + vec![Arc::new(TimestampMillisecondArray::from(vec![ + source.timestamp, + ]))]; + columns.extend( + source + .tags + .iter() + .map(|(_, value)| Arc::new(StringArray::from(vec![*value])) as Arc), + ); + columns.push(source.value.array()); + vec![vec![RecordBatch::try_new(schema.clone(), columns).unwrap()]] + }; + Arc::new(MemTable::try_new(schema, partitions).unwrap()) + } + + fn scan(source: &DirectOrSource) -> LogicalPlan { + LogicalPlanBuilder::scan(source.name, provider_as_source(table(source)), None) + .unwrap() + .build() + .unwrap() + } + + fn direct_or_context(qualifier: &str, tags: &[&str], field: &str) -> PromPlannerContext { + PromPlannerContext { + table_name: Some(qualifier.to_string()), + time_index_column: Some("ts".to_string()), + field_columns: vec![field.to_string()], + tag_columns: tags.iter().map(|tag| (*tag).to_string()).collect(), + ..Default::default() + } + } + + fn or_modifier(expr: &str) -> Option { + let PromExpr::Binary(expr) = parser::parse(expr).unwrap() else { + unreachable!() + }; + expr.modifier + } + + async fn plan_direct_or( + left: LogicalPlan, + right: LogicalPlan, + left_context: PromPlannerContext, + right_context: PromPlannerContext, + modifier: &Option, + ) -> LogicalPlan { + let table_provider = build_test_table_provider_with_fields( + &[(DEFAULT_SCHEMA_NAME.to_string(), "dummy".to_string())], + &[], + ) + .await; + let mut planner = PromPlanner { + table_provider, + ctx: PromPlannerContext::default(), + }; + planner + .or_operator( + left, + right, + left_context.tag_columns.iter().cloned().collect(), + right_context.tag_columns.iter().cloned().collect(), + left_context, + right_context, + modifier, + ) + .unwrap() + } + + async fn execute( + plan: LogicalPlan, + state: &QueryEngineState, + ) -> (LogicalPlan, Vec) { + let context = QueryEngineContext::new(state.session_state(), QueryContext::arc()); + let optimized = state.optimize_by_extension_rules(plan, &context).unwrap(); + let physical = state + .session_state() + .create_physical_plan(&optimized) + .await + .unwrap(); + let batches = + datafusion::physical_plan::collect(physical, state.session_state().task_ctx()) + .await + .unwrap(); + (optimized, batches) + } + + async fn run( + left: &DirectOrSource, + right: &DirectOrSource, + left_context: PromPlannerContext, + right_context: PromPlannerContext, + modifier: &Option, + ) -> (LogicalPlan, Vec) { + let plan = plan_direct_or( + scan(left), + scan(right), + left_context, + right_context, + modifier, + ) + .await; + execute(plan, &build_query_engine_state()).await + } + + fn assert_no_internal_or_keys(schema: &DFSchema) { + assert!( + schema + .fields() + .iter() + .all(|field| !field.name().starts_with("__promql_or_match_")), + "{schema:?}" + ); + } + + fn values(batches: &[RecordBatch], column: &str) -> Vec { + batches + .iter() + .flat_map(|batch| { + batch + .column_by_name(column) + .unwrap() + .as_any() + .downcast_ref::() + .unwrap() + .iter() + .flatten() + }) + .collect() + } + + fn rows(batches: &[RecordBatch]) -> Vec<(f64, Option)> { + let mut rows = batches + .iter() + .flat_map(|batch| { + let values = batch + .column_by_name("v") + .unwrap() + .as_any() + .downcast_ref::() + .unwrap(); + let labels = batch + .column_by_name("k") + .map(|column| column.as_any().downcast_ref::().unwrap()); + (0..batch.num_rows()).map(move |i| { + ( + values.value(i), + labels.and_then(|labels| { + (!labels.is_null(i)).then(|| labels.value(i).to_string()) + }), + ) + }) + }) + .collect::>(); + rows.sort_by(|left, right| left.0.total_cmp(&right.0)); + rows + } + + fn matrix_source( + name: &'static str, + k: Option>, + timestamp: i64, + value: f64, + ) -> DirectOrSource { + let mut tags = vec![("job", Some("job"))]; + if let Some(k) = k { + tags.push(("k", k)); + } + source(name, false, timestamp, tags, DirectOrValue::Float64(value)) + } + + fn matrix_context(name: &str, k: Option>) -> PromPlannerContext { + direct_or_context( + name, + if k.is_some() { &["job", "k"] } else { &["job"] }, + "v", + ) + } + async fn build_test_table_provider( table_name_tuples: &[(String, String)], num_tag: usize, @@ -8191,47 +8667,33 @@ Projection: count(prometheus_tsdb_head_series.greptime_value) AS my_series, prom #[tokio::test] async fn test_or_not_exists_table_label() { - let mut eval_stmt = EvalStmt { - expr: PromExpr::NumberLiteral(NumberLiteral { val: 1.0 }), - start: UNIX_EPOCH, - end: UNIX_EPOCH - .checked_add(Duration::from_secs(100_000)) - .unwrap(), - interval: Duration::from_secs(5), - lookback_delta: Duration::from_secs(1), - }; - let case = r#"sum by (job, tag0, tag2) (metric_exists) or sum by (job, tag0, tag2) (metric_not_exists)"#; - - let prom_expr = parser::parse(case).unwrap(); - eval_stmt.expr = prom_expr; - let table_provider = build_test_table_provider_with_fields( - &[(DEFAULT_SCHEMA_NAME.to_string(), "metric_exists".to_string())], + let state = build_query_engine_state(); + let provider = build_test_table_provider_with_fields( + &[(DEFAULT_SCHEMA_NAME.to_string(), "normal_metric".to_string())], &["job"], ) .await; - - let plan = - PromPlanner::stmt_to_plan(table_provider, &eval_stmt, &build_query_engine_state()) - .await - .unwrap(); - let expected = r#"UnionDistinctOn: on col=[["job"]], ts_col=[greptime_timestamp] [greptime_timestamp:Timestamp(ms), job:Utf8, sum(metric_exists.greptime_value):Float64;N] - SubqueryAlias: metric_exists [greptime_timestamp:Timestamp(ms), job:Utf8, sum(metric_exists.greptime_value):Float64;N] - Projection: metric_exists.greptime_timestamp, metric_exists.job, sum(metric_exists.greptime_value) [greptime_timestamp:Timestamp(ms), job:Utf8, sum(metric_exists.greptime_value):Float64;N] - Sort: metric_exists.job ASC NULLS LAST, metric_exists.greptime_timestamp ASC NULLS LAST [job:Utf8, greptime_timestamp:Timestamp(ms), sum(metric_exists.greptime_value):Float64;N] - Aggregate: groupBy=[[metric_exists.job, metric_exists.greptime_timestamp]], aggr=[[sum(metric_exists.greptime_value)]] [job:Utf8, greptime_timestamp:Timestamp(ms), sum(metric_exists.greptime_value):Float64;N] - PromInstantManipulate: range=[0..100000000], lookback=[1000], interval=[5000], time index=[greptime_timestamp] [job:Utf8, greptime_timestamp:Timestamp(ms), greptime_value:Float64;N] - PromSeriesDivide: tags=["job"] [job:Utf8, greptime_timestamp:Timestamp(ms), greptime_value:Float64;N] - Sort: metric_exists.job ASC NULLS FIRST, metric_exists.greptime_timestamp ASC NULLS FIRST [job:Utf8, greptime_timestamp:Timestamp(ms), greptime_value:Float64;N] - Filter: metric_exists.greptime_timestamp >= TimestampMillisecond(-999, None) AND metric_exists.greptime_timestamp <= TimestampMillisecond(100000000, None) [job:Utf8, greptime_timestamp:Timestamp(ms), greptime_value:Float64;N] - TableScan: metric_exists [job:Utf8, greptime_timestamp:Timestamp(ms), greptime_value:Float64;N] - SubqueryAlias: [greptime_timestamp:Timestamp(ms), job:Utf8;N, sum(.value):Float64;N] - Projection: .time AS greptime_timestamp, Utf8(NULL) AS job, sum(.value) [greptime_timestamp:Timestamp(ms), job:Utf8;N, sum(.value):Float64;N] - Sort: .time ASC NULLS LAST [time:Timestamp(ms), sum(.value):Float64;N] - Aggregate: groupBy=[[.time]], aggr=[[sum(.value)]] [time:Timestamp(ms), sum(.value):Float64;N] - EmptyMetric: range=[0..-1], interval=[5000] [time:Timestamp(ms), value:Float64;N] - TableScan: dummy [time:Timestamp(ms), value:Float64;N]"#; - - assert_eq!(plan.display_indent_schema().to_string(), expected); + let raw = PromPlanner::stmt_to_plan( + provider, + &build_eval_stmt(r#"missing_metric or on(absent_label) normal_metric"#), + &state, + ) + .await + .unwrap(); + assert!( + raw.display_indent_schema() + .to_string() + .contains("__promql_or_match_0@") + ); + let (optimized, batches) = execute(raw, &state).await; + assert_no_internal_or_keys(optimized.schema()); + assert!(batches.iter().all(|batch| { + batch + .schema() + .fields() + .iter() + .all(|field| !field.name().starts_with("__promql_or_match_")) + })); } #[tokio::test] @@ -8283,4 +8745,250 @@ Projection: count(prometheus_tsdb_head_series.greptime_value) AS my_series, prom _ => panic!("Expected EmptyRelation, but got: {:?}", plan), } } + + #[tokio::test] + async fn test_direct_or_normalizes_missing_match_labels() { + type Case<'a> = ( + Option>, + Option>, + i64, + i64, + &'a [(f64, Option<&'a str>)], + ); + + let modifier = or_modifier("lhs or on(k) rhs"); + #[rustfmt::skip] + let cases: &[Case<'_>] = &[ + (None, None, 1, 1, &[(1.0, None)]), + (None, Some(Some("")), 1, 1, &[(1.0, None)]), + (Some(Some("")), None, 1, 1, &[(1.0, Some(""))]), + (None, Some(Some("r")), 1, 1, &[(1.0, None), (2.0, Some("r"))]), + (Some(Some("l")), None, 1, 1, &[(1.0, Some("l")), (2.0, None)]), + (Some(None), Some(Some("")), 1, 1, &[(1.0, None)]), + (Some(None), Some(Some("r")), 1, 1, &[(1.0, None), (2.0, Some("r"))]), + (Some(Some("same")), Some(Some("same")), 1, 2, &[(1.0, Some("same")), (2.0, Some("same"))]), + ]; + for &(left, right, left_ts, right_ts, expected) in cases { + let (optimized, batches) = run( + &matrix_source("lhs", left, left_ts, 1.0), + &matrix_source("rhs", right, right_ts, 2.0), + matrix_context("lhs", left), + matrix_context("rhs", right), + &modifier, + ) + .await; + assert_no_internal_or_keys(optimized.schema()); + assert_eq!( + rows(&batches), + expected + .iter() + .map(|(value, label)| (*value, label.map(str::to_string))) + .collect::>() + ); + } + } + + #[tokio::test] + async fn test_direct_or_match_modifiers() { + for (modifier, left, right, expected) in [ + (None, "left", "right", 2), + (or_modifier("lhs or on(k) rhs"), "same", "same", 1), + (or_modifier("lhs or on() rhs"), "left", "right", 1), + (or_modifier("lhs or ignoring(k) rhs"), "left", "right", 1), + ] { + let (_, batches) = run( + &matrix_source("lhs", Some(Some(left)), 1, 1.0), + &matrix_source("rhs", Some(Some(right)), 1, 2.0), + direct_or_context("lhs", &["job", "k"], "v"), + direct_or_context("rhs", &["job", "k"], "v"), + &modifier, + ) + .await; + assert_eq!( + batches.iter().map(RecordBatch::num_rows).sum::(), + expected + ); + } + } + + #[tokio::test] + async fn test_direct_or_nested_projection_uses_left_context() { + let left = matrix_source("lhs", Some(Some("k")), 1, 1.0); + let right = matrix_source("rhs", Some(Some("k")), 1, 2.0); + let raw = plan_direct_or( + scan(&left), + scan(&right), + direct_or_context("lhs", &["job", "k"], "v"), + direct_or_context("rhs", &["job", "k"], "v"), + &or_modifier("lhs or on(k) rhs"), + ) + .await; + assert!(raw.schema().iter().any(|(qualifier, field)| { + qualifier.as_ref().is_some_and(|q| q.to_string() == "lhs") && field.name() == "v" + })); + let nested = LogicalPlanBuilder::from(raw) + .project(vec![ + DfExpr::BinaryExpr(BinaryExpr { + left: Box::new(DfExpr::Column(Column::new( + Some(TableReference::bare("lhs")), + "v", + ))), + op: Operator::Plus, + right: Box::new(lit(1.0)), + }) + .alias("v_plus"), + ]) + .unwrap() + .build() + .unwrap(); + let (_, batches) = execute(nested, &build_query_engine_state()).await; + assert_eq!(values(&batches, "v_plus"), vec![2.0]); + } + + #[tokio::test] + async fn test_direct_or_skips_user_internal_key_name() { + const USER_TAG: &str = "__promql_or_match_0"; + let left = tagged_source( + "lhs", + false, + (USER_TAG, Some("left")), + DirectOrValue::Float64(1.0), + ); + let right = tagged_source( + "rhs", + false, + (USER_TAG, Some("right")), + DirectOrValue::Float64(2.0), + ); + let raw = plan_direct_or( + scan(&left), + scan(&right), + direct_or_context("lhs", &["job", USER_TAG], "v"), + direct_or_context("rhs", &["job", USER_TAG], "v"), + &or_modifier("lhs or on(missing_label) rhs"), + ) + .await; + assert!( + raw.display_indent_schema() + .to_string() + .contains("__promql_or_match_1@") + ); + let (_, batches) = execute(raw, &build_query_engine_state()).await; + assert!( + batches + .iter() + .all(|batch| batch.column_by_name(USER_TAG).is_some()) + ); + } + + #[tokio::test] + async fn test_direct_or_substrait_round_trip_with_normalized_key() { + let state = build_query_engine_state(); + let ctx = SessionContext::new_with_state(state.session_state()); + let catalog = Arc::new(MemoryCatalogProvider::new()); + catalog + .register_schema("public", Arc::new(MemorySchemaProvider::new())) + .unwrap(); + ctx.register_catalog("datafusion", catalog); + let left = matrix_source("lhs", Some(Some("")), 1, 1.0); + let right = matrix_source("rhs", None, 1, 2.0); + ctx.register_table( + TableReference::full("datafusion", "public", "lhs"), + table(&left), + ) + .unwrap(); + ctx.register_table( + TableReference::full("datafusion", "public", "rhs"), + table(&right), + ) + .unwrap(); + let raw = plan_direct_or( + ctx.table("datafusion.public.lhs") + .await + .unwrap() + .into_unoptimized_plan(), + ctx.table("datafusion.public.rhs") + .await + .unwrap() + .into_unoptimized_plan(), + direct_or_context("lhs", &["job", "k"], "v"), + direct_or_context("rhs", &["job"], "v"), + &or_modifier("lhs or on(k) rhs"), + ) + .await; + let decoded = DFLogicalSubstraitConvertor + .decode( + DFLogicalSubstraitConvertor + .encode(&raw, DefaultSerializer) + .unwrap(), + ctx.state(), + ) + .await + .unwrap(); + let (optimized, batches) = execute(decoded, &state).await; + assert_no_internal_or_keys(optimized.schema()); + assert!(batches.iter().all(|batch| { + batch + .schema() + .fields() + .iter() + .all(|field| !field.name().starts_with("__promql_or_match_")) + })); + assert_eq!(values(&batches, "v"), vec![1.0]); + } + + #[tokio::test] + async fn test_direct_or_numeric_value_types() { + let left = tagged_source("lhs", true, ("k", Some("lhs")), DirectOrValue::Int64(0)); + let right = tagged_source( + "rhs", + false, + ("k", Some("rhs")), + DirectOrValue::Float64(0.5), + ); + let (optimized, batches) = run( + &left, + &right, + direct_or_context("lhs", &["job", "k"], "v"), + direct_or_context("rhs", &["job", "k"], "v"), + &or_modifier("lhs or on(k) rhs"), + ) + .await; + assert_eq!( + optimized + .schema() + .field_with_name(None, "v") + .unwrap() + .data_type(), + &ArrowDataType::Float64 + ); + assert_eq!(values(&batches, "v"), vec![0.5]); + let provider = build_test_table_provider_with_fields( + &[(DEFAULT_SCHEMA_NAME.to_string(), "dummy".to_string())], + &[], + ) + .await; + let mut planner = PromPlanner { + table_provider: provider, + ctx: PromPlannerContext::default(), + }; + let left_context = direct_or_context("lhs", &["job"], "v"); + let right_context = direct_or_context("rhs", &["job"], "v"); + let error = planner + .or_operator( + scan(&job_source("lhs", DirectOrValue::Utf8("x"))), + scan(&job_source("rhs", DirectOrValue::Float64(1.0))), + left_context.tag_columns.iter().cloned().collect(), + right_context.tag_columns.iter().cloned().collect(), + left_context, + right_context, + &or_modifier("lhs or on() rhs"), + ) + .unwrap_err(); + assert!( + error + .to_string() + .contains("OR value fields have incompatible types") + ); + } } diff --git a/tests/cases/standalone/common/promql/set_operation.result b/tests/cases/standalone/common/promql/set_operation.result index 83af6fc875..50a0b6cf3b 100644 --- a/tests/cases/standalone/common/promql/set_operation.result +++ b/tests/cases/standalone/common/promql/set_operation.result @@ -466,12 +466,12 @@ tql eval (0, 2000, '400') sum(t1{job="a"}); -- SQLNESS SORT_RESULT 3 1 tql eval (0, 2000, '400') sum(t1{job="a"}) - sum(t1{job="e"} or vector(1)); -+---------------------+-----------------------------------------------------+ -| ts | t1.sum(t1.greptime_value) - .sum(t1.greptime_value) | -+---------------------+-----------------------------------------------------+ -| 1970-01-01T00:00:00 | 0.0 | -| 1970-01-01T00:20:00 | 2.0 | -+---------------------+-----------------------------------------------------+ ++---------------------+---------------------------------------------------------+ +| ts | lhs.sum(t1.greptime_value) - rhs.sum(t1.greptime_value) | ++---------------------+---------------------------------------------------------+ +| 1970-01-01T00:00:00 | 0.0 | +| 1970-01-01T00:20:00 | 2.0 | ++---------------------+---------------------------------------------------------+ drop table t1; @@ -545,12 +545,12 @@ tql eval (0, 2000, '400') count(max by (namespace) (stats_used_bytes{namespace=~ +---------------------+---------------------------------------------+ | ts | count(max(stats_used_bytes.greptime_value)) | +---------------------+---------------------------------------------+ -| 1970-01-01T00:00:00 | 0 | -| 1970-01-01T00:06:40 | 0 | -| 1970-01-01T00:13:20 | 0 | -| 1970-01-01T00:20:00 | 2 | -| 1970-01-01T00:26:40 | 0 | -| 1970-01-01T00:33:20 | 0 | +| 1970-01-01T00:00:00 | 0.0 | +| 1970-01-01T00:06:40 | 0.0 | +| 1970-01-01T00:13:20 | 0.0 | +| 1970-01-01T00:20:00 | 2.0 | +| 1970-01-01T00:26:40 | 0.0 | +| 1970-01-01T00:33:20 | 0.0 | +---------------------+---------------------------------------------+ -- SQLNESS SORT_RESULT 3 1 diff --git a/tests/cases/standalone/common/promql/tsid_binary_join_regression.result b/tests/cases/standalone/common/promql/tsid_binary_join_regression.result index d5313ba41c..8c90d97a32 100644 --- a/tests/cases/standalone/common/promql/tsid_binary_join_regression.result +++ b/tests/cases/standalone/common/promql/tsid_binary_join_regression.result @@ -512,10 +512,10 @@ TQL ANALYZE (0, 5, '5s') (tsid_binary_join_left or tsid_binary_join_right) / tsi +-+-+-+ | stage | node | plan_| +-+-+-+ -| 0_| 0_|_ProjectionExec: expr=[host@2 as host, job@3 as job, ts@5 as ts, __tsid@4 as __tsid, greptime_value@0 / greptime_value@1 as tsid_binary_join_right.greptime_value / tsid_binary_join_left.greptime_value] REDACTED +| 0_| 0_|_ProjectionExec: expr=[host@2 as host, job@3 as job, ts@5 as ts, __tsid@4 as __tsid, greptime_value@0 / greptime_value@1 as lhs.greptime_value / rhs.greptime_value] REDACTED |_|_|_HashJoinExec: mode=CollectLeft, join_type=Inner, on=[(__tsid@1, __tsid@3), (ts@0, ts@4)], projection=[greptime_value@2, greptime_value@3, host@4, job@5, __tsid@6, ts@7], NullsEqual: true REDACTED |_|_|_ProjectionExec: expr=[ts@0 as ts, __tsid@1 as __tsid, greptime_value@2 as greptime_value] REDACTED -|_|_|_UnionDistinctOnExec: on col=[["host", "job"]], ts_col=[ts] REDACTED +|_|_|_UnionDistinctOnExec: on col=[5, 6], ts_col=0 REDACTED |_|_|_CoalescePartitionsExec REDACTED |_|_|_MergeScanExec: REDACTED |_|_|_CoalescePartitionsExec REDACTED @@ -523,14 +523,14 @@ TQL ANALYZE (0, 5, '5s') (tsid_binary_join_left or tsid_binary_join_right) / tsi |_|_|_CooperativeExec REDACTED |_|_|_MergeScanExec: REDACTED |_|_|_| -| 1_| 0_|_ProjectionExec: expr=[ts@4 as ts, __tsid@3 as __tsid, greptime_value@0 as greptime_value, host@1 as host, job@2 as job] REDACTED +| 1_| 0_|_ProjectionExec: expr=[ts@4 as ts, __tsid@3 as __tsid, greptime_value@0 as greptime_value, host@1 as host, job@2 as job, CASE WHEN host@1 IS NOT NULL THEN host@1 ELSE_END as __promql_or_match_0, CASE WHEN job@2 IS NOT NULL THEN job@2 ELSE_END as __promql_or_match_1] REDACTED |_|_|_PromInstantManipulateExec: range=[0..5000], lookback=[300000], interval=[5000], time index=[ts] REDACTED |_|_|_PromSeriesDivideExec: tags=["__tsid"] REDACTED |_|_|_ProjectionExec: expr=[greptime_value@1 as greptime_value, host@3 as host, job@4 as job, __tsid@2 as __tsid, ts@0 as ts] REDACTED |_|_|_CooperativeExec REDACTED |_|_|_SeriesScan: region=REDACTED, "partition_count":REDACTED, "distribution":"PerSeries" REDACTED |_|_|_| -| 1_| 0_|_ProjectionExec: expr=[ts@4 as ts, __tsid@3 as __tsid, greptime_value@0 as greptime_value, host@1 as host, job@2 as job] REDACTED +| 1_| 0_|_ProjectionExec: expr=[ts@4 as ts, __tsid@3 as __tsid, greptime_value@0 as greptime_value, host@1 as host, job@2 as job, CASE WHEN host@1 IS NOT NULL THEN host@1 ELSE_END as __promql_or_match_0, CASE WHEN job@2 IS NOT NULL THEN job@2 ELSE_END as __promql_or_match_1] REDACTED |_|_|_PromInstantManipulateExec: range=[0..5000], lookback=[300000], interval=[5000], time index=[ts] REDACTED |_|_|_PromSeriesDivideExec: tags=["__tsid"] REDACTED |_|_|_ProjectionExec: expr=[greptime_value@1 as greptime_value, host@3 as host, job@4 as job, __tsid@2 as __tsid, ts@0 as ts] REDACTED @@ -672,14 +672,14 @@ TQL EVAL (0, 5, '5s') ((tsid_binary_join_left > bool tsid_binary_join_right) + t -- SQLNESS SORT_RESULT 3 1 TQL EVAL (0, 5, '5s') (tsid_binary_join_left or tsid_binary_join_right) / tsid_binary_join_left; -+-------+------+---------------------+------------------------------------------------------------------------------+ -| host | job | ts | tsid_binary_join_right.greptime_value / tsid_binary_join_left.greptime_value | -+-------+------+---------------------+------------------------------------------------------------------------------+ -| host1 | job1 | 1970-01-01T00:00:00 | 1.0 | -| host1 | job1 | 1970-01-01T00:00:05 | 1.0 | -| host2 | job2 | 1970-01-01T00:00:00 | 1.0 | -| host2 | job2 | 1970-01-01T00:00:05 | 1.0 | -+-------+------+---------------------+------------------------------------------------------------------------------+ ++-------+------+---------------------+-----------------------------------------+ +| host | job | ts | lhs.greptime_value / rhs.greptime_value | ++-------+------+---------------------+-----------------------------------------+ +| host1 | job1 | 1970-01-01T00:00:00 | 1.0 | +| host1 | job1 | 1970-01-01T00:00:05 | 1.0 | +| host2 | job2 | 1970-01-01T00:00:00 | 1.0 | +| host2 | job2 | 1970-01-01T00:00:05 | 1.0 | ++-------+------+---------------------+-----------------------------------------+ -- Range functions are outside the first island version; the range selector subtree must -- remain a barrier.