Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions src/catalog/table.rs
Original file line number Diff line number Diff line change
Expand Up @@ -125,6 +125,9 @@ impl TableCatalog {
) -> Result<DmlTableSnapshot<'_>, DatabaseError> {
let index_metas = self
.indexes()
.filter(|index_meta| {
!matches!(arena.index(**index_meta).ty, IndexType::PrimaryKey { .. })
})
.map(|index_meta| {
Ok((
*index_meta,
Expand Down
4 changes: 2 additions & 2 deletions src/db.rs
Original file line number Diff line number Diff line change
Expand Up @@ -589,7 +589,7 @@ impl<S: Storage> State<S> {
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,
Expand Down Expand Up @@ -665,7 +665,7 @@ impl<S: Storage> State<S> {
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,
Expand Down
4 changes: 2 additions & 2 deletions src/db/prepared.rs
Original file line number Diff line number Diff line change
Expand Up @@ -323,7 +323,7 @@ mod tests {
}
let mut execution_arena = crate::execution::ExecArena::<
<crate::storage::memory::MemoryStorage as Storage>::TransactionType<'_>,
>::new();
>::with_capacity(0);
let executor = crate::execution::dql::index_scan::IndexScan::new(
scan_op,
info,
Expand Down Expand Up @@ -551,7 +551,7 @@ mod tests {
};
let mut execution_arena = crate::execution::ExecArena::<
<crate::storage::memory::MemoryStorage as Storage>::TransactionType<'_>,
>::new();
>::with_capacity(0);
let ranges = crate::execution::dql::index_scan::IndexScan::new(
scan_op,
info,
Expand Down
3 changes: 1 addition & 2 deletions src/execution/ddl/add_column.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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(
Expand Down
5 changes: 3 additions & 2 deletions src/execution/ddl/change_column.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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()));
}
Expand Down
11 changes: 6 additions & 5 deletions src/execution/ddl/create_index.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)?;
}
}

Expand Down
11 changes: 6 additions & 5 deletions src/execution/dml/analyze.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)?;
}
Expand Down
9 changes: 5 additions & 4 deletions src/execution/dml/copy_from_file.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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)?;
Expand Down
9 changes: 1 addition & 8 deletions src/execution/dml/copy_to_file.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::<Vec<_>>(),
)?;
writer.write_record(arena.result_tuple().values.iter().map(|v| v.to_string()))?;
}
writer.flush().map_err(DatabaseError::from)?;

Expand Down
13 changes: 9 additions & 4 deletions src/execution/dml/delete.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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> {
Expand Down Expand Up @@ -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((
Expand All @@ -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);
Expand Down
33 changes: 12 additions & 21 deletions src/execution/dml/insert.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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());
Expand All @@ -121,28 +122,24 @@ 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);
}

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()));
Expand All @@ -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)?;
Expand All @@ -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(())
}
}
49 changes: 21 additions & 28 deletions src/execution/dml/update.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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;

Expand All @@ -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,
Expand All @@ -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 {
Expand All @@ -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)?;
}

Expand All @@ -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(())
}
}
3 changes: 2 additions & 1 deletion src/execution/dql/aggregate/count.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}

Expand Down
2 changes: 1 addition & 1 deletion src/execution/dql/aggregate/hash_agg.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
}

Expand Down
Loading
Loading