From 5311e84dbc9d73f3c7fb514b618736d3702a4a6f Mon Sep 17 00:00:00 2001 From: kould Date: Sun, 4 Oct 2026 06:22:10 +0800 Subject: [PATCH 1/2] fix: preserve scalar cardinality and outer join match semantics Validate scalar results before caching, track outer join matches per build row, and use a consistent TopK tie-break ordering. Add regressions to existing SLT suites. --- src/execution/dql/join/hash/full_join.rs | 33 +++--- src/execution/dql/join/hash/inner_join.rs | 1 - src/execution/dql/join/hash/left_join.rs | 17 +-- src/execution/dql/join/hash/mod.rs | 2 - src/execution/dql/join/hash/right_join.rs | 6 - src/execution/dql/join/hash_join.rs | 3 - src/execution/dql/scalar_apply.rs | 48 +++++++- src/execution/dql/scalar_subquery.rs | 15 ++- src/execution/dql/top_k.rs | 86 +++++++++++++- tests/slt/join.slt | 132 ++++++++++++++++++++++ tests/slt/order_by.slt | 41 +++++++ tests/slt/subquery.slt | 51 +++++++++ 12 files changed, 378 insertions(+), 57 deletions(-) diff --git a/src/execution/dql/join/hash/full_join.rs b/src/execution/dql/join/hash/full_join.rs index 64a96948..1fa9c4d7 100644 --- a/src/execution/dql/join/hash/full_join.rs +++ b/src/execution/dql/join/hash/full_join.rs @@ -63,7 +63,7 @@ impl JoinProbeState for FullJoinState { ))); }; - if probe_state.index < build_state.tuples.len() { + while probe_state.index < build_state.tuples.len() { let (i, Tuple { values, pk }) = &build_state.tuples[probe_state.index]; probe_state.index += 1; @@ -71,12 +71,7 @@ impl JoinProbeState for FullJoinState { let full_values = SplitTupleRef::from_slices(values, &probe_state.probe_tuple.values); if !filter(&full_values, filter_expr, plan_arena)? { - probe_state.has_filtered = true; - self.bits.insert(*i); - return Ok(Some(Self::full_right_row( - self.left_schema_len, - &probe_state.probe_tuple, - ))); + continue; } } let full_values = Vec::from_iter( @@ -85,17 +80,22 @@ impl JoinProbeState for FullJoinState { .chain(probe_state.probe_tuple.values.iter()) .cloned(), ); - build_state.is_used = true; - build_state.has_filted = probe_state.has_filtered; + self.bits.insert(*i); + probe_state.produced = true; return Ok(Some(Tuple::new( pk.as_ref().or(probe_state.probe_tuple.pk.as_ref()).cloned(), full_values, ))); } - build_state.is_used = !probe_state.has_filtered; - build_state.has_filted = probe_state.has_filtered; probe_state.finished = true; + if !probe_state.produced && !probe_state.emitted_unmatched { + probe_state.emitted_unmatched = true; + return Ok(Some(Self::full_right_row( + self.left_schema_len, + &probe_state.probe_tuple, + ))); + } Ok(None) } @@ -108,12 +108,9 @@ impl JoinProbeState for FullJoinState { let full_schema_len = self.right_schema_len + self.left_schema_len; loop { - if let Some(LeftDropTuples { - tuples, has_filted, .. - }) = left_drop_state.current.as_mut() - { + if let Some(LeftDropTuples { tuples }) = left_drop_state.current.as_mut() { for (i, mut left_tuple) in tuples.by_ref() { - if !self.bits.contains(i) && *has_filted { + if self.bits.contains(i) { continue; } left_tuple.values.resize(full_schema_len, DataValue::Null); @@ -126,12 +123,8 @@ impl JoinProbeState for FullJoinState { return Ok(None); }; - if state.is_used { - continue; - } left_drop_state.current = Some(LeftDropTuples { tuples: state.tuples.into_iter(), - has_filted: state.has_filted, }); } } diff --git a/src/execution/dql/join/hash/inner_join.rs b/src/execution/dql/join/hash/inner_join.rs index 95066e84..9df19866 100644 --- a/src/execution/dql/join/hash/inner_join.rs +++ b/src/execution/dql/join/hash/inner_join.rs @@ -39,7 +39,6 @@ impl JoinProbeState for InnerJoinState { return Ok(None); }; - build_state.is_used = true; while probe_state.index < build_state.tuples.len() { let (_, Tuple { values, pk }) = &build_state.tuples[probe_state.index]; probe_state.index += 1; diff --git a/src/execution/dql/join/hash/left_join.rs b/src/execution/dql/join/hash/left_join.rs index dd99951d..5d4ebc82 100644 --- a/src/execution/dql/join/hash/left_join.rs +++ b/src/execution/dql/join/hash/left_join.rs @@ -55,8 +55,6 @@ impl JoinProbeState for LeftJoinState { let full_values = SplitTupleRef::from_slices(values, &probe_state.probe_tuple.values); if !filter(&full_values, filter_expr, plan_arena)? { - probe_state.has_filtered = true; - self.bits.insert(*i); continue; } } @@ -66,15 +64,13 @@ impl JoinProbeState for LeftJoinState { .chain(probe_state.probe_tuple.values.iter()) .cloned(), ); - build_state.is_used = true; + self.bits.insert(*i); return Ok(Some(Tuple::new( pk.as_ref().or(probe_state.probe_tuple.pk.as_ref()).cloned(), full_values, ))); } - build_state.is_used = !probe_state.has_filtered; - build_state.has_filted = probe_state.has_filtered; probe_state.finished = true; Ok(None) } @@ -88,12 +84,9 @@ impl JoinProbeState for LeftJoinState { let full_schema_len = self.right_schema_len + self.left_schema_len; loop { - if let Some(LeftDropTuples { - tuples, has_filted, .. - }) = left_drop_state.current.as_mut() - { + if let Some(LeftDropTuples { tuples }) = left_drop_state.current.as_mut() { for (i, mut left_tuple) in tuples.by_ref() { - if !self.bits.contains(i) && *has_filted { + if self.bits.contains(i) { continue; } left_tuple.values.resize(full_schema_len, DataValue::Null); @@ -106,12 +99,8 @@ impl JoinProbeState for LeftJoinState { return Ok(None); }; - if state.is_used { - continue; - } left_drop_state.current = Some(LeftDropTuples { tuples: state.tuples.into_iter(), - has_filted: state.has_filted, }); } } diff --git a/src/execution/dql/join/hash/mod.rs b/src/execution/dql/join/hash/mod.rs index a93478f1..228cd36f 100644 --- a/src/execution/dql/join/hash/mod.rs +++ b/src/execution/dql/join/hash/mod.rs @@ -34,7 +34,6 @@ pub(crate) struct ProbeState { pub(crate) is_keys_has_null: bool, pub(crate) probe_tuple: Tuple, pub(crate) index: usize, - pub(crate) has_filtered: bool, pub(crate) produced: bool, pub(crate) finished: bool, pub(crate) emitted_unmatched: bool, @@ -47,7 +46,6 @@ pub(crate) struct LeftDropState { pub(crate) struct LeftDropTuples { pub(crate) tuples: std::vec::IntoIter<(usize, Tuple)>, - pub(crate) has_filted: bool, } pub(crate) trait JoinProbeState { diff --git a/src/execution/dql/join/hash/right_join.rs b/src/execution/dql/join/hash/right_join.rs index 0f40c4cc..0c0f3276 100644 --- a/src/execution/dql/join/hash/right_join.rs +++ b/src/execution/dql/join/hash/right_join.rs @@ -66,7 +66,6 @@ impl JoinProbeState for RightJoinState { let full_values = SplitTupleRef::from_slices(values, &probe_state.probe_tuple.values); if !filter(&full_values, filter_expr, plan_arena)? { - probe_state.has_filtered = true; continue; } } @@ -77,17 +76,12 @@ impl JoinProbeState for RightJoinState { .cloned(), ); probe_state.produced = true; - build_state.is_used = true; - build_state.has_filted = probe_state.has_filtered; return Ok(Some(Tuple::new( pk.as_ref().or(probe_state.probe_tuple.pk.as_ref()).cloned(), full_values, ))); } - build_state.is_used = probe_state.produced; - build_state.has_filted = probe_state.has_filtered; - if !probe_state.produced && !probe_state.emitted_unmatched { probe_state.emitted_unmatched = true; probe_state.finished = true; diff --git a/src/execution/dql/join/hash_join.rs b/src/execution/dql/join/hash_join.rs index 222ffac8..94b581ef 100644 --- a/src/execution/dql/join/hash_join.rs +++ b/src/execution/dql/join/hash_join.rs @@ -223,8 +223,6 @@ impl HashJoin { #[derive(Default, Debug)] pub(crate) struct BuildState { pub(crate) tuples: Vec<(usize, Tuple)>, - pub(crate) is_used: bool, - pub(crate) has_filted: bool, } impl<'a, T: Transaction + 'a> ReadExecutor<'a, T> for HashJoin { @@ -290,7 +288,6 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for HashJoin { is_keys_has_null: probe_buf.iter().any(DataValue::is_null), probe_tuple: tuple, index: 0, - has_filtered: false, produced: false, finished: false, emitted_unmatched: false, diff --git a/src/execution/dql/scalar_apply.rs b/src/execution/dql/scalar_apply.rs index f0c74c2e..b506d68f 100644 --- a/src/execution/dql/scalar_apply.rs +++ b/src/execution/dql/scalar_apply.rs @@ -84,7 +84,13 @@ impl ScalarApply { "scalar apply right input returned no rows".to_string(), )); } - Ok(cached_right.insert(arena.materialize_tuple())) + let first = arena.materialize_tuple(); + if arena.next_tuple(right_input, plan_arena)? { + return Err(DatabaseError::InvalidValue( + "scalar apply right input returned more than one row".to_string(), + )); + } + Ok(cached_right.insert(first)) } } } @@ -190,6 +196,46 @@ mod tests { Ok(()) } + #[test] + fn scalar_subquery_checks_cardinality_on_second_call() -> Result<(), DatabaseError> { + for rows in [ + vec![], + vec![vec![DataValue::Int32(7)]], + vec![vec![DataValue::Int32(7)], vec![DataValue::Int32(8)]], + ] { + let row_count = rows.len(); + let table_arena = crate::planner::TableArenaCell::default(); + let mut plan_arena = crate::planner::PlanArena::new(&table_arena); + let mut input = build_values(&mut plan_arena, "right_c1", rows); + input.populate_output_schema_recursive(&mut plan_arena); + let (table_cache, view_cache, meta_cache, _temp_dir, storage) = build_test_storage()?; + let transaction = storage.transaction()?; + let mut executor = + execute_input::<_, crate::execution::dql::scalar_subquery::ScalarSubquery>( + (&ScalarSubqueryOperator, &input), + crate::execution::empty_context(&table_cache, &view_cache, &meta_cache), + plan_arena, + &transaction, + ); + let expected = if row_count == 0 { + DataValue::Null + } else { + DataValue::Int32(7) + }; + assert_eq!(executor.next_tuple()?.unwrap().values, vec![expected]); + if row_count > 1 { + assert!( + matches!(executor.next_tuple(), Err(DatabaseError::InvalidValue(message)) + if message == "scalar subquery returned more than one row") + ); + } else { + assert!(executor.next_tuple()?.is_none()); + assert!(executor.next_tuple()?.is_none()); + } + } + Ok(()) + } + #[test] fn scalar_apply_repeats_null_scalar_result_for_each_left_row() -> Result<(), DatabaseError> { let table_arena = crate::planner::TableArenaCell::default(); diff --git a/src/execution/dql/scalar_subquery.rs b/src/execution/dql/scalar_subquery.rs index f59b5115..822f81cb 100644 --- a/src/execution/dql/scalar_subquery.rs +++ b/src/execution/dql/scalar_subquery.rs @@ -54,14 +54,19 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for ScalarSubquery { arena: &mut ExecArena<'a, T>, plan_arena: &mut (dyn MetaArena + 'a), ) -> Result<(), DatabaseError> { + let has_next = arena.next_tuple(self.input, plan_arena)?; if self.returned { + if has_next { + return Err(DatabaseError::InvalidValue( + "scalar subquery returned more than one row".to_string(), + )); + } arena.finish(); return Ok(()); } self.returned = true; - let has_first = arena.next_tuple(self.input, plan_arena)?; - if !has_first { + if !has_next { let output = arena.result_tuple_mut(); output.pk = None; output.values.clear(); @@ -72,12 +77,6 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for ScalarSubquery { return Ok(()); } - if arena.next_tuple(self.input, plan_arena)? { - return Err(DatabaseError::InvalidValue( - "scalar subquery returned more than one row".to_string(), - )); - } - arena.resume(); Ok(()) } diff --git a/src/execution/dql/top_k.rs b/src/execution/dql/top_k.rs index df4eef8a..55cc9168 100644 --- a/src/execution/dql/top_k.rs +++ b/src/execution/dql/top_k.rs @@ -29,15 +29,26 @@ use std::cmp::Ordering; use std::collections::{btree_set::IntoIter as BTreeSetIntoIter, BTreeSet}; use std::mem::transmute; -#[derive(Eq, PartialEq, Debug)] +#[derive(Debug)] struct CmpItem<'a> { key: BumpVec<'a, u8>, + sequence: usize, tuple: Tuple, } +impl PartialEq for CmpItem<'_> { + fn eq(&self, other: &Self) -> bool { + self.key == other.key && self.sequence == other.sequence + } +} + +impl Eq for CmpItem<'_> {} + impl Ord for CmpItem<'_> { fn cmp(&self, other: &Self) -> Ordering { - self.key.cmp(&other.key).then_with(|| Ordering::Greater) + self.key + .cmp(&other.key) + .then_with(|| self.sequence.cmp(&other.sequence)) } } @@ -54,6 +65,7 @@ fn top_sort<'a>( heap: &mut BTreeSet>, tuple: Tuple, keep_count: usize, + sequence: usize, plan_arena: &(dyn MetaArena + '_), ) -> Result<(), DatabaseError> { let mut full_key = BumpBytes::new_in(arena); @@ -79,6 +91,7 @@ fn top_sort<'a>( if heap.len() < keep_count { heap.insert(CmpItem { key: full_key, + sequence, tuple, }); } else if let Some(cmp_item) = heap.last() { @@ -86,6 +99,7 @@ fn top_sort<'a>( heap.pop_last(); heap.insert(CmpItem { key: full_key, + sequence, tuple, }); } @@ -142,6 +156,7 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for TopK<'a> { #[allow(clippy::mutable_key_type)] let mut set = BTreeSet::new(); + let mut sequence = 0; while arena.next_tuple(self.input, plan_arena)? { top_sort( &self.arena, @@ -149,8 +164,10 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for TopK<'a> { &mut set, arena.materialize_tuple(), keep_count, + sequence, plan_arena, )?; + sequence += 1; } let offset = self.offset.unwrap_or(0); @@ -188,6 +205,35 @@ mod test { use bumpalo::Bump; use std::collections::BTreeSet; + #[test] + fn top_k_equal_keys_have_consistent_ordering() { + let arena = Bump::new(); + let make_item = |sequence, value| { + let mut key = crate::storage::table_codec::BumpBytes::new_in(&arena); + key.push(1); + CmpItem { + key, + sequence, + tuple: Tuple::new(None, vec![DataValue::Int32(value)]), + } + }; + let first = make_item(0, 10); + let same_key_and_sequence = make_item(0, 20); + let second = make_item(1, 30); + assert_eq!(first.cmp(&first), std::cmp::Ordering::Equal); + assert_eq!(first, same_key_and_sequence); + assert_eq!(first.cmp(&same_key_and_sequence), std::cmp::Ordering::Equal); + assert_eq!(first.cmp(&second), std::cmp::Ordering::Less); + assert_eq!(second.cmp(&first), std::cmp::Ordering::Greater); + + let mut set = BTreeSet::new(); + assert!(set.insert(second)); + assert!(set.insert(first)); + assert!(!set.insert(same_key_and_sequence)); + assert_eq!(set.pop_first().unwrap().sequence, 0); + assert_eq!(set.pop_first().unwrap().sequence, 1); + } + #[test] fn test_top_k_sort() -> Result<(), DatabaseError> { let table_arena = crate::planner::TableArenaCell::default(); @@ -267,6 +313,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Null]), 2, + 0, &plan_arena, )?; top_sort( @@ -275,6 +322,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Int32(0)]), 2, + 1, &plan_arena, )?; top_sort( @@ -283,6 +331,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Int32(1)]), 2, + 2, &plan_arena, )?; fn_asc_and_nulls_first_eq(indices); @@ -295,6 +344,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Null]), 2, + 3, &plan_arena, )?; top_sort( @@ -303,6 +353,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Int32(0)]), 2, + 4, &plan_arena, )?; top_sort( @@ -311,6 +362,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Int32(1)]), 2, + 5, &plan_arena, )?; fn_asc_and_nulls_last_eq(indices); @@ -323,6 +375,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Null]), 2, + 6, &plan_arena, )?; top_sort( @@ -331,6 +384,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Int32(0)]), 2, + 7, &plan_arena, )?; top_sort( @@ -339,6 +393,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Int32(1)]), 2, + 8, &plan_arena, )?; fn_desc_and_nulls_first_eq(indices); @@ -351,6 +406,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Null]), 2, + 9, &plan_arena, )?; top_sort( @@ -359,6 +415,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Int32(0)]), 2, + 10, &plan_arena, )?; top_sort( @@ -367,6 +424,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Int32(1)]), 2, + 11, &plan_arena, )?; fn_desc_and_nulls_last_eq(indices); @@ -556,6 +614,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Null, DataValue::Null]), 4, + 12, &plan_arena, )?; top_sort( @@ -564,6 +623,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Int32(0), DataValue::Null]), 4, + 13, &plan_arena, )?; top_sort( @@ -572,6 +632,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Int32(1), DataValue::Null]), 4, + 14, &plan_arena, )?; top_sort( @@ -580,6 +641,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Null, DataValue::Int32(0)]), 4, + 15, &plan_arena, )?; top_sort( @@ -588,6 +650,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Int32(0), DataValue::Int32(0)]), 4, + 16, &plan_arena, )?; top_sort( @@ -596,6 +659,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Int32(1), DataValue::Int32(0)]), 4, + 17, &plan_arena, )?; fn_asc_1_and_nulls_first_1_and_asc_2_and_nulls_first_2_eq(indices); @@ -608,6 +672,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Null, DataValue::Null]), 4, + 18, &plan_arena, )?; top_sort( @@ -616,6 +681,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Int32(0), DataValue::Null]), 4, + 19, &plan_arena, )?; top_sort( @@ -624,6 +690,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Int32(1), DataValue::Null]), 4, + 20, &plan_arena, )?; top_sort( @@ -632,6 +699,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Null, DataValue::Int32(0)]), 4, + 21, &plan_arena, )?; top_sort( @@ -640,6 +708,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Int32(0), DataValue::Int32(0)]), 4, + 22, &plan_arena, )?; top_sort( @@ -648,6 +717,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Int32(1), DataValue::Int32(0)]), 4, + 23, &plan_arena, )?; fn_asc_1_and_nulls_last_1_and_asc_2_and_nulls_first_2_eq(indices); @@ -660,6 +730,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Null, DataValue::Null]), 4, + 24, &plan_arena, )?; top_sort( @@ -668,6 +739,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Int32(0), DataValue::Null]), 4, + 25, &plan_arena, )?; top_sort( @@ -676,6 +748,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Int32(1), DataValue::Null]), 4, + 26, &plan_arena, )?; top_sort( @@ -684,6 +757,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Null, DataValue::Int32(0)]), 4, + 27, &plan_arena, )?; top_sort( @@ -692,6 +766,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Int32(0), DataValue::Int32(0)]), 4, + 28, &plan_arena, )?; top_sort( @@ -700,6 +775,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Int32(1), DataValue::Int32(0)]), 4, + 29, &plan_arena, )?; fn_desc_1_and_nulls_first_1_and_asc_2_and_nulls_first_2_eq(indices); @@ -712,6 +788,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Null, DataValue::Null]), 4, + 30, &plan_arena, )?; top_sort( @@ -720,6 +797,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Int32(0), DataValue::Null]), 4, + 31, &plan_arena, )?; top_sort( @@ -728,6 +806,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Int32(1), DataValue::Null]), 4, + 32, &plan_arena, )?; top_sort( @@ -736,6 +815,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Null, DataValue::Int32(0)]), 4, + 33, &plan_arena, )?; top_sort( @@ -744,6 +824,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Int32(0), DataValue::Int32(0)]), 4, + 34, &plan_arena, )?; top_sort( @@ -752,6 +833,7 @@ mod test { &mut indices, Tuple::new(None, vec![DataValue::Int32(1), DataValue::Int32(0)]), 4, + 35, &plan_arena, )?; fn_desc_1_and_nulls_last_1_and_asc_2_and_nulls_first_2_eq(indices); diff --git a/tests/slt/join.slt b/tests/slt/join.slt index 5496a260..a1ef0c12 100644 --- a/tests/slt/join.slt +++ b/tests/slt/join.slt @@ -325,3 +325,135 @@ drop table lookup_inner; statement ok drop table lookup_outer; + +# A build row matched by any probe must never be emitted again with NULL padding. +# A probe is unmatched only if no build row passes the complete ON condition. +statement ok +create table hj_left (id int primary key, k int); + +statement ok +create table hj_right (id int primary key, k int); + +statement ok +insert into hj_left values (1, 5), (4, 5), (6, 5), (7, null), (8, 8); + +statement ok +insert into hj_right values (3, 5), (5, 5), (9, 5), (10, null), (11, 11), (0, 5); + +query II rowsort +select coalesce(a.id, -1), coalesce(b.id, -1) from hj_left a full join hj_right b on a.k = b.k and a.id < b.id; +---- +-1 0 +-1 10 +-1 11 +1 3 +1 5 +1 9 +4 5 +4 9 +6 9 +7 -1 +8 -1 + +query II rowsort +select coalesce(a.id, -1), coalesce(b.id, -1) from hj_left a left join hj_right b on a.k = b.k and a.id < b.id; +---- +1 3 +1 5 +1 9 +4 5 +4 9 +6 9 +7 -1 +8 -1 + +query II rowsort +select coalesce(a.id, -1), coalesce(b.id, -1) from hj_left a right join hj_right b on a.k = b.k and a.id < b.id; +---- +-1 0 +-1 10 +-1 11 +1 3 +1 5 +1 9 +4 5 +4 9 +6 9 + +# Multiple left rows in one bucket: one matches, another must be padded. +query II rowsort +select coalesce(a.id, -1), coalesce(b.id, -1) from hj_left a full join hj_right b on a.k = b.k and a.id < b.id and b.id = 3; +---- +-1 0 +-1 10 +-1 11 +-1 5 +-1 9 +1 3 +4 -1 +6 -1 +7 -1 +8 -1 + +query II rowsort +select coalesce(a.id, -1), coalesce(b.id, -1) from hj_left a left join hj_right b on a.k = b.k and a.id < b.id and b.id = 3; +---- +1 3 +4 -1 +6 -1 +7 -1 +8 -1 + +# NULL residual predicate rejects every candidate, but each unmatched row emits once. +query II rowsort +select coalesce(a.id, -1), coalesce(b.id, -1) from hj_left a full join hj_right b on a.k = b.k and a.id < b.k + null; +---- +-1 0 +-1 10 +-1 11 +-1 3 +-1 5 +-1 9 +1 -1 +4 -1 +6 -1 +7 -1 +8 -1 + +# No residual predicate: ordinary equal-key matches and NULL keys. +query I +select count(*) from hj_left a full join hj_right b on a.k = b.k; +---- +16 + +# Empty probe/build inputs must still emit each retained-side row once. +statement ok +create table hj_empty (id int primary key, k int); + +query II rowsort +select coalesce(a.id, -1), coalesce(b.id, -1) from hj_left a full join hj_empty b on a.k = b.k and a.id < b.id; +---- +1 -1 +4 -1 +6 -1 +7 -1 +8 -1 + +query II rowsort +select coalesce(a.id, -1), coalesce(b.id, -1) from hj_empty a full join hj_right b on a.k = b.k and a.id < b.id; +---- +-1 0 +-1 10 +-1 11 +-1 3 +-1 5 +-1 9 + +statement ok +drop table hj_empty; + +statement ok +drop table hj_right; + +statement ok +drop table hj_left; diff --git a/tests/slt/order_by.slt b/tests/slt/order_by.slt index 5ac306f5..2ac16bea 100644 --- a/tests/slt/order_by.slt +++ b/tests/slt/order_by.slt @@ -106,3 +106,44 @@ select v1 as a from t order by a statement ok drop table t + +# TopK retains duplicate keys and handles replacement followed by equal keys. +statement ok +create table tk_duplicates (id int primary key, k int); + +statement ok +insert into tk_duplicates values (1, 9), (2, 1), (3, 1), (4, 1), (5, 0), (6, 0), (7, null), (8, null); + +query I +select k from tk_duplicates order by k asc nulls last limit 4; +---- +0 +0 +1 +1 + +query I +select k from tk_duplicates order by k asc nulls last limit 3 offset 1; +---- +0 +1 +1 + +query I +select k from tk_duplicates order by k desc nulls last limit 4; +---- +9 +1 +1 +1 + +query I +select coalesce(k, -1) from tk_duplicates order by k asc nulls first limit 4; +---- +-1 +-1 +0 +0 + +statement ok +drop table tk_duplicates; diff --git a/tests/slt/subquery.slt b/tests/slt/subquery.slt index ee69c26d..1c2b958c 100644 --- a/tests/slt/subquery.slt +++ b/tests/slt/subquery.slt @@ -499,3 +499,54 @@ drop table users; statement ok drop table orders; + +# ScalarApply must preserve its cached first result while validating scalar cardinality. +statement ok +create table scalar_eof (id int primary key, v int, keep_row boolean); + +statement ok +insert into scalar_eof values (1, 10, true), (2, 20, false), (3, 30, false); + +query I +select (select v from scalar_eof where keep_row = true); +---- +10 + +query I +select (select v from scalar_eof where keep_row = true order by v desc); +---- +10 + +query I +select coalesce((select v from scalar_eof where id < 0), -1); +---- +-1 + +statement error +select (select v from scalar_eof where id <= 2); + +# Reuse the validated scalar result for every outer row, even when checking EOF +# scans later rows that fail the inner filter. +query II +select outer_row.id, (select v from scalar_eof where keep_row = true) +from scalar_eof outer_row order by outer_row.id; +---- +1 10 +2 10 +3 10 + +query II +select outer_row.id, (select v from scalar_eof where id < 0) +from scalar_eof outer_row order by outer_row.id; +---- +1 null +2 null +3 null + +# Validate before accepting the cached scalar, including a single outer-row limit. +statement error +select outer_row.id, (select v from scalar_eof where id <= 2) +from scalar_eof outer_row order by outer_row.id limit 1; + +statement ok +drop table scalar_eof; From 4dcd1118175bf75155cad855b4f127ce0034d495 Mon Sep 17 00:00:00 2001 From: kould Date: Sun, 4 Oct 2026 07:33:14 +0800 Subject: [PATCH 2/2] perf: simplify executor tuple ownership and reuse key buffers --- src/catalog/table.rs | 3 + src/db.rs | 4 +- src/db/prepared.rs | 4 +- src/execution/ddl/add_column.rs | 3 +- src/execution/ddl/change_column.rs | 5 +- src/execution/ddl/create_index.rs | 11 +- src/execution/dml/analyze.rs | 11 +- src/execution/dml/copy_from_file.rs | 9 +- src/execution/dml/copy_to_file.rs | 9 +- src/execution/dml/delete.rs | 13 +- src/execution/dml/insert.rs | 33 ++-- src/execution/dml/update.rs | 49 +++--- src/execution/dql/aggregate/count.rs | 3 +- src/execution/dql/aggregate/hash_agg.rs | 2 +- src/execution/dql/join/hash/full_join.rs | 27 ++- src/execution/dql/join/hash/right_join.rs | 12 +- src/execution/dql/join/hash_join.rs | 3 +- src/execution/dql/join/nested_loop_join.rs | 183 +++++++------------- src/execution/dql/mark_apply.rs | 22 +-- src/execution/dql/recursive_cte.rs | 2 +- src/execution/dql/set_membership.rs | 9 +- src/execution/dql/sort.rs | 2 +- src/execution/dql/top_k.rs | 184 ++++++++++----------- src/execution/dql/window.rs | 9 +- src/execution/mod.rs | 20 ++- src/planner/mod.rs | 22 ++- src/types/tuple_builder.rs | 20 ++- 27 files changed, 310 insertions(+), 364 deletions(-) diff --git a/src/catalog/table.rs b/src/catalog/table.rs index d77f4512..ac523df1 100644 --- a/src/catalog/table.rs +++ b/src/catalog/table.rs @@ -125,6 +125,9 @@ impl TableCatalog { ) -> Result, DatabaseError> { let index_metas = self .indexes() + .filter(|index_meta| { + !matches!(arena.index(**index_meta).ty, IndexType::PrimaryKey { .. }) + }) .map(|index_meta| { Ok(( *index_meta, diff --git a/src/db.rs b/src/db.rs index 24c10588..4e58cf54 100644 --- a/src/db.rs +++ b/src/db.rs @@ -589,7 +589,7 @@ impl State { let keeper = PlanKeeper::new(plan); let plan = keeper.plan(); let schema = plan.read_schema().clone(); - let mut arena = ExecArena::new(); + let mut arena = ExecArena::with_capacity(plan.exec_capacity_hint()); arena.set_statement_stamp(statement_stamp(transaction, plan)?); let read_context = ExecutionContext::new( &self.table_cache, @@ -665,7 +665,7 @@ impl State { let keeper = PlanKeeper::new(PlanInput::Owned(plan)); let plan = keeper.plan(); let schema = plan.read_schema().clone(); - let mut arena = ExecArena::new(); + let mut arena = ExecArena::with_capacity(plan.exec_capacity_hint()); arena.set_statement_stamp(statement_stamp(transaction, plan)?); let cache = ExecutionContext::new( table_cache, diff --git a/src/db/prepared.rs b/src/db/prepared.rs index 4a4d0981..380f61c1 100644 --- a/src/db/prepared.rs +++ b/src/db/prepared.rs @@ -323,7 +323,7 @@ mod tests { } let mut execution_arena = crate::execution::ExecArena::< ::TransactionType<'_>, - >::new(); + >::with_capacity(0); let executor = crate::execution::dql::index_scan::IndexScan::new( scan_op, info, @@ -551,7 +551,7 @@ mod tests { }; let mut execution_arena = crate::execution::ExecArena::< ::TransactionType<'_>, - >::new(); + >::with_capacity(0); let ranges = crate::execution::dql::index_scan::IndexScan::new( scan_op, info, diff --git a/src/execution/ddl/add_column.rs b/src/execution/ddl/add_column.rs index 3417db45..878eff21 100644 --- a/src/execution/ddl/add_column.rs +++ b/src/execution/ddl/add_column.rs @@ -95,7 +95,6 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for AddColumn<'a> { (unique_meta, DDLApply::upsert_table(table, false)) }; arena.push_ddl_apply(apply); - let default_for_index = default_value.clone(); let mut state = arena.local_state(plan_arena); let plan_arena = state.plan_arena; @@ -128,7 +127,7 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for AddColumn<'a> { |transaction, table_codec, tuple| { if let (Some(unique_index_id), Some(value), Some(tuple_id)) = ( unique_index_id.as_ref(), - default_for_index.as_ref(), + default_value.as_ref(), tuple.pk.as_ref(), ) { let index = Index::new( diff --git a/src/execution/ddl/change_column.rs b/src/execution/ddl/change_column.rs index 8bca2787..7405111b 100644 --- a/src/execution/ddl/change_column.rs +++ b/src/execution/ddl/change_column.rs @@ -21,6 +21,7 @@ use crate::iter_ext::Itertools; use crate::planner::operator::alter_table::change_column::{ChangeColumnOperator, NotNullChange}; use crate::planner::MetaArena; use crate::storage::Transaction; +use crate::types::value::DataValue; pub struct ChangeColumn<'a> { op: &'a ChangeColumnOperator, @@ -127,8 +128,8 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for ChangeColumn<'a> { }) }, |tuple| { - tuple.values[column_index] = - tuple.values[column_index].clone().cast(&target_data_type)?; + let value = std::mem::replace(&mut tuple.values[column_index], DataValue::Null); + tuple.values[column_index] = value.cast(&target_data_type)?; if needs_not_null_validation && tuple.values[column_index].is_null() { return Err(DatabaseError::not_null_column(target_column_name.clone())); } diff --git a/src/execution/ddl/create_index.rs b/src/execution/ddl/create_index.rs index 1c2ed39a..7705f1ec 100644 --- a/src/execution/ddl/create_index.rs +++ b/src/execution/ddl/create_index.rs @@ -116,15 +116,16 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for CreateIndex<'a> { }; while arena.next_tuple(self.input, plan_arena)? { - let Some(tuple_pk) = arena.result_tuple().pk.clone() else { + if arena.result_tuple().pk.is_none() { continue; - }; + } arena.rewrite(&column_exprs, plan_arena, None)?; { let mut state = arena.local_state(plan_arena); - let (values, transaction, table_codec) = state.index_values_transaction_codec_mut(); - let index = Index::new(index_id, values, *ty); - transaction.add_index(table_codec, table_name.as_ref(), index, &tuple_pk)?; + let (tuple, transaction, table_codec) = state.tuple_transaction_codec_mut(); + let tuple_pk = tuple.pk.as_ref().ok_or(DatabaseError::PrimaryKeyNotFound)?; + let index = Index::new(index_id, &tuple.values, *ty); + transaction.add_index(table_codec, table_name.as_ref(), index, tuple_pk)?; } } diff --git a/src/execution/dml/analyze.rs b/src/execution/dml/analyze.rs index f53fb017..213d87a1 100644 --- a/src/execution/dml/analyze.rs +++ b/src/execution/dml/analyze.rs @@ -107,11 +107,12 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for Analyze<'a> { while arena.next_tuple(input, plan_arena)? { let tuple = arena.materialize_tuple(); for State { exprs, builder, .. } in builders.iter_mut() { - arena.rewrite(exprs, plan_arena, Some(&tuple))?; - let key = arena.materialize_tuple().values; - let value = match <[DataValue; 1]>::try_from(key) { - Ok([value]) => value, - Err(key) => DataValue::Tuple(key), + let value = match exprs.as_slice() { + [expr] => expr.eval(plan_arena, Some(&tuple))?.into_owned(), + _ => { + arena.rewrite(exprs, plan_arena, Some(&tuple))?; + DataValue::Tuple(arena.materialize_tuple().values) + } }; builder.append(value)?; } diff --git a/src/execution/dml/copy_from_file.rs b/src/execution/dml/copy_from_file.rs index 50f88dc6..c1d25aa8 100644 --- a/src/execution/dml/copy_from_file.rs +++ b/src/execution/dml/copy_from_file.rs @@ -21,6 +21,7 @@ use crate::iter_ext::Itertools; use crate::planner::operator::copy_from_file::CopyFromFileOperator; use crate::planner::MetaArena; use crate::storage::Transaction; +use crate::types::tuple::Tuple; use crate::types::tuple_builder::TupleBuilder; use std::fs::File; use std::io::BufReader; @@ -85,16 +86,16 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for CopyFromFile<'a> { let column_count = op.schema_ref.len(); let tuple_builder = TupleBuilder::new(column_types, Some(table.primary_key_indices())); - for record in reader.records() { - let record = record?; - + let mut record = csv::StringRecord::new(); + let mut chunk = Tuple::new(None, Vec::with_capacity(column_count)); + while reader.read_record(&mut record)? { if !(record.len() == column_count || record.len() == column_count + 1 && record.get(column_count) == Some("")) { return Err(DatabaseError::MisMatch("columns", "values")); } - let chunk = tuple_builder.build_with_row(record.iter())?; + tuple_builder.build_with_row(&mut chunk, record.iter().take(column_count))?; let mut state = arena.local_state(plan_arena); let (transaction, table_codec) = state.transaction_codec_mut(); transaction.append_tuple(table_codec, &table_name, &chunk, &serializers, false)?; diff --git a/src/execution/dml/copy_to_file.rs b/src/execution/dml/copy_to_file.rs index c4511c42..07e02865 100644 --- a/src/execution/dml/copy_to_file.rs +++ b/src/execution/dml/copy_to_file.rs @@ -96,14 +96,7 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for CopyToFile<'a> { let mut writer = self.create_writer()?; while arena.next_tuple(input, plan_arena)? { - let tuple = arena.materialize_tuple(); - writer.write_record( - tuple - .values - .iter() - .map(|v| v.to_string()) - .collect::>(), - )?; + writer.write_record(arena.result_tuple().values.iter().map(|v| v.to_string()))?; } writer.flush().map_err(DatabaseError::from)?; diff --git a/src/execution/dml/delete.rs b/src/execution/dml/delete.rs index 3f25092b..ecdff0c9 100644 --- a/src/execution/dml/delete.rs +++ b/src/execution/dml/delete.rs @@ -21,7 +21,7 @@ use crate::planner::operator::delete::DeleteOperator; use crate::planner::LogicalPlan; use crate::planner::MetaArena; use crate::storage::Transaction; -use crate::types::index::Index; +use crate::types::index::{Index, IndexType}; use crate::types::tuple_builder::TupleBuilder; pub struct Delete<'a> { @@ -68,6 +68,12 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for Delete<'a> { .ok_or(DatabaseError::TableNotFound)?; table .indexes() + .filter(|index_meta| { + !matches!( + plan_arena.index(**index_meta).ty, + IndexType::PrimaryKey { .. } + ) + }) .map(|index_meta| { let index_meta = plan_arena.index(*index_meta); Ok(( @@ -83,11 +89,10 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for Delete<'a> { let mut deleted_count = 0; while arena.next_tuple(input, plan_arena)? { - let Some(tuple_id) = arena.result_tuple().pk.clone() else { + let tuple = arena.materialize_tuple(); + let Some(tuple_id) = tuple.pk.as_ref() else { continue; }; - - let tuple = arena.materialize_tuple(); for (index_id, index_ty, exprs) in index_templates.iter() { arena.rewrite(exprs, plan_arena, Some(&tuple))?; let mut state = arena.local_state(plan_arena); diff --git a/src/execution/dml/insert.rs b/src/execution/dml/insert.rs index 9a40b8cb..edc9e84d 100644 --- a/src/execution/dml/insert.rs +++ b/src/execution/dml/insert.rs @@ -110,6 +110,7 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for Insert<'a> { .map(|table| table.dml_snapshot(plan_arena)) .transpose()? }; + let mut inserted_count = 0; if let Some(table_snapshot) = table_snapshot { if table_snapshot.primary_key_indices.is_empty() { return Err(DatabaseError::not_null()); @@ -121,11 +122,10 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for Insert<'a> { .map(|column| plan_arena.column(*column).datatype().serializable()) .collect_vec(); let mut tuple = Tuple::new(None, Vec::with_capacity(table_snapshot.columns_len)); - let mut inserted_count = 0; while arena.next_tuple(input, plan_arena)? { let mut tuple_map = HashMap::with_capacity(self.input_schema.len()); - for (i, value) in arena.materialize_tuple().values.into_iter().enumerate() { + for (i, value) in arena.result_tuple_mut().values.drain(..).enumerate() { let column = plan_arena.column(self.input_schema[i]); tuple_map.insert(Self::column_key(column, self.is_mapping_by_name), value); } @@ -133,16 +133,13 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for Insert<'a> { tuple.values.clear(); for column in table_snapshot.columns.iter() { let column = plan_arena.column(*column); - let mut value = { - let mut value = - tuple_map.remove(&Self::column_key(column, self.is_mapping_by_name)); - - if value.is_none() { - value = column.default_value(plan_arena)?; - } - value.unwrap_or(DataValue::Null) - }; - value = value.cast(column.datatype())?; + let value = match tuple_map + .remove(&Self::column_key(column, self.is_mapping_by_name)) + { + Some(value) => value, + None => column.default_value(plan_arena)?.unwrap_or(DataValue::Null), + } + .cast(column.datatype())?; value.check_len(column.datatype())?; if value.is_null() && !column.nullable() { return Err(DatabaseError::not_null_column(column.name().to_string())); @@ -153,7 +150,6 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for Insert<'a> { table_snapshot.primary_key_indices, &tuple.values, )); - for (index_meta, exprs) in table_snapshot.index_metas.iter() { let index_meta = plan_arena.index(*index_meta); let tuple_id = tuple.pk.as_ref().ok_or(DatabaseError::PrimaryKeyNotFound)?; @@ -175,14 +171,9 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for Insert<'a> { )?; inserted_count += 1; } - - arena.produce_tuple(TupleBuilder::build_result(inserted_count.to_string())); - arena.resume(); - Ok(()) - } else { - arena.produce_tuple(TupleBuilder::build_result("0".to_string())); - arena.resume(); - Ok(()) } + arena.produce_tuple(TupleBuilder::build_result(inserted_count.to_string())); + arena.resume(); + Ok(()) } } diff --git a/src/execution/dml/update.rs b/src/execution/dml/update.rs index b993ebd1..3b052a9e 100644 --- a/src/execution/dml/update.rs +++ b/src/execution/dml/update.rs @@ -117,6 +117,8 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for Update<'a> { .map(|table| table.dml_snapshot(plan_arena)) .transpose()? }; + let mut updated_count = 0; + if let Some(table_snapshot) = table_snapshot { let updates_primary_key = table_snapshot.primary_key_indices.iter().any(|index| { table_snapshot @@ -128,11 +130,9 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for Update<'a> { let serializers = self .input_schema .iter() - .map(|column| plan_arena.column(*column).datatype().serializable()) + .map(|column: &ColumnRef| plan_arena.column(*column).datatype().serializable()) .collect_vec(); - let mut updated_count = 0; - while arena.next_tuple(input, plan_arena)? { let mut is_overwrite = true; @@ -141,10 +141,7 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for Update<'a> { continue; }; - let mut old_index_values = Vec::new(); - for (index_offset, (index_meta, exprs)) in - table_snapshot.index_metas.iter().enumerate() - { + for (index_meta, exprs) in table_snapshot.index_metas.iter() { let index_meta = plan_arena.index(*index_meta); if !Self::index_needs_update( index_meta, @@ -155,8 +152,11 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for Update<'a> { } arena.rewrite(exprs, plan_arena, Some(&tuple))?; - let values = arena.materialize_tuple().values; - old_index_values.push((index_offset, values)); + let mut state = arena.local_state(plan_arena); + let (values, transaction, table_codec) = + state.index_values_transaction_codec_mut(); + let old_index = Index::new(index_meta.id, values, index_meta.ty); + transaction.del_index(table_codec, self.table_name, &old_index, &old_pk)?; } for (i, column) in self.input_schema.iter().enumerate() { let Some(column_id) = plan_arena.column(*column).id() else { @@ -180,22 +180,20 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for Update<'a> { is_overwrite = false; } - for (index_offset, old_value) in old_index_values { - let (index_meta, exprs) = &table_snapshot.index_metas[index_offset]; + for (index_meta, exprs) in table_snapshot.index_metas.iter() { let index_meta = plan_arena.index(*index_meta); - let index_id = index_meta.id; - let index_ty = index_meta.ty; + if !Self::index_needs_update( + index_meta, + &updated_column_ids, + updates_primary_key, + ) { + continue; + } arena.rewrite(exprs, plan_arena, Some(&tuple))?; let mut state = arena.local_state(plan_arena); let (values, transaction, table_codec) = state.index_values_transaction_codec_mut(); - if !primary_key_changed && old_value == values { - continue; - } - - let old_index = Index::new(index_id, &old_value, index_ty); - let new_index = Index::new(index_id, values, index_ty); - transaction.del_index(table_codec, self.table_name, &old_index, &old_pk)?; + let new_index = Index::new(index_meta.id, values, index_meta.ty); transaction.add_index(table_codec, self.table_name, new_index, &new_pk)?; } @@ -214,14 +212,9 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for Update<'a> { })?; updated_count += 1; } - - arena.produce_tuple(TupleBuilder::build_result(updated_count.to_string())); - arena.resume(); - Ok(()) - } else { - arena.produce_tuple(TupleBuilder::build_result("0".to_string())); - arena.resume(); - Ok(()) } + arena.produce_tuple(TupleBuilder::build_result(updated_count.to_string())); + arena.resume(); + Ok(()) } } diff --git a/src/execution/dql/aggregate/count.rs b/src/execution/dql/aggregate/count.rs index df3cbbbb..38c0d743 100644 --- a/src/execution/dql/aggregate/count.rs +++ b/src/execution/dql/aggregate/count.rs @@ -66,7 +66,8 @@ impl DistinctCountAccumulator { impl Accumulator for DistinctCountAccumulator { fn update_value(&mut self, value: &DataValue) -> Result<(), DatabaseError> { - if !value.is_null() && self.distinct_values.insert(value.clone()) { + if !value.is_null() && !self.distinct_values.contains(value) { + self.distinct_values.insert(value.clone()); self.result = DataValue::Int32(self.distinct_values.len() as i32); } diff --git a/src/execution/dql/aggregate/hash_agg.rs b/src/execution/dql/aggregate/hash_agg.rs index 7e338d41..5028e343 100644 --- a/src/execution/dql/aggregate/hash_agg.rs +++ b/src/execution/dql/aggregate/hash_agg.rs @@ -92,7 +92,7 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for HashAggExecutor<'a> { } else { let mut accs = create_accumulators(self.agg_calls, plan_arena)?; update_accumulators(&mut accs, self.agg_calls, tuple, plan_arena)?; - group_hash_accs.insert(group_keys.clone(), accs); + group_hash_accs.insert(std::mem::take(&mut group_keys), accs); } } diff --git a/src/execution/dql/join/hash/full_join.rs b/src/execution/dql/join/hash/full_join.rs index 1fa9c4d7..ff9de82a 100644 --- a/src/execution/dql/join/hash/full_join.rs +++ b/src/execution/dql/join/hash/full_join.rs @@ -44,9 +44,9 @@ impl JoinProbeState for FullJoinState { } probe_state.emitted_unmatched = true; probe_state.finished = true; - return Ok(Some(Self::full_right_row( + return Ok(Some(Self::take_full_right_row( self.left_schema_len, - &probe_state.probe_tuple, + &mut probe_state.probe_tuple, ))); } @@ -57,9 +57,9 @@ impl JoinProbeState for FullJoinState { } probe_state.emitted_unmatched = true; probe_state.finished = true; - return Ok(Some(Self::full_right_row( + return Ok(Some(Self::take_full_right_row( self.left_schema_len, - &probe_state.probe_tuple, + &mut probe_state.probe_tuple, ))); }; @@ -91,9 +91,9 @@ impl JoinProbeState for FullJoinState { probe_state.finished = true; if !probe_state.produced && !probe_state.emitted_unmatched { probe_state.emitted_unmatched = true; - return Ok(Some(Self::full_right_row( + return Ok(Some(Self::take_full_right_row( self.left_schema_len, - &probe_state.probe_tuple, + &mut probe_state.probe_tuple, ))); } Ok(None) @@ -131,13 +131,12 @@ impl JoinProbeState for FullJoinState { } impl FullJoinState { - pub(crate) fn full_right_row(left_schema_len: usize, probe_tuple: &Tuple) -> Tuple { - let full_values = Vec::from_iter( - (0..left_schema_len) - .map(|_| DataValue::Null) - .chain(probe_tuple.values.iter().cloned()), - ); - - Tuple::new(probe_tuple.pk.clone(), full_values) + pub(crate) fn take_full_right_row(left_schema_len: usize, probe_tuple: &mut Tuple) -> Tuple { + let mut tuple = std::mem::take(probe_tuple); + tuple + .values + .resize(tuple.values.len() + left_schema_len, DataValue::Null); + tuple.values.rotate_right(left_schema_len); + tuple } } diff --git a/src/execution/dql/join/hash/right_join.rs b/src/execution/dql/join/hash/right_join.rs index 0c0f3276..cee9034d 100644 --- a/src/execution/dql/join/hash/right_join.rs +++ b/src/execution/dql/join/hash/right_join.rs @@ -39,9 +39,9 @@ impl JoinProbeState for RightJoinState { } probe_state.emitted_unmatched = true; probe_state.finished = true; - return Ok(Some(FullJoinState::full_right_row( + return Ok(Some(FullJoinState::take_full_right_row( self.left_schema_len, - &probe_state.probe_tuple, + &mut probe_state.probe_tuple, ))); } @@ -52,9 +52,9 @@ impl JoinProbeState for RightJoinState { } probe_state.emitted_unmatched = true; probe_state.finished = true; - return Ok(Some(FullJoinState::full_right_row( + return Ok(Some(FullJoinState::take_full_right_row( self.left_schema_len, - &probe_state.probe_tuple, + &mut probe_state.probe_tuple, ))); }; @@ -85,9 +85,9 @@ impl JoinProbeState for RightJoinState { if !probe_state.produced && !probe_state.emitted_unmatched { probe_state.emitted_unmatched = true; probe_state.finished = true; - return Ok(Some(FullJoinState::full_right_row( + return Ok(Some(FullJoinState::take_full_right_row( self.left_schema_len, - &probe_state.probe_tuple, + &mut probe_state.probe_tuple, ))); } diff --git a/src/execution/dql/join/hash_join.rs b/src/execution/dql/join/hash_join.rs index 94b581ef..11d6bfcc 100644 --- a/src/execution/dql/join/hash_join.rs +++ b/src/execution/dql/join/hash_join.rs @@ -166,8 +166,9 @@ impl HashJoin { match build_map.get_mut(&build_buf) { None => { + let key = std::mem::replace(&mut build_buf, BumpVec::new_in(&self.bump)); build_map.insert( - Self::own_bump_vec(build_buf.clone()), + Self::own_bump_vec(key), BuildState { tuples: vec![(build_count, tuple)], ..Default::default() diff --git a/src/execution/dql/join/nested_loop_join.rs b/src/execution/dql/join/nested_loop_join.rs index 5c9385a2..fb24a828 100644 --- a/src/execution/dql/join/nested_loop_join.rs +++ b/src/execution/dql/join/nested_loop_join.rs @@ -23,7 +23,6 @@ use crate::execution::dql::join::RowBitmap; use crate::execution::{ build_read, ExecArena, ExecId, ExecNode, ExecutionContext, ExecutorNode, ReadExecutor, }; -use crate::iter_ext::Itertools; use crate::planner::operator::join::{JoinCondition, JoinOperator, JoinType}; use crate::planner::ExprRef; use crate::storage::Transaction; @@ -196,72 +195,37 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for NestedLoopJoin<'a> { mut right_bitmap, } => { while arena.next_tuple(active_left.right_input, plan_arena)? { - let right_tuple = arena.materialize_tuple(); + let right_tuple = arena.result_tuple(); let idx = active_left.right_index; active_left.right_index += 1; - let tuple = match ( - self.filter.as_ref(), - self.eq_cond.equals( - &active_left.left_tuple, - &right_tuple, - plan_arena, - )?, - ) { - (None, true) if matches!(self.ty, JoinType::RightOuter) => { - active_left.has_matched = true; - Self::emit_tuple( - &right_tuple, - &active_left.left_tuple, - self.ty, - true, - ) - } - (None, true) => { - active_left.has_matched = true; - Self::emit_tuple( - &active_left.left_tuple, - &right_tuple, - self.ty, - true, - ) - } - (Some(filter), true) => { - let values = if matches!(self.ty, JoinType::RightOuter) { - SplitTupleRef::new(&right_tuple, &active_left.left_tuple) - } else { - SplitTupleRef::new(&active_left.left_tuple, &right_tuple) - }; - let value = plan_arena - .expression(*filter) - .eval(plan_arena, Some(&values))?; - match &*value { - DataValue::Boolean(true) => { - let tuple = match self.ty { - JoinType::RightOuter => Self::emit_tuple( - &right_tuple, - &active_left.left_tuple, - self.ty, - true, - ), - _ => Self::emit_tuple( - &active_left.left_tuple, - &right_tuple, - self.ty, - true, - ), - }; - active_left.has_matched = true; - tuple - } - DataValue::Boolean(false) | DataValue::Null => None, - _ => return Err(DatabaseError::InvalidType), - } + if !self + .eq_cond + .equals(&active_left.left_tuple, right_tuple, plan_arena)? + { + continue; + } + if let Some(filter) = self.filter { + let values = if matches!(self.ty, JoinType::RightOuter) { + SplitTupleRef::new(right_tuple, &active_left.left_tuple) + } else { + SplitTupleRef::new(&active_left.left_tuple, right_tuple) + }; + let value = plan_arena + .expression(filter) + .eval(plan_arena, Some(&values))?; + match &*value { + DataValue::Boolean(true) => (), + DataValue::Boolean(false) | DataValue::Null => continue, + _ => return Err(DatabaseError::InvalidType), } - _ => None, - }; - - if let Some(tuple) = tuple { + } + active_left.has_matched = true; + if Self::emit_tuple( + &active_left.left_tuple, + arena.result_tuple_mut(), + self.ty, + ) { if matches!(self.ty, JoinType::Full) { if let Some(bits) = right_bitmap.as_mut() { bits.insert(idx); @@ -274,7 +238,7 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for NestedLoopJoin<'a> { active_left, right_bitmap, }; - arena.produce_tuple(tuple); + arena.resume(); return Ok(()); } } @@ -293,34 +257,22 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for NestedLoopJoin<'a> { } } let right_schema_len = self.eq_cond.right_len; - let tuple = match self.ty { + let should_emit = match self.ty { JoinType::LeftOuter | JoinType::RightOuter | JoinType::Full if !active_left.has_matched => { - let right_tuple = - Tuple::new(None, vec![DataValue::Null; right_schema_len]); - if matches!(self.ty, JoinType::RightOuter) { - Self::emit_tuple( - &right_tuple, - &active_left.left_tuple, - self.ty, - false, - ) - } else { - Self::emit_tuple( - &active_left.left_tuple, - &right_tuple, - self.ty, - false, - ) - } + let right_tuple = arena.result_tuple_mut(); + right_tuple.pk = None; + right_tuple.values.clear(); + right_tuple.values.resize(right_schema_len, DataValue::Null); + Self::emit_tuple(&active_left.left_tuple, right_tuple, self.ty) } - _ => None, + _ => false, }; self.state = NestedLoopJoinState::PullLeft { right_bitmap }; - if let Some(tuple) = tuple { - arena.produce_tuple(tuple); + if should_emit { + arena.resume(); return Ok(()); } state = std::mem::replace(&mut self.state, NestedLoopJoinState::End); @@ -331,7 +283,6 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for NestedLoopJoin<'a> { mut right_emit_index, } => { while arena.next_tuple(right_input, plan_arena)? { - let mut right_tuple = arena.materialize_tuple(); let idx = right_emit_index; right_emit_index += 1; @@ -341,14 +292,15 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for NestedLoopJoin<'a> { }; if is_unmatched { - let mut values = vec![DataValue::Null; self.eq_cond.left_len]; - values.append(&mut right_tuple.values); + let values = &mut arena.result_tuple_mut().values; + values.resize(values.len() + self.eq_cond.left_len, DataValue::Null); + values.rotate_right(self.eq_cond.left_len); self.state = NestedLoopJoinState::EmitRightUnmatched { right_input, right_bitmap, right_emit_index, }; - arena.produce_tuple(Tuple::new(right_tuple.pk, values)); + arena.resume(); return Ok(()); } } @@ -382,47 +334,23 @@ impl<'a> NestedLoopJoin<'a> { /// Emit a tuple according to the join type. /// - /// `left_tuple`: left tuple to be included. - /// `right_tuple` right tuple to be included. + /// `left_tuple`: retained outer tuple (logical right for RightOuter). + /// `right_tuple`: current inner tuple (logical left for RightOuter), rewritten in place. /// `ty`: the type of join - /// `is_match`: whether [`NestedLoopJoin::left_input`] and [`NestedLoopJoin::right_input`] are matched - fn emit_tuple( - left_tuple: &Tuple, - right_tuple: &Tuple, - ty: JoinType, - is_matched: bool, - ) -> Option { - let left_len = left_tuple.values.len(); - let mut values = left_tuple - .values - .iter() - .cloned() - .chain(right_tuple.values.clone()) - .collect_vec(); - match ty { - JoinType::Inner | JoinType::Cross if !is_matched => values.clear(), - JoinType::LeftOuter | JoinType::Full if !is_matched => { - values - .iter_mut() - .skip(left_len) - .for_each(|v| *v = DataValue::Null); + fn emit_tuple(left_tuple: &Tuple, right_tuple: &mut Tuple, ty: JoinType) -> bool { + let right_len = right_tuple.values.len(); + right_tuple.values.extend(left_tuple.values.iter().cloned()); + if matches!(ty, JoinType::RightOuter) { + if right_tuple.pk.is_none() { + right_tuple.pk = left_tuple.pk.clone(); } - JoinType::RightOuter if !is_matched => { - (0..left_len).for_each(|i| { - values[i] = DataValue::Null; - }); + } else { + right_tuple.values.rotate_left(right_len); + if left_tuple.pk.is_some() { + right_tuple.pk = left_tuple.pk.clone(); } - _ => (), - }; - - if values.is_empty() { - return None; } - - Some(Tuple::new( - left_tuple.pk.as_ref().or(right_tuple.pk.as_ref()).cloned(), - values, - )) + !right_tuple.values.is_empty() } } @@ -434,6 +362,7 @@ mod test { use crate::execution::dql::test::build_integers; use crate::execution::try_collect; use crate::expression::BinaryOperator; + use crate::iter_ext::Itertools; use crate::optimizer::heuristic::batch::HepBatchStrategy; use crate::optimizer::heuristic::optimizer::HepOptimizerPipeline; use crate::optimizer::rule::normalization::NormalizationRuleImpl; @@ -641,7 +570,7 @@ mod test { let mut plan = cross(left.clone(), cross(left, right)); plan.populate_output_schema_recursive(&mut plan_arena); let context = crate::execution::empty_context(&table_cache, &view_cache, &meta_cache); - let mut arena = ExecArena::new(); + let mut arena = ExecArena::with_capacity(plan.exec_capacity_hint()); arena.init_context(context, &transaction); let root = build_read(&mut arena, &mut plan_arena, &plan, context, &transaction); let count = arena.nodes.items.len(); diff --git a/src/execution/dql/mark_apply.rs b/src/execution/dql/mark_apply.rs index 6f11b25a..f4b7c9cc 100644 --- a/src/execution/dql/mark_apply.rs +++ b/src/execution/dql/mark_apply.rs @@ -98,10 +98,12 @@ impl<'a> MarkApply<'a> { while arena.next_tuple(*root, plan_arena)? { let right = arena.result_tuple(); if Self::predicates_matched(self.op.predicates(), left, right, plan_arena)? { - let mut output = left.clone(); - output.pk = output.pk.or_else(|| right.pk.clone()); - output.values.extend(right.values.iter().cloned()); - arena.produce_tuple(output); + let output = arena.result_tuple_mut(); + let right_len = output.values.len(); + output.values.extend(left.values.iter().cloned()); + output.values.rotate_left(right_len); + output.pk = left.pk.clone().or(output.pk.take()); + arena.resume(); return Ok(()); } } @@ -197,7 +199,7 @@ impl<'a> MarkApply<'a> { Ok(DataValue::Boolean(false)) } MarkApplyKind::Quantified(MarkApplyQuantifier::Any) => { - if let Some(probe_value) = self.parameterized_probe_value(left_tuple, plan_arena)? { + if let Some(probe_value) = probe { if !probe_value.is_null() { let right_input = self.build_right_input(arena, plan_arena, Some(probe_value)); @@ -507,7 +509,7 @@ mod tests { let (table_cache, view_cache, meta_cache, _temp_dir, storage) = build_test_storage()?; let transaction = storage.transaction()?; let cache = crate::execution::empty_context(&table_cache, &view_cache, &meta_cache); - let mut arena = ExecArena::new(); + let mut arena = ExecArena::with_capacity(0); arena.init_context(cache, &transaction); let root = >::into_executor( (&op, &left, &right), @@ -597,7 +599,7 @@ mod tests { let (table_cache, view_cache, meta_cache, _temp_dir, storage) = build_test_storage()?; let transaction = storage.transaction()?; let context = crate::execution::empty_context(&table_cache, &view_cache, &meta_cache); - let mut arena = ExecArena::new(); + let mut arena = ExecArena::with_capacity(0); arena.init_context(context, &transaction); let root = >::into_executor( (&op, &left, &right), @@ -724,7 +726,7 @@ mod tests { let (table_cache, view_cache, meta_cache, _temp_dir, storage) = build_test_storage()?; let transaction = storage.transaction()?; - let mut arena = ExecArena::new(); + let mut arena = ExecArena::with_capacity(0); arena.init_context( crate::execution::empty_context(&table_cache, &view_cache, &meta_cache), &transaction, @@ -777,7 +779,7 @@ mod tests { let (table_cache, view_cache, meta_cache, _temp_dir, storage) = build_test_storage()?; let transaction = storage.transaction()?; - let mut arena = ExecArena::new(); + let mut arena = ExecArena::with_capacity(0); arena.init_context( crate::execution::empty_context(&table_cache, &view_cache, &meta_cache), &transaction, @@ -830,7 +832,7 @@ mod tests { let (table_cache, view_cache, meta_cache, _temp_dir, storage) = build_test_storage()?; let transaction = storage.transaction()?; - let mut arena = ExecArena::new(); + let mut arena = ExecArena::with_capacity(0); arena.init_context( crate::execution::empty_context(&table_cache, &view_cache, &meta_cache), &transaction, diff --git a/src/execution/dql/recursive_cte.rs b/src/execution/dql/recursive_cte.rs index a59441f1..9563baef 100644 --- a/src/execution/dql/recursive_cte.rs +++ b/src/execution/dql/recursive_cte.rs @@ -234,7 +234,7 @@ impl<'a, T: Transaction + 'a> ReadExecutor<'a, T> for RecursiveCte<'a, T> { cache: ExecutionContext<'_>, transaction: &T, ) -> ExecId { - let mut recursive_arena = ExecArena::new(); + let mut recursive_arena = ExecArena::with_capacity(recursive_plan.exec_capacity_hint()); recursive_arena.init_context(arena.context(), arena.transaction()); let anchor_input = build_read(arena, plan_arena, anchor_plan, cache, transaction); arena.push(ExecNode::RecursiveCte(Self::new( diff --git a/src/execution/dql/set_membership.rs b/src/execution/dql/set_membership.rs index 7f64f148..75c38726 100644 --- a/src/execution/dql/set_membership.rs +++ b/src/execution/dql/set_membership.rs @@ -61,10 +61,11 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for SetMembership { ) -> Result<(), DatabaseError> { if !self.built { while arena.next_tuple(self.right_input, plan_arena)? { - *self - .right_counts - .entry(arena.materialize_tuple()) - .or_insert(0) += 1; + if let Some(count) = self.right_counts.get_mut(arena.result_tuple()) { + *count += 1; + } else { + self.right_counts.insert(arena.materialize_tuple(), 1); + } } self.built = true; } diff --git a/src/execution/dql/sort.rs b/src/execution/dql/sort.rs index 569178f6..85257185 100644 --- a/src/execution/dql/sort.rs +++ b/src/execution/dql/sort.rs @@ -86,7 +86,7 @@ pub(crate) fn sort_tuples( ) -> Result<(), DatabaseError> { // Extract the results of calculating SortFields to avoid double calculation // of data during comparison. - let mut eval_values = vec![Vec::with_capacity(tuples.len()); sort_fields.len()]; + let mut eval_values = vec![Vec::new(); sort_fields.len()]; for (x, SortField { expr, .. }) in sort_fields.iter().enumerate() { for (_, tuple) in tuples.iter() { diff --git a/src/execution/dql/top_k.rs b/src/execution/dql/top_k.rs index 55cc9168..108a772b 100644 --- a/src/execution/dql/top_k.rs +++ b/src/execution/dql/top_k.rs @@ -60,49 +60,46 @@ impl PartialOrd for CmpItem<'_> { #[allow(clippy::mutable_key_type)] fn top_sort<'a>( - arena: &'a Bump, + full_key: &mut BumpBytes<'a>, sort_fields: &[SortField], heap: &mut BTreeSet>, - tuple: Tuple, + tuple: &mut Tuple, keep_count: usize, sequence: usize, plan_arena: &(dyn MetaArena + '_), ) -> Result<(), DatabaseError> { - let mut full_key = BumpBytes::new_in(arena); + full_key.clear(); for SortField { expr, nulls_first, asc, } in sort_fields { - let mut key = BumpBytes::new_in(arena); + let start = full_key.len(); plan_arena .expression(*expr) - .eval(plan_arena, Some(&tuple))? - .memcomparable_encode_with_null_order(&mut key, *nulls_first)?; - if !asc && key.len() > 1 { - for byte in key.iter_mut().skip(1) { + .eval(plan_arena, Some(&*tuple))? + .memcomparable_encode_with_null_order(full_key, *nulls_first)?; + if !asc { + for byte in &mut full_key[start + 1..] { *byte ^= 0xFF; } } - full_key.extend(key); } if heap.len() < keep_count { heap.insert(CmpItem { - key: full_key, + key: std::mem::replace(full_key, BumpBytes::new_in(full_key.bump())), sequence, - tuple, + tuple: std::mem::take(tuple), }); - } else if let Some(cmp_item) = heap.last() { + } else if let Some(mut cmp_item) = heap.pop_last() { if full_key.as_slice() < cmp_item.key.as_slice() { - heap.pop_last(); - heap.insert(CmpItem { - key: full_key, - sequence, - tuple, - }); + std::mem::swap(full_key, &mut cmp_item.key); + cmp_item.sequence = sequence; + cmp_item.tuple = std::mem::take(tuple); } + heap.insert(cmp_item); } Ok(()) } @@ -157,12 +154,13 @@ impl<'a, T: Transaction + 'a> ExecutorNode<'a, T> for TopK<'a> { let mut set = BTreeSet::new(); let mut sequence = 0; + let mut key_scratch = BumpBytes::new_in(&self.arena); while arena.next_tuple(self.input, plan_arena)? { top_sort( - &self.arena, + &mut key_scratch, self.sort_fields, &mut set, - arena.materialize_tuple(), + arena.result_tuple_mut(), keep_count, sequence, plan_arena, @@ -255,6 +253,7 @@ mod test { }] }; let arena = Bump::new(); + let mut key_scratch = crate::storage::table_codec::BumpBytes::new_in(&arena); let fn_asc_and_nulls_last_eq = |mut heap: BTreeSet>| { if let Some(reverse) = heap.pop_first() { @@ -308,28 +307,28 @@ mod test { let mut indices = BTreeSet::new(); top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(true, true), &mut indices, - Tuple::new(None, vec![DataValue::Null]), + &mut Tuple::new(None, vec![DataValue::Null]), 2, 0, &plan_arena, )?; top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(true, true), &mut indices, - Tuple::new(None, vec![DataValue::Int32(0)]), + &mut Tuple::new(None, vec![DataValue::Int32(0)]), 2, 1, &plan_arena, )?; top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(true, true), &mut indices, - Tuple::new(None, vec![DataValue::Int32(1)]), + &mut Tuple::new(None, vec![DataValue::Int32(1)]), 2, 2, &plan_arena, @@ -339,28 +338,28 @@ mod test { let mut indices = BTreeSet::new(); top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(true, false), &mut indices, - Tuple::new(None, vec![DataValue::Null]), + &mut Tuple::new(None, vec![DataValue::Null]), 2, 3, &plan_arena, )?; top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(true, false), &mut indices, - Tuple::new(None, vec![DataValue::Int32(0)]), + &mut Tuple::new(None, vec![DataValue::Int32(0)]), 2, 4, &plan_arena, )?; top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(true, false), &mut indices, - Tuple::new(None, vec![DataValue::Int32(1)]), + &mut Tuple::new(None, vec![DataValue::Int32(1)]), 2, 5, &plan_arena, @@ -370,28 +369,28 @@ mod test { let mut indices = BTreeSet::new(); top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(false, true), &mut indices, - Tuple::new(None, vec![DataValue::Null]), + &mut Tuple::new(None, vec![DataValue::Null]), 2, 6, &plan_arena, )?; top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(false, true), &mut indices, - Tuple::new(None, vec![DataValue::Int32(0)]), + &mut Tuple::new(None, vec![DataValue::Int32(0)]), 2, 7, &plan_arena, )?; top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(false, true), &mut indices, - Tuple::new(None, vec![DataValue::Int32(1)]), + &mut Tuple::new(None, vec![DataValue::Int32(1)]), 2, 8, &plan_arena, @@ -401,28 +400,28 @@ mod test { let mut indices = BTreeSet::new(); top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(false, false), &mut indices, - Tuple::new(None, vec![DataValue::Null]), + &mut Tuple::new(None, vec![DataValue::Null]), 2, 9, &plan_arena, )?; top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(false, false), &mut indices, - Tuple::new(None, vec![DataValue::Int32(0)]), + &mut Tuple::new(None, vec![DataValue::Int32(0)]), 2, 10, &plan_arena, )?; top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(false, false), &mut indices, - Tuple::new(None, vec![DataValue::Int32(1)]), + &mut Tuple::new(None, vec![DataValue::Int32(1)]), 2, 11, &plan_arena, @@ -470,6 +469,7 @@ mod test { ] }; let arena = Bump::new(); + let mut key_scratch = crate::storage::table_codec::BumpBytes::new_in(&arena); let fn_asc_1_and_nulls_first_1_and_asc_2_and_nulls_first_2_eq = |mut heap: BTreeSet>| { @@ -609,55 +609,55 @@ mod test { let mut indices = BTreeSet::new(); top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(true, true, true, true), &mut indices, - Tuple::new(None, vec![DataValue::Null, DataValue::Null]), + &mut Tuple::new(None, vec![DataValue::Null, DataValue::Null]), 4, 12, &plan_arena, )?; top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(true, true, true, true), &mut indices, - Tuple::new(None, vec![DataValue::Int32(0), DataValue::Null]), + &mut Tuple::new(None, vec![DataValue::Int32(0), DataValue::Null]), 4, 13, &plan_arena, )?; top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(true, true, true, true), &mut indices, - Tuple::new(None, vec![DataValue::Int32(1), DataValue::Null]), + &mut Tuple::new(None, vec![DataValue::Int32(1), DataValue::Null]), 4, 14, &plan_arena, )?; top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(true, true, true, true), &mut indices, - Tuple::new(None, vec![DataValue::Null, DataValue::Int32(0)]), + &mut Tuple::new(None, vec![DataValue::Null, DataValue::Int32(0)]), 4, 15, &plan_arena, )?; top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(true, true, true, true), &mut indices, - Tuple::new(None, vec![DataValue::Int32(0), DataValue::Int32(0)]), + &mut Tuple::new(None, vec![DataValue::Int32(0), DataValue::Int32(0)]), 4, 16, &plan_arena, )?; top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(true, true, true, true), &mut indices, - Tuple::new(None, vec![DataValue::Int32(1), DataValue::Int32(0)]), + &mut Tuple::new(None, vec![DataValue::Int32(1), DataValue::Int32(0)]), 4, 17, &plan_arena, @@ -667,55 +667,55 @@ mod test { let mut indices = BTreeSet::new(); top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(true, false, true, true), &mut indices, - Tuple::new(None, vec![DataValue::Null, DataValue::Null]), + &mut Tuple::new(None, vec![DataValue::Null, DataValue::Null]), 4, 18, &plan_arena, )?; top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(true, false, true, true), &mut indices, - Tuple::new(None, vec![DataValue::Int32(0), DataValue::Null]), + &mut Tuple::new(None, vec![DataValue::Int32(0), DataValue::Null]), 4, 19, &plan_arena, )?; top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(true, false, true, true), &mut indices, - Tuple::new(None, vec![DataValue::Int32(1), DataValue::Null]), + &mut Tuple::new(None, vec![DataValue::Int32(1), DataValue::Null]), 4, 20, &plan_arena, )?; top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(true, false, true, true), &mut indices, - Tuple::new(None, vec![DataValue::Null, DataValue::Int32(0)]), + &mut Tuple::new(None, vec![DataValue::Null, DataValue::Int32(0)]), 4, 21, &plan_arena, )?; top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(true, false, true, true), &mut indices, - Tuple::new(None, vec![DataValue::Int32(0), DataValue::Int32(0)]), + &mut Tuple::new(None, vec![DataValue::Int32(0), DataValue::Int32(0)]), 4, 22, &plan_arena, )?; top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(true, false, true, true), &mut indices, - Tuple::new(None, vec![DataValue::Int32(1), DataValue::Int32(0)]), + &mut Tuple::new(None, vec![DataValue::Int32(1), DataValue::Int32(0)]), 4, 23, &plan_arena, @@ -725,55 +725,55 @@ mod test { let mut indices = BTreeSet::new(); top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(false, true, true, true), &mut indices, - Tuple::new(None, vec![DataValue::Null, DataValue::Null]), + &mut Tuple::new(None, vec![DataValue::Null, DataValue::Null]), 4, 24, &plan_arena, )?; top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(false, true, true, true), &mut indices, - Tuple::new(None, vec![DataValue::Int32(0), DataValue::Null]), + &mut Tuple::new(None, vec![DataValue::Int32(0), DataValue::Null]), 4, 25, &plan_arena, )?; top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(false, true, true, true), &mut indices, - Tuple::new(None, vec![DataValue::Int32(1), DataValue::Null]), + &mut Tuple::new(None, vec![DataValue::Int32(1), DataValue::Null]), 4, 26, &plan_arena, )?; top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(false, true, true, true), &mut indices, - Tuple::new(None, vec![DataValue::Null, DataValue::Int32(0)]), + &mut Tuple::new(None, vec![DataValue::Null, DataValue::Int32(0)]), 4, 27, &plan_arena, )?; top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(false, true, true, true), &mut indices, - Tuple::new(None, vec![DataValue::Int32(0), DataValue::Int32(0)]), + &mut Tuple::new(None, vec![DataValue::Int32(0), DataValue::Int32(0)]), 4, 28, &plan_arena, )?; top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(false, true, true, true), &mut indices, - Tuple::new(None, vec![DataValue::Int32(1), DataValue::Int32(0)]), + &mut Tuple::new(None, vec![DataValue::Int32(1), DataValue::Int32(0)]), 4, 29, &plan_arena, @@ -783,55 +783,55 @@ mod test { let mut indices = BTreeSet::new(); top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(false, false, true, true), &mut indices, - Tuple::new(None, vec![DataValue::Null, DataValue::Null]), + &mut Tuple::new(None, vec![DataValue::Null, DataValue::Null]), 4, 30, &plan_arena, )?; top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(false, false, true, true), &mut indices, - Tuple::new(None, vec![DataValue::Int32(0), DataValue::Null]), + &mut Tuple::new(None, vec![DataValue::Int32(0), DataValue::Null]), 4, 31, &plan_arena, )?; top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(false, false, true, true), &mut indices, - Tuple::new(None, vec![DataValue::Int32(1), DataValue::Null]), + &mut Tuple::new(None, vec![DataValue::Int32(1), DataValue::Null]), 4, 32, &plan_arena, )?; top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(false, false, true, true), &mut indices, - Tuple::new(None, vec![DataValue::Null, DataValue::Int32(0)]), + &mut Tuple::new(None, vec![DataValue::Null, DataValue::Int32(0)]), 4, 33, &plan_arena, )?; top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(false, false, true, true), &mut indices, - Tuple::new(None, vec![DataValue::Int32(0), DataValue::Int32(0)]), + &mut Tuple::new(None, vec![DataValue::Int32(0), DataValue::Int32(0)]), 4, 34, &plan_arena, )?; top_sort( - &arena, + &mut key_scratch, &fn_sort_fields(false, false, true, true), &mut indices, - Tuple::new(None, vec![DataValue::Int32(1), DataValue::Int32(0)]), + &mut Tuple::new(None, vec![DataValue::Int32(1), DataValue::Int32(0)]), 4, 35, &plan_arena, diff --git a/src/execution/dql/window.rs b/src/execution/dql/window.rs index a2011b26..cf5453da 100644 --- a/src/execution/dql/window.rs +++ b/src/execution/dql/window.rs @@ -131,17 +131,16 @@ impl<'a> Window<'a> { for (index, field) in self.sort_fields.iter().enumerate() { let value = plan_arena .expression(field.expr) - .eval(plan_arena, Some(tuple))? - .into_owned(); - if self.state.started && self.state.sort_values[index] != value { + .eval(plan_arena, Some(tuple))?; + if self.state.started && self.state.sort_values[index] != *value { if index < self.partition_by_len { boundary = Some(Boundary::Partition); } else if boundary.is_none() { boundary = Some(Boundary::Peer); } - self.state.sort_values[index] = value; + self.state.sort_values[index] = value.into_owned(); } else if !self.state.started { - self.state.sort_values.push(value); + self.state.sort_values.push(value.into_owned()); } } self.state.started = true; diff --git a/src/execution/mod.rs b/src/execution/mod.rs index 3176cfb8..e7d87195 100644 --- a/src/execution/mod.rs +++ b/src/execution/mod.rs @@ -469,6 +469,16 @@ impl<'b, 'a, T: Transaction + 'a> ExecArenaLocalState<'b, 'a, T> { unsafe { (&mut *self.transaction, &mut *self.table_codec) } } + pub(crate) fn tuple_transaction_codec_mut(&mut self) -> (&Tuple, &mut T, &mut TableCodec) { + unsafe { + ( + &self.result.tuple, + &mut *self.transaction, + &mut *self.table_codec, + ) + } + } + pub(crate) fn index_values_transaction_codec_mut( &mut self, ) -> (&[DataValue], &mut T, &mut TableCodec) { @@ -503,10 +513,10 @@ impl<'a, T: Transaction + 'a> ExecArena<'a, T> { self.table_codec.set_stamp(stamp); } - pub(crate) fn new() -> Self { + pub(crate) fn with_capacity(node_capacity: usize) -> Self { Self { nodes: ExecNodes { - items: Vec::new(), + items: Vec::with_capacity(node_capacity), pos: 0, executing: 0, }, @@ -955,7 +965,7 @@ mod test_utils { T: Transaction + 'a, E: ReadExecutor<'a, T>, { - let mut arena = ExecArena::new(); + let mut arena = ExecArena::with_capacity(0); arena.init_context(cache, transaction); let root = >::into_executor( input, @@ -981,7 +991,7 @@ mod test_utils { T: Transaction + 'a, E: WriteExecutor<'a, T>, { - let mut arena = ExecArena::new(); + let mut arena = ExecArena::with_capacity(0); arena.init_context(cache, transaction); let root = >::into_executor( input, @@ -1023,7 +1033,7 @@ mod test { fn active_nodes_cannot_be_overwritten_or_relocated() { let table_arena = crate::planner::TableArenaCell::default(); let mut plan_arena = crate::planner::PlanArena::new(&table_arena); - let mut arena = ExecArena::<'_, MemoryTransaction>::new(); + let mut arena = ExecArena::<'_, MemoryTransaction>::with_capacity(0); arena.push(ExecNode::Dummy(Dummy::default())); let slot = &arena.nodes.items[0] as *const std::cell::RefCell>; diff --git a/src/planner/mod.rs b/src/planner/mod.rs index 484e5585..155f7bac 100644 --- a/src/planner/mod.rs +++ b/src/planner/mod.rs @@ -156,6 +156,7 @@ pub struct LogicalPlan { pub(crate) childrens: Box, pub(crate) physical_option: Option, output_schema: Option, + exec_capacity_hint: usize, } impl LogicalPlan { @@ -165,6 +166,7 @@ impl LogicalPlan { childrens: Box::new(childrens), physical_option: None, output_schema: None, + exec_capacity_hint: 0, } } @@ -359,20 +361,30 @@ impl LogicalPlan { .expect("output schema must be computed before it is read") } + pub(crate) fn exec_capacity_hint(&self) -> usize { + self.exec_capacity_hint + } + pub(crate) fn populate_output_schema_recursive(&mut self, arena: &mut (dyn MetaArena + '_)) { - match self.childrens.as_mut() { - Childrens::Only(child) => child.populate_output_schema_recursive(arena), + let child_nodes = match self.childrens.as_mut() { + Childrens::Only(child) => { + child.populate_output_schema_recursive(arena); + child.exec_capacity_hint + } Childrens::Twins { left, right } => { left.populate_output_schema_recursive(arena); right.populate_output_schema_recursive(arena); + left.exec_capacity_hint + right.exec_capacity_hint } - Childrens::None => (), - } + Childrens::None => 0, + }; + self.exec_capacity_hint = 1 + child_nodes; self.output_schema(arena); } pub fn reset_output_schema_cache(&mut self) { self.output_schema = None; + self.exec_capacity_hint = 0; } pub fn reset_output_schema_cache_recursive(&mut self) { @@ -420,6 +432,7 @@ impl Clone for LogicalPlan { childrens: self.childrens.clone(), physical_option: self.physical_option.clone(), output_schema: self.output_schema.clone(), + exec_capacity_hint: self.exec_capacity_hint, } } } @@ -504,6 +517,7 @@ impl crate::serdes::ReferenceSerialization for LogicalPlan { childrens, physical_option, output_schema: None, + exec_capacity_hint: 0, }) } } diff --git a/src/types/tuple_builder.rs b/src/types/tuple_builder.rs index f836ec2a..b5d9cd9e 100644 --- a/src/types/tuple_builder.rs +++ b/src/types/tuple_builder.rs @@ -53,12 +53,14 @@ impl<'a> TupleBuilder<'a> { pub fn build_with_row<'b>( &self, + tuple: &mut Tuple, row: impl IntoIterator, - ) -> Result { - let mut values = Vec::with_capacity(self.column_types.len()); + ) -> Result<(), DatabaseError> { + tuple.pk = None; + tuple.values.clear(); for (i, value) in row.into_iter().enumerate() { - values.push( + tuple.values.push( DataValue::Utf8 { value: value.to_string(), ty: Utf8Type::Variable(None), @@ -67,15 +69,15 @@ impl<'a> TupleBuilder<'a> { .cast(&self.column_types[i])?, ); } - if values.len() != self.column_types.len() { + if tuple.values.len() != self.column_types.len() { return Err(DatabaseError::MisMatch("types", "values")); } - let pk = self + tuple.pk = self .pk_indices - .map(|indices| Tuple::primary_projection(indices, &values)); + .map(|indices| Tuple::primary_projection(indices, &tuple.values)); - Ok(Tuple::new(pk, values)) + Ok(()) } } @@ -103,7 +105,7 @@ mod tests { ], Some(&pk_indices), ); - let tuple = builder.build_with_row(["7", "kite"]).unwrap(); + builder.build_with_row(&mut tuple, ["7", "kite"]).unwrap(); assert_eq!( tuple.pk, Some(DataValue::Tuple(vec![ @@ -117,7 +119,7 @@ mod tests { ); assert!(matches!( - builder.build_with_row(["7"]), + builder.build_with_row(&mut tuple, ["7"]), Err(DatabaseError::MisMatch("types", "values")) )); }