| // Licensed to the Apache Software Foundation (ASF) under one |
| // or more contributor license agreements. See the NOTICE file |
| // distributed with this work for additional information |
| // regarding copyright ownership. The ASF licenses this file |
| // to you under the Apache License, Version 2.0 (the |
| // "License"); you may not use this file except in compliance |
| // with the License. You may obtain a copy of the License at |
| // |
| // http://www.apache.org/licenses/LICENSE-2.0 |
| // |
| // Unless required by applicable law or agreed to in writing, |
| // software distributed under the License is distributed on an |
| // "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY |
| // KIND, either express or implied. See the License for the |
| // specific language governing permissions and limitations |
| // under the License. |
| |
| //! TableWrite for writing Arrow data to Paimon tables. |
| //! |
| //! Reference: [pypaimon TableWrite](https://github.com/apache/paimon/blob/master/paimon-python/pypaimon/write/table_write.py) |
| //! and [pypaimon FileStoreWrite](https://github.com/apache/paimon/blob/master/paimon-python/pypaimon/write/file_store_write.py) |
| |
| use crate::arrow::build_target_arrow_schema; |
| use crate::spec::PartitionComputer; |
| use crate::spec::{ |
| first_row_supports_changelog_producer, BinaryRow, ChangelogProducer, CoreOptions, DataField, |
| DataType, MergeEngine, EMPTY_SERIALIZED_ROW, POSTPONE_BUCKET, |
| }; |
| use crate::table::blob_file_writer::AppendBlobFileWriter; |
| use crate::table::bucket_assigner::{BucketAssignerEnum, PartitionBucketKey}; |
| use crate::table::bucket_assigner_constant::ConstantBucketAssigner; |
| use crate::table::bucket_assigner_cross::CrossPartitionAssigner; |
| use crate::table::bucket_assigner_dynamic::DynamicBucketAssigner; |
| use crate::table::bucket_assigner_fixed::FixedBucketAssigner; |
| use crate::table::bucket_function::validate_bucket_function; |
| use crate::table::commit_message::CommitMessage; |
| use crate::table::data_file_writer::DataFileWriter; |
| use crate::table::kv_file_writer::{KeyValueFileWriter, KeyValueWriteConfig}; |
| use crate::table::partition_filter::PartitionFilter; |
| use crate::table::postpone_file_writer::{PostponeFileWriter, PostponeWriteConfig}; |
| use crate::table::prepared_files::PreparedFiles; |
| use crate::table::{SnapshotManager, Table, TableScan}; |
| use crate::Result; |
| use arrow_array::RecordBatch; |
| use std::collections::{HashMap, HashSet}; |
| use std::sync::Arc; |
| |
| /// Enum to hold either an append-only writer, a key-value writer, or a postpone writer. |
| enum FileWriter { |
| Append(DataFileWriter), |
| AppendBlob(AppendBlobFileWriter), |
| KeyValue(KeyValueFileWriter), |
| Postpone(PostponeFileWriter), |
| } |
| |
| impl FileWriter { |
| async fn write(&mut self, batch: &RecordBatch) -> Result<()> { |
| match self { |
| FileWriter::Append(w) => w.write(batch).await, |
| FileWriter::AppendBlob(w) => w.write(batch).await, |
| FileWriter::KeyValue(w) => w.write(batch).await, |
| FileWriter::Postpone(w) => w.write(batch).await, |
| } |
| } |
| |
| async fn prepare_commit(mut self) -> Result<PreparedFiles> { |
| match self { |
| FileWriter::Append(ref mut w) => w.prepare_commit().await.map(PreparedFiles::data), |
| FileWriter::AppendBlob(ref mut w) => w.prepare_commit().await.map(PreparedFiles::data), |
| FileWriter::KeyValue(ref mut w) => w.prepare_commit().await, |
| FileWriter::Postpone(ref mut w) => w.prepare_commit().await.map(PreparedFiles::data), |
| } |
| } |
| } |
| |
| /// TableWrite writes Arrow RecordBatches to Paimon data files. |
| /// |
| /// Each (partition, bucket) pair gets its own writer held in a HashMap. |
| /// Batches are routed to the correct writer based on partition/bucket. |
| /// |
| /// Call `prepare_commit()` to close all writers and collect |
| /// `CommitMessage`s for use with `TableCommit`. |
| /// |
| /// Reference: [pypaimon BatchTableWrite](https://github.com/apache/paimon/blob/master/paimon-python/pypaimon/write/table_write.py) |
| pub struct TableWrite { |
| table: Table, |
| partition_writers: HashMap<PartitionBucketKey, FileWriter>, |
| partition_computer: PartitionComputer, |
| partition_keys: Vec<String>, |
| schema_id: i64, |
| target_file_size: i64, |
| blob_target_file_size: i64, |
| file_compression: String, |
| file_compression_zstd_level: i32, |
| write_buffer_size: i64, |
| file_format: String, |
| primary_key_indices: Vec<usize>, |
| primary_key_types: Vec<DataType>, |
| sequence_field_indices: Vec<usize>, |
| merge_engine: MergeEngine, |
| changelog_producer: ChangelogProducer, |
| changelog_file_prefix: String, |
| changelog_file_format: String, |
| changelog_file_compression: String, |
| partition_seq_cache: HashMap<Vec<u8>, HashMap<i32, i64>>, |
| commit_user: String, |
| /// Bucket assignment strategy (fixed, dynamic, or cross-partition). |
| bucket_assigner: BucketAssignerEnum, |
| /// Whether this is an overwrite operation (skip seq/index restore). |
| is_overwrite: bool, |
| /// Blob descriptor fields (stored inline in parquet, not as separate .blob files). |
| blob_descriptor_fields: HashSet<String>, |
| /// Whether the table has non-descriptor blob fields requiring AppendBlobFileWriter. |
| has_blob_fields: bool, |
| } |
| |
| impl TableWrite { |
| pub(crate) fn new(table: &Table, commit_user: String) -> crate::Result<Self> { |
| let is_overwrite = false; |
| let schema = table.schema(); |
| let core_options = CoreOptions::new(schema.options()); |
| let blob_descriptor_fields = core_options.blob_descriptor_fields(); |
| |
| for name in &blob_descriptor_fields { |
| match schema.fields().iter().find(|f| f.name() == name) { |
| None => { |
| return Err(crate::Error::DataInvalid { |
| message: format!("blob-descriptor-field '{name}' does not exist in schema"), |
| source: None, |
| }); |
| } |
| Some(f) if !f.data_type().is_blob_type() => { |
| return Err(crate::Error::DataInvalid { |
| message: format!( |
| "blob-descriptor-field '{name}' is not a top-level BLOB field" |
| ), |
| source: None, |
| }); |
| } |
| _ => {} |
| } |
| } |
| |
| let total_buckets = core_options.bucket(); |
| let has_primary_keys = !schema.primary_keys().is_empty(); |
| let is_dynamic_bucket = has_primary_keys && total_buckets == -1; |
| let changelog_producer = core_options.try_changelog_producer()?; |
| |
| let is_dynamic_cross_partition = |
| is_dynamic_bucket && !schema.partition_keys().is_empty() && { |
| let pk_set: HashSet<&str> = |
| schema.primary_keys().iter().map(String::as_str).collect(); |
| schema |
| .partition_keys() |
| .iter() |
| .any(|p| !pk_set.contains(p.as_str())) |
| }; |
| |
| if has_primary_keys |
| && !is_dynamic_bucket |
| && total_buckets < 1 |
| && total_buckets != POSTPONE_BUCKET |
| { |
| return Err(crate::Error::Unsupported { |
| message: format!( |
| "KeyValueFileWriter does not support bucket={total_buckets}, only fixed bucket (>= 1), -1 (dynamic), or -2 (postpone) is supported" |
| ), |
| }); |
| } |
| |
| if !has_primary_keys && total_buckets != -1 && core_options.bucket_key().is_none() { |
| return Err(crate::Error::Unsupported { |
| message: "Append tables with fixed bucket must configure 'bucket-key'".to_string(), |
| }); |
| } |
| let target_file_size = core_options.target_file_size(); |
| let blob_target_file_size = core_options.blob_target_file_size(); |
| let file_compression = core_options.file_compression().to_string(); |
| let file_compression_zstd_level = core_options.file_compression_zstd_level(); |
| let file_format = core_options.file_format().to_string(); |
| let changelog_file_prefix = core_options.changelog_file_prefix().to_string(); |
| let changelog_file_format = core_options.changelog_file_format().to_string(); |
| let changelog_file_compression = core_options.changelog_file_compression().to_string(); |
| let write_buffer_size = core_options.write_parquet_buffer_size(); |
| let partition_keys: Vec<String> = schema.partition_keys().to_vec(); |
| let fields = schema.fields(); |
| |
| let partition_field_indices: Vec<usize> = partition_keys |
| .iter() |
| .filter_map(|pk| fields.iter().position(|f| f.name() == pk)) |
| .collect(); |
| |
| // Bucket keys: resolved by TableSchema |
| let bucket_keys = schema.bucket_keys(); |
| |
| let bucket_key_indices: Vec<usize> = bucket_keys |
| .iter() |
| .filter_map(|bk| fields.iter().position(|f| f.name() == bk)) |
| .collect(); |
| |
| let partition_computer = PartitionComputer::new( |
| &partition_keys, |
| fields, |
| core_options.partition_default_name(), |
| core_options.legacy_partition_name(), |
| ) |
| .unwrap(); |
| |
| let primary_key_indices: Vec<usize> = schema |
| .trimmed_primary_keys() |
| .iter() |
| .filter_map(|pk| fields.iter().position(|f| f.name() == pk)) |
| .collect(); |
| |
| let primary_key_types: Vec<DataType> = primary_key_indices |
| .iter() |
| .map(|&idx| fields[idx].data_type().clone()) |
| .collect(); |
| |
| let sequence_field_indices: Vec<usize> = core_options |
| .sequence_fields() |
| .iter() |
| .filter_map(|sf| fields.iter().position(|f| f.name() == *sf)) |
| .collect(); |
| |
| let merge_engine = core_options.merge_engine()?; |
| |
| if merge_engine == MergeEngine::FirstRow |
| && !first_row_supports_changelog_producer(changelog_producer) |
| { |
| return Err(crate::Error::Unsupported { |
| message: format!( |
| "Table '{}' has incompatible table options: merge-engine=first-row only supports changelog-producer=none or lookup, but found changelog-producer={}", |
| table.identifier().full_name(), |
| changelog_producer.as_str() |
| ), |
| }); |
| } |
| |
| if is_dynamic_cross_partition && merge_engine == MergeEngine::PartialUpdate { |
| return Err(crate::Error::Unsupported { |
| message: |
| "merge-engine=partial-update with cross-partition update is not supported yet" |
| .to_string(), |
| }); |
| } |
| |
| if is_dynamic_cross_partition && merge_engine == MergeEngine::Aggregation { |
| return Err(crate::Error::Unsupported { |
| message: |
| "merge-engine=aggregation with cross-partition update is not supported yet" |
| .to_string(), |
| }); |
| } |
| |
| if has_primary_keys && core_options.rowkind_field().is_some() { |
| return Err(crate::Error::Unsupported { |
| message: "KeyValueFileWriter does not support rowkind.field".to_string(), |
| }); |
| } |
| |
| let target_bucket_row_number = core_options.dynamic_bucket_target_row_num(); |
| let bucket_function_type = core_options.bucket_function_type()?; |
| |
| let bucket_assigner = if is_dynamic_cross_partition { |
| BucketAssignerEnum::CrossPartition(Box::new(CrossPartitionAssigner::new( |
| table.clone(), |
| partition_field_indices, |
| primary_key_indices.clone(), |
| target_bucket_row_number, |
| merge_engine, |
| ))) |
| } else if is_dynamic_bucket { |
| BucketAssignerEnum::Dynamic(DynamicBucketAssigner::new( |
| partition_field_indices, |
| primary_key_indices.clone(), |
| schema.fields().to_vec(), |
| target_bucket_row_number, |
| table.file_io().clone(), |
| table.location().to_string(), |
| is_overwrite, |
| )) |
| } else if total_buckets == POSTPONE_BUCKET { |
| BucketAssignerEnum::Constant(ConstantBucketAssigner::new( |
| partition_field_indices, |
| POSTPONE_BUCKET, |
| )) |
| } else if total_buckets <= 1 || bucket_key_indices.is_empty() { |
| BucketAssignerEnum::Constant(ConstantBucketAssigner::new(partition_field_indices, 0)) |
| } else { |
| let bucket_key_fields: Vec<DataField> = bucket_key_indices |
| .iter() |
| .map(|&idx| fields[idx].clone()) |
| .collect(); |
| validate_bucket_function(bucket_function_type, &bucket_key_fields)?; |
| BucketAssignerEnum::Fixed(FixedBucketAssigner::new( |
| partition_field_indices, |
| bucket_key_indices, |
| bucket_function_type, |
| total_buckets, |
| )) |
| }; |
| |
| let has_blob_fields = schema |
| .fields() |
| .iter() |
| .any(|f| f.data_type().is_blob_type() && !blob_descriptor_fields.contains(f.name())); |
| |
| Ok(Self { |
| table: table.clone(), |
| partition_writers: HashMap::new(), |
| partition_computer, |
| partition_keys, |
| schema_id: schema.id(), |
| target_file_size, |
| blob_target_file_size, |
| file_compression, |
| file_compression_zstd_level, |
| write_buffer_size, |
| file_format, |
| primary_key_indices, |
| primary_key_types, |
| sequence_field_indices, |
| merge_engine, |
| changelog_producer, |
| changelog_file_prefix, |
| changelog_file_format, |
| changelog_file_compression, |
| partition_seq_cache: HashMap::new(), |
| commit_user, |
| bucket_assigner, |
| is_overwrite, |
| blob_descriptor_fields, |
| has_blob_fields, |
| }) |
| } |
| |
| /// Scan the latest snapshot for a specific partition and return a map of |
| /// bucket → (max_sequence_number + 1) for each bucket in that partition. |
| async fn scan_partition_sequence_numbers( |
| table: &Table, |
| partition_bytes: &[u8], |
| ) -> crate::Result<HashMap<i32, i64>> { |
| let snapshot_manager = |
| SnapshotManager::new(table.file_io().clone(), table.location().to_string()); |
| let latest_snapshot = snapshot_manager.get_latest_snapshot().await?; |
| let mut bucket_seq: HashMap<i32, i64> = HashMap::new(); |
| if let Some(snapshot) = latest_snapshot { |
| let partition_filter = Self::build_partition_filter(table, partition_bytes)?; |
| let scan = TableScan::new(table, partition_filter, vec![], None, None, None) |
| .with_scan_all_files(); |
| let entries = scan.plan_manifest_entries(&snapshot).await?; |
| for entry in &entries { |
| let bucket = entry.bucket(); |
| let max_seq = entry.file().max_sequence_number; |
| let current = bucket_seq.entry(bucket).or_insert(0); |
| if max_seq + 1 > *current { |
| *current = max_seq + 1; |
| } |
| } |
| } |
| Ok(bucket_seq) |
| } |
| |
| /// Build a partition filter from serialized partition bytes. |
| /// |
| /// Uses `PartitionSet` for O(1) byte-level matching when partition fields exist. |
| fn build_partition_filter( |
| table: &Table, |
| partition_bytes: &[u8], |
| ) -> crate::Result<Option<PartitionFilter>> { |
| let partition_fields = table.schema().partition_fields(); |
| if partition_fields.is_empty() { |
| return Ok(None); |
| } |
| let partitions = HashSet::from([partition_bytes.to_vec()]); |
| Ok(Some(PartitionFilter::from_partition_set( |
| partitions, |
| &partition_fields, |
| )?)) |
| } |
| |
| /// Mark this write as an overwrite operation. |
| /// |
| /// Overwrite skips restoring sequence numbers and bucket index loading |
| /// since old data will be fully replaced at commit time. |
| pub fn with_overwrite(mut self) -> Self { |
| self.is_overwrite = true; |
| self.bucket_assigner.set_overwrite(true); |
| self |
| } |
| |
| /// Write an Arrow RecordBatch. Rows are routed to the correct partition and bucket. |
| pub async fn write_arrow_batch(&mut self, batch: &RecordBatch) -> Result<()> { |
| if batch.num_rows() == 0 { |
| return Ok(()); |
| } |
| |
| let grouped = self.divide_by_partition_bucket(batch).await?; |
| for ((partition_bytes, bucket), sub_batch) in grouped { |
| self.write_bucket(partition_bytes, bucket, sub_batch) |
| .await?; |
| } |
| Ok(()) |
| } |
| |
| /// Group rows by (partition_bytes, bucket) and return sub-batches. |
| /// |
| /// In cross-partition mode, also generates DELETE sub-batches for keys that |
| /// migrated from one partition to another. |
| async fn divide_by_partition_bucket( |
| &mut self, |
| batch: &RecordBatch, |
| ) -> Result<Vec<(PartitionBucketKey, RecordBatch)>> { |
| // Fast path: constant bucket with no partitions — skip per-row routing |
| if let BucketAssignerEnum::Constant(ref a) = self.bucket_assigner { |
| if self.partition_keys.is_empty() { |
| return Ok(vec![( |
| (EMPTY_SERIALIZED_ROW.clone(), a.bucket()), |
| batch.clone(), |
| )]); |
| } |
| } |
| |
| let fields = self.table.schema().fields().to_vec(); |
| let output = self.bucket_assigner.assign_batch(batch, &fields).await?; |
| |
| let mut groups: HashMap<PartitionBucketKey, Vec<usize>> = HashMap::new(); |
| let skip_set: HashSet<usize> = output.skips.into_iter().collect(); |
| for row_idx in 0..batch.num_rows() { |
| if skip_set.contains(&row_idx) { |
| continue; |
| } |
| groups |
| .entry(( |
| output.partition_bytes[row_idx].clone(), |
| output.buckets[row_idx], |
| )) |
| .or_default() |
| .push(row_idx); |
| } |
| |
| let mut result = Vec::with_capacity(groups.len()); |
| // Cross-partition writers must always include _VALUE_KIND to keep the |
| // Arrow schema stable across batches (some batches may have deletes, |
| // others may not — KeyValueFileWriter's concat_batches requires a |
| // consistent schema). |
| let needs_value_kind = |
| matches!(self.bucket_assigner, BucketAssignerEnum::CrossPartition(_)) |
| || !output.deletes.is_empty(); |
| for (key, row_indices) in groups { |
| let sub_batch = Self::take_rows(batch, &row_indices)?; |
| let sub_batch = if needs_value_kind { |
| Self::add_value_kind_column(&sub_batch, 0)? |
| } else { |
| sub_batch |
| }; |
| result.push((key, sub_batch)); |
| } |
| |
| if !output.deletes.is_empty() { |
| let mut delete_groups: HashMap<PartitionBucketKey, Vec<usize>> = HashMap::new(); |
| for (row_idx, old_partition, old_bucket) in &output.deletes { |
| delete_groups |
| .entry((old_partition.clone(), *old_bucket)) |
| .or_default() |
| .push(*row_idx); |
| } |
| for (key, row_indices) in delete_groups { |
| let sub_batch = Self::take_rows(batch, &row_indices)?; |
| let delete_batch = Self::add_value_kind_column(&sub_batch, 1)?; |
| result.push((key, delete_batch)); |
| } |
| } |
| |
| Ok(result) |
| } |
| |
| /// Extract rows from a batch by indices. |
| fn take_rows(batch: &RecordBatch, row_indices: &[usize]) -> Result<RecordBatch> { |
| if row_indices.len() == batch.num_rows() { |
| return Ok(batch.clone()); |
| } |
| let indices = arrow_array::UInt32Array::from( |
| row_indices.iter().map(|&i| i as u32).collect::<Vec<_>>(), |
| ); |
| let columns: Vec<Arc<dyn arrow_array::Array>> = batch |
| .columns() |
| .iter() |
| .map(|col| arrow_select::take::take(col.as_ref(), &indices, None)) |
| .collect::<std::result::Result<Vec<_>, _>>() |
| .map_err(|e| crate::Error::DataInvalid { |
| message: format!("Failed to take rows: {e}"), |
| source: None, |
| })?; |
| RecordBatch::try_new(batch.schema(), columns).map_err(|e| crate::Error::DataInvalid { |
| message: format!("Failed to create sub-batch: {e}"), |
| source: None, |
| }) |
| } |
| |
| /// Add a `_VALUE_KIND` column to a batch with the given value for all rows. |
| fn add_value_kind_column(batch: &RecordBatch, value_kind: i8) -> Result<RecordBatch> { |
| use arrow_array::Int8Array; |
| use arrow_schema::{DataType as ArrowDataType, Field as ArrowField}; |
| |
| let vk_array = Arc::new(Int8Array::from(vec![value_kind; batch.num_rows()])); |
| let vk_field = Arc::new(ArrowField::new( |
| crate::spec::VALUE_KIND_FIELD_NAME, |
| ArrowDataType::Int8, |
| false, |
| )); |
| |
| let mut fields = batch.schema().fields().to_vec(); |
| let mut columns: Vec<Arc<dyn arrow_array::Array>> = batch.columns().to_vec(); |
| fields.push(vk_field); |
| columns.push(vk_array); |
| |
| let schema = Arc::new(arrow_schema::Schema::new(fields)); |
| RecordBatch::try_new(schema, columns).map_err(|e| crate::Error::DataInvalid { |
| message: format!("Failed to add _VALUE_KIND column: {e}"), |
| source: None, |
| }) |
| } |
| |
| /// Write a batch directly to the writer for the given (partition, bucket). |
| async fn write_bucket( |
| &mut self, |
| partition_bytes: Vec<u8>, |
| bucket: i32, |
| batch: RecordBatch, |
| ) -> Result<()> { |
| let key = (partition_bytes, bucket); |
| if !self.partition_writers.contains_key(&key) { |
| self.create_writer(key.0.clone(), key.1).await?; |
| } |
| let writer = self.partition_writers.get_mut(&key).unwrap(); |
| writer.write(&batch).await |
| } |
| |
| /// Write multiple Arrow RecordBatches. |
| pub async fn write_arrow(&mut self, batches: &[RecordBatch]) -> Result<()> { |
| for batch in batches { |
| self.write_arrow_batch(batch).await?; |
| } |
| Ok(()) |
| } |
| |
| /// Close all writers and collect CommitMessages for use with TableCommit. |
| /// Writers are cleared after this call, allowing the TableWrite to be reused. |
| pub async fn prepare_commit(&mut self) -> Result<Vec<CommitMessage>> { |
| let writers: Vec<(PartitionBucketKey, FileWriter)> = |
| self.partition_writers.drain().collect(); |
| |
| let futures: Vec<_> = writers |
| .into_iter() |
| .map(|((partition_bytes, bucket), writer)| async move { |
| let files = writer.prepare_commit().await?; |
| Ok::<_, crate::Error>((partition_bytes, bucket, files)) |
| }) |
| .collect(); |
| |
| let results = futures::future::try_join_all(futures).await?; |
| |
| // Collect index files from bucket assigner |
| let file_io = self.table.file_io(); |
| let index_dir = format!("{}/index", self.table.location()); |
| let mut index_files_by_key = self |
| .bucket_assigner |
| .prepare_commit_index(file_io, &index_dir) |
| .await?; |
| |
| let mut messages = Vec::new(); |
| for (partition_bytes, bucket, files) in results { |
| let key = (partition_bytes.clone(), bucket); |
| let index_files = index_files_by_key.remove(&key).unwrap_or_default(); |
| if !files.data_files.is_empty() |
| || !files.changelog_files.is_empty() |
| || !index_files.is_empty() |
| { |
| let mut msg = CommitMessage::new(partition_bytes, bucket, files.data_files); |
| msg.new_changelog_files = files.changelog_files; |
| msg.new_index_files = index_files; |
| messages.push(msg); |
| } |
| } |
| // Emit index-only messages for (partition, bucket) pairs that had no data writer |
| // (e.g., old buckets where keys migrated away in cross-partition mode). |
| for ((partition_bytes, bucket), idx_files) in index_files_by_key { |
| if !idx_files.is_empty() { |
| let mut msg = CommitMessage::new(partition_bytes, bucket, vec![]); |
| msg.new_index_files = idx_files; |
| messages.push(msg); |
| } |
| } |
| Ok(messages) |
| } |
| |
| async fn create_writer(&mut self, partition_bytes: Vec<u8>, bucket: i32) -> Result<()> { |
| let partition_path = self.resolve_partition_path(&partition_bytes)?; |
| |
| let writer = if self.primary_key_indices.is_empty() { |
| self.create_append_writer(partition_path, bucket)? |
| } else if bucket == POSTPONE_BUCKET { |
| self.create_postpone_writer(partition_path, bucket) |
| } else { |
| self.create_kv_writer(partition_path, bucket, &partition_bytes) |
| .await? |
| }; |
| |
| self.partition_writers |
| .insert((partition_bytes, bucket), writer); |
| Ok(()) |
| } |
| |
| fn resolve_partition_path(&self, partition_bytes: &[u8]) -> Result<String> { |
| if self.partition_keys.is_empty() { |
| Ok(String::new()) |
| } else { |
| let row = BinaryRow::from_serialized_bytes(partition_bytes)?; |
| self.partition_computer.generate_partition_path(&row) |
| } |
| } |
| |
| /// Create an append-only writer for non-PK tables. |
| fn create_append_writer(&self, partition_path: String, bucket: i32) -> Result<FileWriter> { |
| if self.has_blob_fields { |
| let fields = self.table.schema().fields(); |
| let input_schema = build_target_arrow_schema(fields)?; |
| Ok(FileWriter::AppendBlob(AppendBlobFileWriter::new( |
| self.table.file_io().clone(), |
| self.table.location().to_string(), |
| partition_path, |
| bucket, |
| self.schema_id, |
| self.target_file_size, |
| self.blob_target_file_size, |
| self.file_compression.clone(), |
| self.file_compression_zstd_level, |
| self.write_buffer_size, |
| self.file_format.clone(), |
| &input_schema, |
| fields, |
| &self.blob_descriptor_fields, |
| ))) |
| } else { |
| Ok(FileWriter::Append(DataFileWriter::new( |
| self.table.file_io().clone(), |
| self.table.location().to_string(), |
| partition_path, |
| bucket, |
| self.schema_id, |
| self.target_file_size, |
| self.file_compression.clone(), |
| self.file_compression_zstd_level, |
| self.write_buffer_size, |
| self.file_format.clone(), |
| self.table.schema().fields().to_vec(), |
| Some(0), |
| None, |
| None, |
| ))) |
| } |
| } |
| |
| /// Create a postpone writer (KV format, no sorting/dedup, special file naming). |
| fn create_postpone_writer(&self, partition_path: String, bucket: i32) -> FileWriter { |
| let data_file_prefix = format!("data-u-{}-s-0-w-", self.commit_user); |
| FileWriter::Postpone(PostponeFileWriter::new( |
| self.table.file_io().clone(), |
| PostponeWriteConfig { |
| table_location: self.table.location().to_string(), |
| partition_path, |
| bucket, |
| schema_id: self.schema_id, |
| target_file_size: self.target_file_size, |
| file_compression: self.file_compression.clone(), |
| file_compression_zstd_level: self.file_compression_zstd_level, |
| write_buffer_size: self.write_buffer_size, |
| file_format: self.file_format.clone(), |
| data_file_prefix, |
| }, |
| )) |
| } |
| |
| /// Create a key-value writer for PK tables with normal buckets. |
| async fn create_kv_writer( |
| &mut self, |
| partition_path: String, |
| bucket: i32, |
| partition_bytes: &[u8], |
| ) -> Result<FileWriter> { |
| // Lazily scan partition sequence numbers on first writer creation per partition. |
| // Overwrite mode skips this — old data will be replaced, so seq starts at 0. |
| if !self.is_overwrite && !self.partition_seq_cache.contains_key(partition_bytes) { |
| let bucket_seq = |
| Self::scan_partition_sequence_numbers(&self.table, partition_bytes).await?; |
| self.partition_seq_cache |
| .insert(partition_bytes.to_vec(), bucket_seq); |
| } |
| let next_seq = self |
| .partition_seq_cache |
| .get(partition_bytes) |
| .and_then(|m| m.get(&bucket)) |
| .copied() |
| .unwrap_or(0); |
| |
| Ok(FileWriter::KeyValue(KeyValueFileWriter::new( |
| self.table.file_io().clone(), |
| KeyValueWriteConfig { |
| table_name: self.table.identifier().full_name(), |
| table_options: self.table.schema().options().clone(), |
| table_location: self.table.location().to_string(), |
| partition_path, |
| bucket, |
| schema_id: self.schema_id, |
| file_compression: self.file_compression.clone(), |
| file_compression_zstd_level: self.file_compression_zstd_level, |
| write_buffer_size: self.write_buffer_size, |
| file_format: self.file_format.clone(), |
| input_changelog: self.changelog_producer == ChangelogProducer::Input |
| && !self.is_overwrite, |
| changelog_file_prefix: self.changelog_file_prefix.clone(), |
| changelog_file_compression: self.changelog_file_compression.clone(), |
| changelog_file_format: self.changelog_file_format.clone(), |
| primary_key_indices: self.primary_key_indices.clone(), |
| primary_key_types: self.primary_key_types.clone(), |
| sequence_field_indices: self.sequence_field_indices.clone(), |
| merge_engine: self.merge_engine, |
| deletion_vectors_enabled: CoreOptions::new(self.table.schema().options()) |
| .deletion_vectors_enabled(), |
| }, |
| next_seq, |
| )?)) |
| } |
| } |
| |
| #[cfg(test)] |
| mod tests { |
| use super::*; |
| use crate::arrow::format::create_format_reader; |
| use crate::catalog::Identifier; |
| use crate::io::{FileIO, FileIOBuilder}; |
| use crate::spec::{ |
| bucket_dir_name, BigIntType, BinaryRowBuilder, BlobType, DataField, DataType, DecimalType, |
| FileKind, IndexManifest, IntType, LocalZonedTimestampType, Manifest, ManifestList, Schema, |
| TableSchema, TimestampType, TinyIntType, VarCharType, SEQUENCE_NUMBER_FIELD_ID, |
| SEQUENCE_NUMBER_FIELD_NAME, VALUE_KIND_FIELD_ID, VALUE_KIND_FIELD_NAME, |
| }; |
| use crate::table::{SnapshotManager, TableCommit}; |
| use arrow_array::RecordBatchReader as _; |
| use arrow_array::{Int32Array, Int64Array, Int8Array, StringArray}; |
| use arrow_schema::{ |
| DataType as ArrowDataType, Field as ArrowField, Schema as ArrowSchema, TimeUnit, |
| }; |
| use std::sync::Arc; |
| |
| fn test_file_io() -> FileIO { |
| FileIOBuilder::new("memory").build().unwrap() |
| } |
| |
| fn test_schema() -> TableSchema { |
| let schema = Schema::builder() |
| .column("id", DataType::Int(IntType::new())) |
| .column("value", DataType::Int(IntType::new())) |
| .build() |
| .unwrap(); |
| TableSchema::new(0, &schema) |
| } |
| |
| fn test_partitioned_schema() -> TableSchema { |
| let schema = Schema::builder() |
| .column("pt", DataType::VarChar(VarCharType::string_type())) |
| .column("id", DataType::Int(IntType::new())) |
| .partition_keys(["pt"]) |
| .build() |
| .unwrap(); |
| TableSchema::new(0, &schema) |
| } |
| |
| fn test_table(file_io: &FileIO, table_path: &str) -> Table { |
| Table::new( |
| file_io.clone(), |
| Identifier::new("default", "test_table"), |
| table_path.to_string(), |
| test_schema(), |
| None, |
| ) |
| } |
| |
| fn test_partitioned_table(file_io: &FileIO, table_path: &str) -> Table { |
| Table::new( |
| file_io.clone(), |
| Identifier::new("default", "test_table"), |
| table_path.to_string(), |
| test_partitioned_schema(), |
| None, |
| ) |
| } |
| |
| fn test_blob_table_schema() -> TableSchema { |
| let schema = Schema::builder() |
| .column("id", DataType::Int(IntType::new())) |
| .column("payload", DataType::Blob(BlobType::new())) |
| .option("data-evolution.enabled", "true") |
| .build() |
| .unwrap(); |
| TableSchema::new(0, &schema) |
| } |
| |
| async fn setup_dirs(file_io: &FileIO, table_path: &str) { |
| file_io |
| .mkdirs(&format!("{table_path}/snapshot/")) |
| .await |
| .unwrap(); |
| file_io |
| .mkdirs(&format!("{table_path}/manifest/")) |
| .await |
| .unwrap(); |
| } |
| |
| fn make_batch(ids: Vec<i32>, values: Vec<i32>) -> RecordBatch { |
| let schema = Arc::new(ArrowSchema::new(vec![ |
| ArrowField::new("id", ArrowDataType::Int32, false), |
| ArrowField::new("value", ArrowDataType::Int32, false), |
| ])); |
| RecordBatch::try_new( |
| schema, |
| vec![ |
| Arc::new(Int32Array::from(ids)), |
| Arc::new(Int32Array::from(values)), |
| ], |
| ) |
| .unwrap() |
| } |
| |
| fn make_batch_with_value_kind( |
| ids: Vec<i32>, |
| values: Vec<i32>, |
| value_kinds: Vec<i8>, |
| ) -> RecordBatch { |
| let schema = Arc::new(ArrowSchema::new(vec![ |
| ArrowField::new("id", ArrowDataType::Int32, false), |
| ArrowField::new("value", ArrowDataType::Int32, false), |
| ArrowField::new( |
| crate::spec::VALUE_KIND_FIELD_NAME, |
| ArrowDataType::Int8, |
| false, |
| ), |
| ])); |
| RecordBatch::try_new( |
| schema, |
| vec![ |
| Arc::new(Int32Array::from(ids)), |
| Arc::new(Int32Array::from(values)), |
| Arc::new(Int8Array::from(value_kinds)), |
| ], |
| ) |
| .unwrap() |
| } |
| |
| fn make_partitioned_batch_with_value_kind( |
| pts: Vec<&str>, |
| ids: Vec<i32>, |
| values: Vec<i32>, |
| value_kinds: Vec<i8>, |
| ) -> RecordBatch { |
| let schema = Arc::new(ArrowSchema::new(vec![ |
| ArrowField::new("pt", ArrowDataType::Utf8, false), |
| ArrowField::new("id", ArrowDataType::Int32, false), |
| ArrowField::new("value", ArrowDataType::Int32, false), |
| ArrowField::new( |
| crate::spec::VALUE_KIND_FIELD_NAME, |
| ArrowDataType::Int8, |
| false, |
| ), |
| ])); |
| RecordBatch::try_new( |
| schema, |
| vec![ |
| Arc::new(arrow_array::StringArray::from(pts)), |
| Arc::new(Int32Array::from(ids)), |
| Arc::new(Int32Array::from(values)), |
| Arc::new(Int8Array::from(value_kinds)), |
| ], |
| ) |
| .unwrap() |
| } |
| |
| fn physical_key_value_fields() -> Vec<DataField> { |
| vec![ |
| DataField::new( |
| SEQUENCE_NUMBER_FIELD_ID, |
| SEQUENCE_NUMBER_FIELD_NAME.to_string(), |
| DataType::BigInt(BigIntType::new()), |
| ), |
| DataField::new( |
| VALUE_KIND_FIELD_ID, |
| VALUE_KIND_FIELD_NAME.to_string(), |
| DataType::TinyInt(TinyIntType::new()), |
| ), |
| DataField::new(0, "id".to_string(), DataType::Int(IntType::new())), |
| DataField::new(1, "value".to_string(), DataType::Int(IntType::new())), |
| ] |
| } |
| |
| async fn read_physical_key_value_batches( |
| file_io: &FileIO, |
| file_path: &str, |
| file_size: i64, |
| ) -> Vec<RecordBatch> { |
| let format_reader = create_format_reader(file_path, false).unwrap(); |
| let input = file_io.new_input(file_path).unwrap(); |
| let file_reader = input.reader().await.unwrap(); |
| let read_fields = physical_key_value_fields(); |
| let stream = format_reader |
| .read_batch_stream( |
| Box::new(file_reader), |
| file_size as u64, |
| &read_fields, |
| None, |
| None, |
| None, |
| ) |
| .await |
| .unwrap(); |
| futures::TryStreamExt::try_collect(stream).await.unwrap() |
| } |
| |
| fn collect_i32(batches: &[RecordBatch], column: usize) -> Vec<i32> { |
| batches |
| .iter() |
| .flat_map(|batch| { |
| batch |
| .column(column) |
| .as_any() |
| .downcast_ref::<Int32Array>() |
| .unwrap() |
| .values() |
| .iter() |
| .copied() |
| }) |
| .collect() |
| } |
| |
| fn collect_i64(batches: &[RecordBatch], column: usize) -> Vec<i64> { |
| batches |
| .iter() |
| .flat_map(|batch| { |
| batch |
| .column(column) |
| .as_any() |
| .downcast_ref::<Int64Array>() |
| .unwrap() |
| .values() |
| .iter() |
| .copied() |
| }) |
| .collect() |
| } |
| |
| fn collect_i8(batches: &[RecordBatch], column: usize) -> Vec<i8> { |
| batches |
| .iter() |
| .flat_map(|batch| { |
| batch |
| .column(column) |
| .as_any() |
| .downcast_ref::<Int8Array>() |
| .unwrap() |
| .values() |
| .iter() |
| .copied() |
| }) |
| .collect() |
| } |
| |
| fn make_partitioned_batch(pts: Vec<&str>, ids: Vec<i32>) -> RecordBatch { |
| let schema = Arc::new(ArrowSchema::new(vec![ |
| ArrowField::new("pt", ArrowDataType::Utf8, false), |
| ArrowField::new("id", ArrowDataType::Int32, false), |
| ])); |
| RecordBatch::try_new( |
| schema, |
| vec![ |
| Arc::new(arrow_array::StringArray::from(pts)), |
| Arc::new(Int32Array::from(ids)), |
| ], |
| ) |
| .unwrap() |
| } |
| |
| #[tokio::test] |
| async fn test_write_and_commit() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_table_write"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let table = test_table(&file_io, table_path); |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| |
| let batch = make_batch(vec![1, 2, 3], vec![10, 20, 30]); |
| table_write.write_arrow_batch(&batch).await.unwrap(); |
| |
| let messages = table_write.prepare_commit().await.unwrap(); |
| assert_eq!(messages.len(), 1); |
| assert_eq!(messages[0].bucket, 0); |
| assert_eq!(messages[0].new_files.len(), 1); |
| assert_eq!(messages[0].new_files[0].row_count, 3); |
| |
| // Commit and verify snapshot |
| let commit = TableCommit::new(table, "test-user".to_string()); |
| commit.commit(messages).await.unwrap(); |
| |
| let snap_manager = SnapshotManager::new(file_io.clone(), table_path.to_string()); |
| let snapshot = snap_manager.get_latest_snapshot().await.unwrap().unwrap(); |
| assert_eq!(snapshot.id(), 1); |
| assert_eq!(snapshot.total_record_count(), Some(3)); |
| } |
| |
| #[tokio::test] |
| async fn test_table_write_row_format_roundtrip() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_table_write_row_format"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let schema = Schema::builder() |
| .column("id", DataType::Int(IntType::new())) |
| .column("value", DataType::Int(IntType::new())) |
| .option("file.format", "row") |
| .build() |
| .unwrap(); |
| let table = Table::new( |
| file_io.clone(), |
| Identifier::new("default", "test_row_table"), |
| table_path.to_string(), |
| TableSchema::new(0, &schema), |
| None, |
| ); |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| table_write |
| .write_arrow_batch(&make_batch(vec![1, 2, 3], vec![10, 20, 30])) |
| .await |
| .unwrap(); |
| |
| let messages = table_write.prepare_commit().await.unwrap(); |
| assert_eq!(messages.len(), 1); |
| assert_eq!(messages[0].new_files.len(), 1); |
| let data_file = &messages[0].new_files[0]; |
| assert!(data_file.file_name.ends_with(".row")); |
| |
| let file_path = format!( |
| "{}/{}/{}", |
| table_path, |
| bucket_dir_name(messages[0].bucket), |
| data_file.file_name |
| ); |
| let format_reader = create_format_reader(&file_path, false).unwrap(); |
| let input = file_io.new_input(&file_path).unwrap(); |
| let stream = format_reader |
| .read_batch_stream( |
| Box::new(input.reader().await.unwrap()), |
| data_file.file_size as u64, |
| table.schema().fields(), |
| None, |
| None, |
| None, |
| ) |
| .await |
| .unwrap(); |
| let batches: Vec<RecordBatch> = futures::TryStreamExt::try_collect(stream).await.unwrap(); |
| |
| assert_eq!(collect_i32(&batches, 0), vec![1, 2, 3]); |
| assert_eq!(collect_i32(&batches, 1), vec![10, 20, 30]); |
| } |
| |
| #[test] |
| fn test_allows_append_blob_table() { |
| let table = Table::new( |
| test_file_io(), |
| Identifier::new("default", "test_blob_table"), |
| "memory:/test_blob_table".to_string(), |
| test_blob_table_schema(), |
| None, |
| ); |
| |
| assert!(TableWrite::new(&table, "test-user".to_string()).is_ok()); |
| } |
| |
| #[tokio::test] |
| async fn test_blob_write_and_commit() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_blob_write"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let table = Table::new( |
| file_io.clone(), |
| Identifier::new("default", "test_blob_table"), |
| table_path.to_string(), |
| test_blob_table_schema(), |
| None, |
| ); |
| |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| |
| let schema = Arc::new(ArrowSchema::new(vec![ |
| ArrowField::new("id", ArrowDataType::Int32, false), |
| ArrowField::new("payload", ArrowDataType::Binary, true), |
| ])); |
| let batch = RecordBatch::try_new( |
| schema, |
| vec![ |
| Arc::new(Int32Array::from(vec![1, 2, 3])), |
| Arc::new(arrow_array::BinaryArray::from(vec![ |
| Some(b"hello" as &[u8]), |
| None, |
| Some(b"world"), |
| ])), |
| ], |
| ) |
| .unwrap(); |
| |
| table_write.write_arrow_batch(&batch).await.unwrap(); |
| let messages = table_write.prepare_commit().await.unwrap(); |
| |
| assert_eq!(messages.len(), 1); |
| // Should have 2 files: 1 parquet (normal cols) + 1 blob (payload) |
| assert_eq!(messages[0].new_files.len(), 2); |
| |
| let parquet_files: Vec<_> = messages[0] |
| .new_files |
| .iter() |
| .filter(|f| f.file_name.ends_with(".parquet")) |
| .collect(); |
| let blob_files: Vec<_> = messages[0] |
| .new_files |
| .iter() |
| .filter(|f| f.file_name.ends_with(".blob")) |
| .collect(); |
| assert_eq!(parquet_files.len(), 1); |
| assert_eq!(blob_files.len(), 1); |
| |
| assert_eq!(parquet_files[0].row_count, 3); |
| assert_eq!(blob_files[0].row_count, 3); |
| assert_eq!(blob_files[0].write_cols, Some(vec!["payload".to_string()])); |
| |
| // Commit and verify snapshot |
| let commit = TableCommit::new(table, "test-user".to_string()); |
| commit.commit(messages).await.unwrap(); |
| |
| let snap_manager = SnapshotManager::new(file_io.clone(), table_path.to_string()); |
| let snapshot = snap_manager.get_latest_snapshot().await.unwrap().unwrap(); |
| assert_eq!(snapshot.id(), 1); |
| } |
| |
| #[test] |
| fn test_allows_partial_update_fixed_bucket_table() { |
| let table = Table::new( |
| test_file_io(), |
| Identifier::new("default", "test_partial_update_table"), |
| "memory:/test_partial_update_table".to_string(), |
| TableSchema::new( |
| 0, |
| &Schema::builder() |
| .column("id", DataType::Int(IntType::new())) |
| .column("value", DataType::Int(IntType::new())) |
| .primary_key(["id"]) |
| .option("bucket", "1") |
| .option("merge-engine", "partial-update") |
| .build() |
| .unwrap(), |
| ), |
| None, |
| ); |
| |
| TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| } |
| |
| #[tokio::test] |
| async fn test_allows_partial_update_dynamic_bucket_table() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_partial_update_dynamic_bucket_table"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let table = Table::new( |
| file_io, |
| Identifier::new("default", "test_partial_update_dynamic_bucket_table"), |
| table_path.to_string(), |
| TableSchema::new( |
| 0, |
| &Schema::builder() |
| .column("id", DataType::Int(IntType::new())) |
| .column("value", DataType::Int(IntType::new())) |
| .primary_key(["id"]) |
| .option("merge-engine", "partial-update") |
| .build() |
| .unwrap(), |
| ), |
| None, |
| ); |
| |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| table_write |
| .write_arrow_batch(&make_batch(vec![1], vec![10])) |
| .await |
| .unwrap(); |
| } |
| |
| #[tokio::test] |
| async fn test_allows_aggregation_dynamic_bucket_table() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_aggregation_dynamic_bucket_table"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let table = Table::new( |
| file_io, |
| Identifier::new("default", "test_aggregation_dynamic_bucket_table"), |
| table_path.to_string(), |
| TableSchema::new( |
| 0, |
| &Schema::builder() |
| .column("id", DataType::Int(IntType::new())) |
| .column("value", DataType::Int(IntType::new())) |
| .primary_key(["id"]) |
| .option("bucket", "-1") |
| .option("merge-engine", "aggregation") |
| .option("fields.value.aggregate-function", "sum") |
| .build() |
| .unwrap(), |
| ), |
| None, |
| ); |
| |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| assert!(matches!( |
| table_write.bucket_assigner, |
| BucketAssignerEnum::Dynamic(_) |
| )); |
| table_write |
| .write_arrow_batch(&make_batch(vec![1], vec![10])) |
| .await |
| .unwrap(); |
| } |
| |
| #[tokio::test] |
| async fn test_rejects_partial_update_with_deletion_vectors_when_creating_writer() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_partial_update_dv_table"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let table = Table::new( |
| file_io, |
| Identifier::new("default", "test_partial_update_dv_table"), |
| table_path.to_string(), |
| TableSchema::new( |
| 0, |
| &Schema::builder() |
| .column("id", DataType::Int(IntType::new())) |
| .column("value", DataType::Int(IntType::new())) |
| .primary_key(["id"]) |
| .option("bucket", "1") |
| .option("merge-engine", "partial-update") |
| .option("deletion-vectors.enabled", "true") |
| .build() |
| .unwrap(), |
| ), |
| None, |
| ); |
| |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| let err = table_write |
| .write_arrow_batch(&make_batch(vec![1], vec![10])) |
| .await |
| .unwrap_err(); |
| assert!( |
| matches!(err, crate::Error::Unsupported { message } if message.contains("deletion-vectors.enabled=true")) |
| ); |
| } |
| |
| #[tokio::test] |
| async fn test_write_partitioned() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_table_write_partitioned"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let table = test_partitioned_table(&file_io, table_path); |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| |
| let batch = make_partitioned_batch(vec!["a", "b", "a"], vec![1, 2, 3]); |
| table_write.write_arrow_batch(&batch).await.unwrap(); |
| |
| let messages = table_write.prepare_commit().await.unwrap(); |
| // Should have 2 commit messages (one per partition) |
| assert_eq!(messages.len(), 2); |
| |
| let total_rows: i64 = messages |
| .iter() |
| .flat_map(|m| &m.new_files) |
| .map(|f| f.row_count) |
| .sum(); |
| assert_eq!(total_rows, 3); |
| |
| // Commit and verify |
| let commit = TableCommit::new(table, "test-user".to_string()); |
| commit.commit(messages).await.unwrap(); |
| |
| let snap_manager = SnapshotManager::new(file_io.clone(), table_path.to_string()); |
| let snapshot = snap_manager.get_latest_snapshot().await.unwrap().unwrap(); |
| assert_eq!(snapshot.id(), 1); |
| assert_eq!(snapshot.total_record_count(), Some(3)); |
| } |
| |
| #[tokio::test] |
| async fn test_write_empty_batch() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_table_write_empty"; |
| let table = test_table(&file_io, table_path); |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| |
| let batch = make_batch(vec![], vec![]); |
| table_write.write_arrow_batch(&batch).await.unwrap(); |
| |
| let messages = table_write.prepare_commit().await.unwrap(); |
| assert!(messages.is_empty()); |
| } |
| |
| #[tokio::test] |
| async fn test_prepare_commit_reusable() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_table_write_reuse"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let table = test_table(&file_io, table_path); |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| |
| // First write + prepare_commit |
| table_write |
| .write_arrow_batch(&make_batch(vec![1, 2], vec![10, 20])) |
| .await |
| .unwrap(); |
| let messages1 = table_write.prepare_commit().await.unwrap(); |
| assert_eq!(messages1.len(), 1); |
| assert_eq!(messages1[0].new_files[0].row_count, 2); |
| |
| // Second write + prepare_commit (reuse) |
| table_write |
| .write_arrow_batch(&make_batch(vec![3, 4, 5], vec![30, 40, 50])) |
| .await |
| .unwrap(); |
| let messages2 = table_write.prepare_commit().await.unwrap(); |
| assert_eq!(messages2.len(), 1); |
| assert_eq!(messages2[0].new_files[0].row_count, 3); |
| |
| // Empty prepare_commit is fine |
| let messages3 = table_write.prepare_commit().await.unwrap(); |
| assert!(messages3.is_empty()); |
| } |
| |
| #[tokio::test] |
| async fn test_write_multiple_batches() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_table_write_multi"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let table = test_table(&file_io, table_path); |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| |
| table_write |
| .write_arrow_batch(&make_batch(vec![1, 2], vec![10, 20])) |
| .await |
| .unwrap(); |
| table_write |
| .write_arrow_batch(&make_batch(vec![3, 4], vec![30, 40])) |
| .await |
| .unwrap(); |
| |
| let messages = table_write.prepare_commit().await.unwrap(); |
| assert_eq!(messages.len(), 1); |
| // Multiple batches accumulate into a single file |
| assert_eq!(messages[0].new_files.len(), 1); |
| |
| let total_rows: i64 = messages[0].new_files.iter().map(|f| f.row_count).sum(); |
| assert_eq!(total_rows, 4); |
| } |
| |
| fn test_bucketed_schema() -> TableSchema { |
| let schema = Schema::builder() |
| .column("id", DataType::Int(IntType::new())) |
| .column("value", DataType::Int(IntType::new())) |
| .option("bucket", "4") |
| .option("bucket-key", "id") |
| .build() |
| .unwrap(); |
| TableSchema::new(0, &schema) |
| } |
| |
| fn test_bucketed_table(file_io: &FileIO, table_path: &str) -> Table { |
| Table::new( |
| file_io.clone(), |
| Identifier::new("default", "test_table"), |
| table_path.to_string(), |
| test_bucketed_schema(), |
| None, |
| ) |
| } |
| |
| /// Build a batch where the bucket-key column ("id") is nullable. |
| fn make_nullable_id_batch(ids: Vec<Option<i32>>, values: Vec<i32>) -> RecordBatch { |
| let schema = Arc::new(ArrowSchema::new(vec![ |
| ArrowField::new("id", ArrowDataType::Int32, true), |
| ArrowField::new("value", ArrowDataType::Int32, false), |
| ])); |
| RecordBatch::try_new( |
| schema, |
| vec![ |
| Arc::new(Int32Array::from(ids)), |
| Arc::new(Int32Array::from(values)), |
| ], |
| ) |
| .unwrap() |
| } |
| |
| #[tokio::test] |
| async fn test_mod_bucket_function_routes_by_floor_mod() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_mod_bucket_function"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let schema = Schema::builder() |
| .column("id", DataType::Int(IntType::new())) |
| .column("value", DataType::Int(IntType::new())) |
| .option("bucket", "5") |
| .option("bucket-key", "id") |
| .option("bucket-function.type", "mod") |
| .build() |
| .unwrap(); |
| let table = Table::new( |
| file_io.clone(), |
| Identifier::new("default", "test_table"), |
| table_path.to_string(), |
| TableSchema::new(0, &schema), |
| None, |
| ); |
| |
| let fields = table.schema().fields().to_vec(); |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| let output = table_write |
| .bucket_assigner |
| .assign_batch(&make_batch(vec![-3, 17, 5], vec![10, 20, 30]), &fields) |
| .await |
| .unwrap(); |
| |
| assert_eq!(output.buckets, vec![2, 2, 0]); |
| } |
| |
| #[tokio::test] |
| async fn test_mod_bucket_function_rejects_non_integral_key() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_mod_bucket_invalid"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let schema = Schema::builder() |
| .column("id", DataType::VarChar(VarCharType::string_type())) |
| .column("value", DataType::Int(IntType::new())) |
| .option("bucket", "5") |
| .option("bucket-key", "id") |
| .option("bucket-function.type", "mod") |
| .build() |
| .unwrap(); |
| let table = Table::new( |
| file_io, |
| Identifier::new("default", "test_table"), |
| table_path.to_string(), |
| TableSchema::new(0, &schema), |
| None, |
| ); |
| |
| let err = match TableWrite::new(&table, "test-user".to_string()) { |
| Ok(_) => panic!("expected mod bucket function to reject non-integral key"), |
| Err(err) => err, |
| }; |
| assert!( |
| matches!(err, crate::Error::ConfigInvalid { ref message } if message.contains("INT or BIGINT")), |
| "expected ConfigInvalid for non-integral mod bucket key, got {err:?}" |
| ); |
| } |
| |
| #[tokio::test] |
| async fn test_hive_bucket_function_routes_by_hive_hash() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_hive_bucket_function"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let schema = Schema::builder() |
| .column("id", DataType::Int(IntType::new())) |
| .column("name", DataType::VarChar(VarCharType::string_type())) |
| .option("bucket", "8") |
| .option("bucket-key", "id,name") |
| .option("bucket-function.type", "hive") |
| .build() |
| .unwrap(); |
| let table = Table::new( |
| file_io.clone(), |
| Identifier::new("default", "test_table"), |
| table_path.to_string(), |
| TableSchema::new(0, &schema), |
| None, |
| ); |
| |
| let arrow_schema = Arc::new(ArrowSchema::new(vec![ |
| ArrowField::new("id", ArrowDataType::Int32, false), |
| ArrowField::new("name", ArrowDataType::Utf8, false), |
| ])); |
| let batch = RecordBatch::try_new( |
| arrow_schema, |
| vec![ |
| Arc::new(Int32Array::from(vec![7])), |
| Arc::new(StringArray::from(vec!["hello"])), |
| ], |
| ) |
| .unwrap(); |
| |
| let fields = table.schema().fields().to_vec(); |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| let output = table_write |
| .bucket_assigner |
| .assign_batch(&batch, &fields) |
| .await |
| .unwrap(); |
| |
| assert_eq!(output.buckets, vec![3]); |
| } |
| |
| #[tokio::test] |
| async fn test_write_bucketed_with_null_bucket_key() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_table_write_null_bk"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let table = test_bucketed_table(&file_io, table_path); |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| |
| // Row with NULL bucket key should not panic |
| let batch = make_nullable_id_batch(vec![None, Some(1), None], vec![10, 20, 30]); |
| table_write.write_arrow_batch(&batch).await.unwrap(); |
| |
| let messages = table_write.prepare_commit().await.unwrap(); |
| let total_rows: i64 = messages |
| .iter() |
| .flat_map(|m| &m.new_files) |
| .map(|f| f.row_count) |
| .sum(); |
| assert_eq!(total_rows, 3); |
| } |
| |
| #[tokio::test] |
| async fn test_null_bucket_key_routes_consistently() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_table_write_null_bk_consistent"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let table = test_bucketed_table(&file_io, table_path); |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| |
| // Two NULLs should land in the same bucket |
| let batch = make_nullable_id_batch(vec![None, None], vec![10, 20]); |
| table_write.write_arrow_batch(&batch).await.unwrap(); |
| |
| let messages = table_write.prepare_commit().await.unwrap(); |
| // Both NULL-key rows must be in the same (partition, bucket) group |
| let null_bucket_rows: i64 = messages |
| .iter() |
| .flat_map(|m| &m.new_files) |
| .map(|f| f.row_count) |
| .sum(); |
| assert_eq!(null_bucket_rows, 2); |
| // All NULL-key rows go to exactly one bucket |
| assert_eq!(messages.len(), 1); |
| } |
| |
| #[tokio::test] |
| async fn test_null_vs_nonnull_bucket_key_differ() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_table_write_null_vs_nonnull"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let table = test_bucketed_table(&file_io, table_path); |
| |
| // Compute bucket for NULL key |
| let fields = table.schema().fields().to_vec(); |
| let mut tw = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| |
| let batch_null = make_nullable_id_batch(vec![None], vec![10]); |
| let out_null = tw |
| .bucket_assigner |
| .assign_batch(&batch_null, &fields) |
| .await |
| .unwrap(); |
| let bucket_null = out_null.buckets[0]; |
| |
| // Compute bucket for key = 0 (the value a null field's fixed bytes happen to be) |
| let batch_zero = make_nullable_id_batch(vec![Some(0)], vec![20]); |
| let out_zero = tw |
| .bucket_assigner |
| .assign_batch(&batch_zero, &fields) |
| .await |
| .unwrap(); |
| let bucket_zero = out_zero.buckets[0]; |
| |
| // A NULL bucket key must produce a BinaryRow with the null bit set, |
| // which hashes differently from a non-null 0 value. |
| // (With 4 buckets they could theoretically collide, but the hash codes differ.) |
| let mut builder_null = BinaryRowBuilder::new(1); |
| builder_null.set_null_at(0); |
| let hash_null = builder_null.build().hash_code(); |
| |
| let mut builder_zero = BinaryRowBuilder::new(1); |
| builder_zero.write_int(0, 0); |
| let hash_zero = builder_zero.build().hash_code(); |
| |
| assert_ne!(hash_null, hash_zero, "NULL and 0 should hash differently"); |
| // If hashes differ, buckets should differ (with 4 buckets, very likely) |
| // But we verify the hash difference is the important invariant |
| let _ = (bucket_null, bucket_zero); |
| } |
| |
| /// Mirrors Java's testUnCompactDecimalAndTimestampNullValueBucketNumber. |
| /// Non-compact types (Decimal(38,18), LocalZonedTimestamp(6), Timestamp(6)) |
| /// use variable-length encoding in BinaryRow — NULL handling must still work. |
| #[tokio::test] |
| async fn test_non_compact_null_bucket_key() { |
| let file_io = test_file_io(); |
| |
| let bucket_cols = ["d", "ltz", "ntz"]; |
| let total_buckets = 16; |
| |
| for bucket_col in &bucket_cols { |
| let table_path = format!("memory:/test_null_bk_{bucket_col}"); |
| setup_dirs(&file_io, &table_path).await; |
| |
| let schema = Schema::builder() |
| .column("d", DataType::Decimal(DecimalType::new(38, 18).unwrap())) |
| .column( |
| "ltz", |
| DataType::LocalZonedTimestamp(LocalZonedTimestampType::new(6).unwrap()), |
| ) |
| .column("ntz", DataType::Timestamp(TimestampType::new(6).unwrap())) |
| .column("k", DataType::Int(IntType::new())) |
| .option("bucket", total_buckets.to_string()) |
| .option("bucket-key", *bucket_col) |
| .build() |
| .unwrap(); |
| let table_schema = TableSchema::new(0, &schema); |
| let table = Table::new( |
| file_io.clone(), |
| Identifier::new("default", "test_table"), |
| table_path.to_string(), |
| table_schema, |
| None, |
| ); |
| |
| let mut tw = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| let fields = table.schema().fields().to_vec(); |
| |
| // Build a batch: d=NULL, ltz=NULL, ntz=NULL, k=1 |
| let arrow_schema = Arc::new(ArrowSchema::new(vec![ |
| ArrowField::new("d", ArrowDataType::Decimal128(38, 18), true), |
| ArrowField::new( |
| "ltz", |
| ArrowDataType::Timestamp(TimeUnit::Microsecond, Some("UTC".into())), |
| true, |
| ), |
| ArrowField::new( |
| "ntz", |
| ArrowDataType::Timestamp(TimeUnit::Microsecond, None), |
| true, |
| ), |
| ArrowField::new("k", ArrowDataType::Int32, false), |
| ])); |
| let batch = RecordBatch::try_new( |
| arrow_schema, |
| vec![ |
| Arc::new( |
| arrow_array::Decimal128Array::from(vec![None::<i128>]) |
| .with_precision_and_scale(38, 18) |
| .unwrap(), |
| ), |
| Arc::new( |
| arrow_array::TimestampMicrosecondArray::from(vec![None::<i64>]) |
| .with_timezone("UTC"), |
| ), |
| Arc::new(arrow_array::TimestampMicrosecondArray::from(vec![ |
| None::<i64>, |
| ])), |
| Arc::new(Int32Array::from(vec![1])), |
| ], |
| ) |
| .unwrap(); |
| |
| let batch_output = tw |
| .bucket_assigner |
| .assign_batch(&batch, &fields) |
| .await |
| .unwrap(); |
| let bucket = batch_output.buckets[0]; |
| |
| // Expected: BinaryRow with 1 field, null at pos 0 |
| let mut builder = BinaryRowBuilder::new(1); |
| builder.set_null_at(0); |
| let expected_bucket = (builder.build().hash_code() % total_buckets).wrapping_abs(); |
| |
| assert_eq!( |
| bucket, expected_bucket, |
| "NULL bucket-key '{bucket_col}' should produce bucket {expected_bucket}, got {bucket}" |
| ); |
| } |
| } |
| |
| #[tokio::test] |
| async fn test_write_rolling_on_target_file_size() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_table_write_rolling"; |
| setup_dirs(&file_io, table_path).await; |
| |
| // Create table with very small target-file-size to trigger rolling |
| let schema = Schema::builder() |
| .column("id", DataType::Int(IntType::new())) |
| .column("value", DataType::Int(IntType::new())) |
| .option("target-file-size", "1b") |
| .build() |
| .unwrap(); |
| let table_schema = TableSchema::new(0, &schema); |
| let table = Table::new( |
| file_io.clone(), |
| Identifier::new("default", "test_table"), |
| table_path.to_string(), |
| table_schema, |
| None, |
| ); |
| |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| |
| // Write multiple batches — each should roll to a new file |
| table_write |
| .write_arrow_batch(&make_batch(vec![1, 2], vec![10, 20])) |
| .await |
| .unwrap(); |
| table_write |
| .write_arrow_batch(&make_batch(vec![3, 4], vec![30, 40])) |
| .await |
| .unwrap(); |
| |
| let messages = table_write.prepare_commit().await.unwrap(); |
| assert_eq!(messages.len(), 1); |
| // With 1-byte target, each batch should produce a separate file |
| assert_eq!(messages[0].new_files.len(), 2); |
| |
| let total_rows: i64 = messages[0].new_files.iter().map(|f| f.row_count).sum(); |
| assert_eq!(total_rows, 4); |
| } |
| |
| #[cfg(feature = "vortex")] |
| #[tokio::test] |
| async fn test_vortex_write_rolling_on_target_file_size() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_vortex_table_write_rolling"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let schema = Schema::builder() |
| .column("id", DataType::Int(IntType::new())) |
| .column("value", DataType::Int(IntType::new())) |
| .option("target-file-size", "1b") |
| .option("file.format", "vortex") |
| .build() |
| .unwrap(); |
| let table_schema = TableSchema::new(0, &schema); |
| let table = Table::new( |
| file_io.clone(), |
| Identifier::new("default", "test_table"), |
| table_path.to_string(), |
| table_schema, |
| None, |
| ); |
| |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| |
| table_write |
| .write_arrow_batch(&make_batch(vec![1, 2], vec![10, 20])) |
| .await |
| .unwrap(); |
| table_write |
| .write_arrow_batch(&make_batch(vec![3, 4], vec![30, 40])) |
| .await |
| .unwrap(); |
| |
| let messages = table_write.prepare_commit().await.unwrap(); |
| assert_eq!(messages.len(), 1); |
| assert_eq!(messages[0].new_files.len(), 2); |
| |
| let total_rows: i64 = messages[0].new_files.iter().map(|f| f.row_count).sum(); |
| assert_eq!(total_rows, 4); |
| } |
| |
| // ----------------------------------------------------------------------- |
| // Primary-key table write tests |
| // ----------------------------------------------------------------------- |
| |
| fn test_pk_schema() -> TableSchema { |
| let schema = Schema::builder() |
| .column("id", DataType::Int(IntType::new())) |
| .column("value", DataType::Int(IntType::new())) |
| .primary_key(["id"]) |
| .option("bucket", "1") |
| .build() |
| .unwrap(); |
| TableSchema::new(0, &schema) |
| } |
| |
| fn test_pk_table(file_io: &FileIO, table_path: &str) -> Table { |
| Table::new( |
| file_io.clone(), |
| Identifier::new("default", "test_pk_table"), |
| table_path.to_string(), |
| test_pk_schema(), |
| None, |
| ) |
| } |
| |
| fn pk_changelog_schema(options: &[(&str, &str)]) -> TableSchema { |
| let mut builder = Schema::builder() |
| .column("id", DataType::Int(IntType::new())) |
| .column("value", DataType::Int(IntType::new())) |
| .primary_key(["id"]) |
| .option("bucket", "1"); |
| for (key, value) in options { |
| builder = builder.option(*key, *value); |
| } |
| TableSchema::new(0, &builder.build().unwrap()) |
| } |
| |
| fn loaded_first_row_schema_with_changelog_producer(producer: &str) -> TableSchema { |
| let schema = Schema::builder() |
| .column("id", DataType::Int(IntType::new())) |
| .column("value", DataType::Int(IntType::new())) |
| .primary_key(["id"]) |
| .option("merge-engine", "first-row") |
| .option("changelog-producer", "lookup") |
| .build() |
| .unwrap(); |
| let table_schema = TableSchema::new(0, &schema); |
| let mut value = serde_json::to_value(&table_schema).unwrap(); |
| value["options"]["changelog-producer"] = serde_json::Value::String(producer.to_string()); |
| serde_json::from_value(value).unwrap() |
| } |
| |
| fn ordinary_dynamic_pk_changelog_schema() -> TableSchema { |
| let schema = Schema::builder() |
| .column("pt", DataType::VarChar(VarCharType::string_type())) |
| .column("id", DataType::Int(IntType::new())) |
| .column("value", DataType::Int(IntType::new())) |
| .partition_keys(["pt"]) |
| .primary_key(["pt", "id"]) |
| .option("changelog-producer", "input") |
| .build() |
| .unwrap(); |
| TableSchema::new(0, &schema) |
| } |
| |
| #[test] |
| fn test_table_write_rejects_loaded_first_row_with_incompatible_changelog_producer() { |
| let file_io = test_file_io(); |
| |
| for producer in ["input", "full-compaction"] { |
| let table = Table::new( |
| file_io.clone(), |
| Identifier::new("default", "test_first_row_changelog"), |
| "memory:/test_first_row_changelog".to_string(), |
| loaded_first_row_schema_with_changelog_producer(producer), |
| None, |
| ); |
| |
| let err = match TableWrite::new(&table, "test-user".to_string()) { |
| Ok(_) => panic!( |
| "first-row should reject changelog-producer={producer} during write setup" |
| ), |
| Err(err) => err, |
| }; |
| |
| assert!( |
| matches!(err, crate::Error::Unsupported { ref message } |
| if message.contains("incompatible table options") |
| && message.contains("merge-engine=first-row") |
| && message.contains(producer)), |
| "first-row runtime guard should reject changelog-producer={producer}, got {err:?}" |
| ); |
| } |
| } |
| |
| #[tokio::test] |
| async fn test_input_changelog_writes_raw_rows_separately_from_data_rows() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_input_changelog_duplicate_pk"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let table = Table::new( |
| file_io.clone(), |
| Identifier::new("default", "test_input_changelog"), |
| table_path.to_string(), |
| pk_changelog_schema(&[ |
| ("changelog-producer", "input"), |
| ("changelog-file.prefix", "custom-changelog-"), |
| ]), |
| None, |
| ); |
| |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| table_write |
| .write_arrow_batch(&make_batch(vec![1, 1], vec![10, 20])) |
| .await |
| .unwrap(); |
| |
| let messages = table_write.prepare_commit().await.unwrap(); |
| assert_eq!(messages.len(), 1); |
| assert_eq!(messages[0].new_files.len(), 1); |
| assert_eq!(messages[0].new_files[0].row_count, 1); |
| assert_eq!(messages[0].new_changelog_files.len(), 1); |
| assert_eq!(messages[0].new_changelog_files[0].row_count, 2); |
| assert!(messages[0].new_files[0].file_name.starts_with("data-")); |
| assert!(messages[0].new_changelog_files[0] |
| .file_name |
| .starts_with("custom-changelog-")); |
| |
| let bucket_dir = bucket_dir_name(messages[0].bucket); |
| let data_file = &messages[0].new_files[0]; |
| let data_file_path = format!("{table_path}/{bucket_dir}/{}", data_file.file_name); |
| let data_batches = |
| read_physical_key_value_batches(&file_io, &data_file_path, data_file.file_size).await; |
| assert_eq!(collect_i64(&data_batches, 0), vec![1]); |
| assert_eq!(collect_i8(&data_batches, 1), vec![0]); |
| assert_eq!(collect_i32(&data_batches, 2), vec![1]); |
| assert_eq!(collect_i32(&data_batches, 3), vec![20]); |
| |
| let changelog_file = &messages[0].new_changelog_files[0]; |
| let changelog_file_path = format!("{table_path}/{bucket_dir}/{}", changelog_file.file_name); |
| let changelog_batches = read_physical_key_value_batches( |
| &file_io, |
| &changelog_file_path, |
| changelog_file.file_size, |
| ) |
| .await; |
| assert_eq!(collect_i64(&changelog_batches, 0), vec![0, 1]); |
| assert_eq!(collect_i8(&changelog_batches, 1), vec![0, 0]); |
| assert_eq!(collect_i32(&changelog_batches, 2), vec![1, 1]); |
| assert_eq!(collect_i32(&changelog_batches, 3), vec![10, 20]); |
| } |
| |
| #[tokio::test] |
| async fn test_input_changelog_metadata_counts_retract_rows() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_input_changelog_retract_rows"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let table = Table::new( |
| file_io, |
| Identifier::new("default", "test_input_changelog"), |
| table_path.to_string(), |
| pk_changelog_schema(&[("changelog-producer", "input")]), |
| None, |
| ); |
| |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| table_write |
| .write_arrow_batch(&make_batch_with_value_kind( |
| vec![1, 2, 3], |
| vec![10, 20, 30], |
| vec![0, 1, 3], |
| )) |
| .await |
| .unwrap(); |
| |
| let messages = table_write.prepare_commit().await.unwrap(); |
| assert_eq!(messages[0].new_files[0].delete_row_count, Some(2)); |
| assert_eq!(messages[0].new_changelog_files[0].delete_row_count, Some(2)); |
| } |
| |
| #[tokio::test] |
| async fn test_input_changelog_rejects_invalid_value_kind() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_input_changelog_invalid_value_kind"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let table = Table::new( |
| file_io, |
| Identifier::new("default", "test_input_changelog"), |
| table_path.to_string(), |
| pk_changelog_schema(&[("changelog-producer", "input")]), |
| None, |
| ); |
| |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| table_write |
| .write_arrow_batch(&make_batch_with_value_kind(vec![1], vec![10], vec![4])) |
| .await |
| .unwrap(); |
| |
| let err = table_write.prepare_commit().await.unwrap_err(); |
| assert!( |
| matches!(err, crate::Error::DataInvalid { message, .. } if message.contains("Invalid RowKind value")) |
| ); |
| } |
| |
| #[tokio::test] |
| async fn test_input_changelog_commit_writes_changelog_manifest_metadata() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_input_changelog_commit"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let table = Table::new( |
| file_io.clone(), |
| Identifier::new("default", "test_input_changelog"), |
| table_path.to_string(), |
| pk_changelog_schema(&[("changelog-producer", "input")]), |
| None, |
| ); |
| |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| table_write |
| .write_arrow_batch(&make_batch(vec![1, 1], vec![10, 20])) |
| .await |
| .unwrap(); |
| let messages = table_write.prepare_commit().await.unwrap(); |
| let data_file_name = messages[0].new_files[0].file_name.clone(); |
| let changelog_file_name = messages[0].new_changelog_files[0].file_name.clone(); |
| |
| let commit = TableCommit::new(table, "test-user".to_string()); |
| commit.commit(messages).await.unwrap(); |
| |
| let snap_manager = SnapshotManager::new(file_io.clone(), table_path.to_string()); |
| let snapshot = snap_manager.get_latest_snapshot().await.unwrap().unwrap(); |
| assert_eq!(snapshot.total_record_count(), Some(1)); |
| assert_eq!(snapshot.delta_record_count(), Some(1)); |
| assert_eq!(snapshot.changelog_record_count(), Some(2)); |
| |
| let manifest_dir = format!("{table_path}/manifest"); |
| let delta_metas = ManifestList::read( |
| &file_io, |
| &format!("{manifest_dir}/{}", snapshot.delta_manifest_list()), |
| ) |
| .await |
| .unwrap(); |
| let delta_entries = Manifest::read( |
| &file_io, |
| &format!("{manifest_dir}/{}", delta_metas[0].file_name()), |
| ) |
| .await |
| .unwrap(); |
| assert_eq!(delta_entries.len(), 1); |
| assert_eq!(*delta_entries[0].kind(), FileKind::Add); |
| assert_eq!(delta_entries[0].file().file_name, data_file_name); |
| |
| let changelog_list = snapshot |
| .changelog_manifest_list() |
| .expect("changelog manifest list"); |
| let changelog_metas = |
| ManifestList::read(&file_io, &format!("{manifest_dir}/{changelog_list}")) |
| .await |
| .unwrap(); |
| let changelog_entries = Manifest::read( |
| &file_io, |
| &format!("{manifest_dir}/{}", changelog_metas[0].file_name()), |
| ) |
| .await |
| .unwrap(); |
| assert_eq!(changelog_entries.len(), 1); |
| assert_eq!(*changelog_entries[0].kind(), FileKind::Add); |
| assert_eq!(changelog_entries[0].file().file_name, changelog_file_name); |
| assert_eq!(changelog_entries[0].file().row_count, 2); |
| } |
| |
| #[tokio::test] |
| async fn test_input_changelog_dynamic_bucket_commits_data_changelog_and_index() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_input_changelog_dynamic_bucket_commit"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let table = Table::new( |
| file_io.clone(), |
| Identifier::new("default", "test_input_changelog_dynamic_bucket"), |
| table_path.to_string(), |
| ordinary_dynamic_pk_changelog_schema(), |
| None, |
| ); |
| |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| assert!(matches!( |
| table_write.bucket_assigner, |
| BucketAssignerEnum::Dynamic(_) |
| )); |
| table_write |
| .write_arrow_batch(&make_partitioned_batch_with_value_kind( |
| vec!["a", "a"], |
| vec![1, 2], |
| vec![10, 20], |
| vec![0, 3], |
| )) |
| .await |
| .unwrap(); |
| |
| let messages = table_write.prepare_commit().await.unwrap(); |
| assert_eq!(messages.len(), 1); |
| assert_eq!(messages[0].new_files.len(), 1); |
| assert_eq!(messages[0].new_files[0].row_count, 2); |
| assert_eq!(messages[0].new_files[0].delete_row_count, Some(1)); |
| assert_eq!(messages[0].new_changelog_files.len(), 1); |
| assert_eq!(messages[0].new_changelog_files[0].row_count, 2); |
| assert_eq!(messages[0].new_changelog_files[0].delete_row_count, Some(1)); |
| assert_eq!(messages[0].new_index_files.len(), 1); |
| assert_eq!(messages[0].new_index_files[0].index_type, "HASH"); |
| assert_eq!(messages[0].new_index_files[0].row_count, 2); |
| |
| let index_file_name = messages[0].new_index_files[0].file_name.clone(); |
| |
| let commit = TableCommit::new(table, "test-user".to_string()); |
| commit.commit(messages).await.unwrap(); |
| |
| let snap_manager = SnapshotManager::new(file_io.clone(), table_path.to_string()); |
| let snapshot = snap_manager.get_latest_snapshot().await.unwrap().unwrap(); |
| assert_eq!(snapshot.total_record_count(), Some(2)); |
| assert_eq!(snapshot.delta_record_count(), Some(2)); |
| assert_eq!(snapshot.changelog_record_count(), Some(2)); |
| |
| let manifest_dir = format!("{table_path}/manifest"); |
| let changelog_list = snapshot |
| .changelog_manifest_list() |
| .expect("changelog manifest list"); |
| let changelog_metas = |
| ManifestList::read(&file_io, &format!("{manifest_dir}/{changelog_list}")) |
| .await |
| .unwrap(); |
| let changelog_partition_stats = changelog_metas[0].partition_stats(); |
| assert!(!changelog_partition_stats.min_values().is_empty()); |
| assert!(!changelog_partition_stats.max_values().is_empty()); |
| assert_eq!( |
| changelog_partition_stats.null_counts().as_slice(), |
| &[Some(0)] |
| ); |
| |
| let index_manifest = snapshot.index_manifest().expect("index manifest"); |
| let index_entries = |
| IndexManifest::read(&file_io, &format!("{manifest_dir}/{index_manifest}")) |
| .await |
| .unwrap(); |
| assert_eq!(index_entries.len(), 1); |
| assert_eq!(index_entries[0].kind, FileKind::Add); |
| assert_eq!(index_entries[0].index_file.file_name, index_file_name); |
| assert_eq!(index_entries[0].index_file.index_type, "HASH"); |
| assert_eq!(index_entries[0].index_file.row_count, 2); |
| } |
| |
| #[tokio::test] |
| async fn test_input_changelog_overwrite_does_not_write_changelog_files() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_input_changelog_overwrite"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let table = Table::new( |
| file_io, |
| Identifier::new("default", "test_input_changelog"), |
| table_path.to_string(), |
| pk_changelog_schema(&[("changelog-producer", "input")]), |
| None, |
| ); |
| |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()) |
| .unwrap() |
| .with_overwrite(); |
| table_write |
| .write_arrow_batch(&make_batch(vec![1], vec![10])) |
| .await |
| .unwrap(); |
| |
| let messages = table_write.prepare_commit().await.unwrap(); |
| assert_eq!(messages.len(), 1); |
| assert_eq!(messages[0].new_files.len(), 1); |
| assert!(messages[0].new_changelog_files.is_empty()); |
| } |
| |
| #[tokio::test] |
| async fn test_pk_write_and_commit() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_pk_write"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let table = test_pk_table(&file_io, table_path); |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| |
| let batch = make_batch(vec![3, 1, 2], vec![30, 10, 20]); |
| table_write.write_arrow_batch(&batch).await.unwrap(); |
| |
| let messages = table_write.prepare_commit().await.unwrap(); |
| assert_eq!(messages.len(), 1); |
| assert_eq!(messages[0].new_files.len(), 1); |
| |
| let file = &messages[0].new_files[0]; |
| assert_eq!(file.row_count, 3); |
| assert_eq!(file.level, 0); |
| assert_eq!(file.min_sequence_number, 0); |
| assert_eq!(file.max_sequence_number, 2); |
| // min_key and max_key should be non-empty (serialized BinaryRow) |
| assert!(!file.min_key.is_empty()); |
| assert!(!file.max_key.is_empty()); |
| |
| // Commit |
| let commit = TableCommit::new(table.clone(), "test-user".to_string()); |
| commit.commit(messages).await.unwrap(); |
| |
| let snap_manager = SnapshotManager::new(file_io.clone(), table_path.to_string()); |
| let snapshot = snap_manager.get_latest_snapshot().await.unwrap().unwrap(); |
| assert_eq!(snapshot.id(), 1); |
| assert_eq!(snapshot.total_record_count(), Some(3)); |
| } |
| |
| #[tokio::test] |
| async fn test_pk_write_sorted_output() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_pk_sorted"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let table = test_pk_table(&file_io, table_path); |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| |
| // Write unsorted data |
| let batch = make_batch(vec![5, 2, 4, 1, 3], vec![50, 20, 40, 10, 30]); |
| table_write.write_arrow_batch(&batch).await.unwrap(); |
| |
| let messages = table_write.prepare_commit().await.unwrap(); |
| let commit = TableCommit::new(table.clone(), "test-user".to_string()); |
| commit.commit(messages).await.unwrap(); |
| |
| // Read back using sort-merge reader — should be sorted by PK |
| let rb = table.new_read_builder(); |
| let scan = rb.new_scan(); |
| let plan = scan.plan().await.unwrap(); |
| let read = rb.new_read().unwrap(); |
| let batches: Vec<RecordBatch> = |
| futures::TryStreamExt::try_collect(read.to_arrow(plan.splits()).unwrap()) |
| .await |
| .unwrap(); |
| |
| let ids: Vec<i32> = batches |
| .iter() |
| .flat_map(|b| { |
| b.column(0) |
| .as_any() |
| .downcast_ref::<Int32Array>() |
| .unwrap() |
| .values() |
| .iter() |
| .copied() |
| }) |
| .collect(); |
| assert_eq!(ids, vec![1, 2, 3, 4, 5]); |
| |
| let values: Vec<i32> = batches |
| .iter() |
| .flat_map(|b| { |
| b.column(1) |
| .as_any() |
| .downcast_ref::<Int32Array>() |
| .unwrap() |
| .values() |
| .iter() |
| .copied() |
| }) |
| .collect(); |
| assert_eq!(values, vec![10, 20, 30, 40, 50]); |
| } |
| |
| #[tokio::test] |
| async fn test_pk_write_dedup_across_commits() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_pk_dedup"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let table = test_pk_table(&file_io, table_path); |
| |
| // First commit: id=1,2,3 |
| let mut tw1 = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| tw1.write_arrow_batch(&make_batch(vec![1, 2, 3], vec![10, 20, 30])) |
| .await |
| .unwrap(); |
| let msgs1 = tw1.prepare_commit().await.unwrap(); |
| let commit = TableCommit::new(table.clone(), "test-user".to_string()); |
| commit.commit(msgs1).await.unwrap(); |
| |
| // Second commit: id=2,3,4 with updated values, higher sequence numbers |
| let mut tw2 = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| tw2.write_arrow_batch(&make_batch(vec![2, 3, 4], vec![200, 300, 400])) |
| .await |
| .unwrap(); |
| let msgs2 = tw2.prepare_commit().await.unwrap(); |
| commit.commit(msgs2).await.unwrap(); |
| |
| // Read back — dedup should keep newer values for id=2,3 |
| let rb = table.new_read_builder(); |
| let scan = rb.new_scan(); |
| let plan = scan.plan().await.unwrap(); |
| let read = rb.new_read().unwrap(); |
| let batches: Vec<RecordBatch> = |
| futures::TryStreamExt::try_collect(read.to_arrow(plan.splits()).unwrap()) |
| .await |
| .unwrap(); |
| |
| let ids: Vec<i32> = batches |
| .iter() |
| .flat_map(|b| { |
| b.column(0) |
| .as_any() |
| .downcast_ref::<Int32Array>() |
| .unwrap() |
| .values() |
| .iter() |
| .copied() |
| }) |
| .collect(); |
| let values: Vec<i32> = batches |
| .iter() |
| .flat_map(|b| { |
| b.column(1) |
| .as_any() |
| .downcast_ref::<Int32Array>() |
| .unwrap() |
| .values() |
| .iter() |
| .copied() |
| }) |
| .collect(); |
| |
| assert_eq!(ids, vec![1, 2, 3, 4]); |
| assert_eq!(values, vec![10, 200, 300, 400]); |
| } |
| |
| #[tokio::test] |
| async fn test_pk_write_sequence_number_in_file() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_pk_seq"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let table = test_pk_table(&file_io, table_path); |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| |
| let batch = make_batch(vec![1, 2], vec![10, 20]); |
| table_write.write_arrow_batch(&batch).await.unwrap(); |
| |
| let messages = table_write.prepare_commit().await.unwrap(); |
| let file = &messages[0].new_files[0]; |
| // Fresh table, seq starts at 0, 2 rows → min=0, max=1 |
| assert_eq!(file.min_sequence_number, 0); |
| assert_eq!(file.max_sequence_number, 1); |
| } |
| |
| // ----------------------------------------------------------------------- |
| // Postpone bucket (bucket = -2) write tests |
| // ----------------------------------------------------------------------- |
| |
| fn test_postpone_pk_schema() -> TableSchema { |
| let schema = Schema::builder() |
| .column("id", DataType::Int(IntType::new())) |
| .column("value", DataType::Int(IntType::new())) |
| .primary_key(["id"]) |
| .option("bucket", "-2") |
| .build() |
| .unwrap(); |
| TableSchema::new(0, &schema) |
| } |
| |
| fn test_postpone_pk_table(file_io: &FileIO, table_path: &str) -> Table { |
| Table::new( |
| file_io.clone(), |
| Identifier::new("default", "test_postpone_table"), |
| table_path.to_string(), |
| test_postpone_pk_schema(), |
| None, |
| ) |
| } |
| |
| fn test_postpone_partitioned_schema() -> TableSchema { |
| let schema = Schema::builder() |
| .column("pt", DataType::VarChar(VarCharType::string_type())) |
| .column("id", DataType::Int(IntType::new())) |
| .column("value", DataType::Int(IntType::new())) |
| .primary_key(["pt", "id"]) |
| .partition_keys(["pt"]) |
| .option("bucket", "-2") |
| .build() |
| .unwrap(); |
| TableSchema::new(0, &schema) |
| } |
| |
| fn test_postpone_partitioned_table(file_io: &FileIO, table_path: &str) -> Table { |
| Table::new( |
| file_io.clone(), |
| Identifier::new("default", "test_postpone_table"), |
| table_path.to_string(), |
| test_postpone_partitioned_schema(), |
| None, |
| ) |
| } |
| |
| fn make_partitioned_batch_3col(pts: Vec<&str>, ids: Vec<i32>, values: Vec<i32>) -> RecordBatch { |
| let schema = Arc::new(ArrowSchema::new(vec![ |
| ArrowField::new("pt", ArrowDataType::Utf8, false), |
| ArrowField::new("id", ArrowDataType::Int32, false), |
| ArrowField::new("value", ArrowDataType::Int32, false), |
| ])); |
| RecordBatch::try_new( |
| schema, |
| vec![ |
| Arc::new(arrow_array::StringArray::from(pts)), |
| Arc::new(Int32Array::from(ids)), |
| Arc::new(Int32Array::from(values)), |
| ], |
| ) |
| .unwrap() |
| } |
| |
| #[tokio::test] |
| async fn test_postpone_write_and_commit() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_postpone_write"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let table = test_postpone_pk_table(&file_io, table_path); |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| |
| let batch = make_batch(vec![3, 1, 2], vec![30, 10, 20]); |
| table_write.write_arrow_batch(&batch).await.unwrap(); |
| |
| let messages = table_write.prepare_commit().await.unwrap(); |
| assert_eq!(messages.len(), 1); |
| assert_eq!(messages[0].bucket, POSTPONE_BUCKET); |
| assert_eq!(messages[0].new_files.len(), 1); |
| assert_eq!(messages[0].new_files[0].row_count, 3); |
| |
| // Commit and verify snapshot |
| let commit = TableCommit::new(table, "test-user".to_string()); |
| commit.commit(messages).await.unwrap(); |
| |
| let snap_manager = SnapshotManager::new(file_io.clone(), table_path.to_string()); |
| let snapshot = snap_manager.get_latest_snapshot().await.unwrap().unwrap(); |
| assert_eq!(snapshot.id(), 1); |
| assert_eq!(snapshot.total_record_count(), Some(3)); |
| } |
| |
| #[tokio::test] |
| async fn test_postpone_write_empty_batch() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_postpone_empty"; |
| let table = test_postpone_pk_table(&file_io, table_path); |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| |
| let batch = make_batch(vec![], vec![]); |
| table_write.write_arrow_batch(&batch).await.unwrap(); |
| |
| let messages = table_write.prepare_commit().await.unwrap(); |
| assert!(messages.is_empty()); |
| } |
| |
| #[tokio::test] |
| async fn test_postpone_write_multiple_batches() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_postpone_multi"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let table = test_postpone_pk_table(&file_io, table_path); |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| |
| table_write |
| .write_arrow_batch(&make_batch(vec![1, 2], vec![10, 20])) |
| .await |
| .unwrap(); |
| table_write |
| .write_arrow_batch(&make_batch(vec![3, 4], vec![30, 40])) |
| .await |
| .unwrap(); |
| |
| let messages = table_write.prepare_commit().await.unwrap(); |
| assert_eq!(messages.len(), 1); |
| assert_eq!(messages[0].bucket, POSTPONE_BUCKET); |
| |
| let total_rows: i64 = messages[0].new_files.iter().map(|f| f.row_count).sum(); |
| assert_eq!(total_rows, 4); |
| } |
| |
| #[tokio::test] |
| async fn test_postpone_write_partitioned() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_postpone_partitioned"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let table = test_postpone_partitioned_table(&file_io, table_path); |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| |
| let batch = |
| make_partitioned_batch_3col(vec!["a", "b", "a"], vec![1, 2, 3], vec![10, 20, 30]); |
| table_write.write_arrow_batch(&batch).await.unwrap(); |
| |
| let messages = table_write.prepare_commit().await.unwrap(); |
| // 2 partitions |
| assert_eq!(messages.len(), 2); |
| // All messages should use POSTPONE_BUCKET |
| for msg in &messages { |
| assert_eq!(msg.bucket, POSTPONE_BUCKET); |
| } |
| |
| let total_rows: i64 = messages |
| .iter() |
| .flat_map(|m| &m.new_files) |
| .map(|f| f.row_count) |
| .sum(); |
| assert_eq!(total_rows, 3); |
| |
| // Commit and verify |
| let commit = TableCommit::new(table, "test-user".to_string()); |
| commit.commit(messages).await.unwrap(); |
| |
| let snap_manager = SnapshotManager::new(file_io.clone(), table_path.to_string()); |
| let snapshot = snap_manager.get_latest_snapshot().await.unwrap().unwrap(); |
| assert_eq!(snapshot.id(), 1); |
| assert_eq!(snapshot.total_record_count(), Some(3)); |
| } |
| |
| #[tokio::test] |
| async fn test_postpone_write_reusable() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_postpone_reuse"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let table = test_postpone_pk_table(&file_io, table_path); |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| |
| // First write + prepare_commit |
| table_write |
| .write_arrow_batch(&make_batch(vec![1, 2], vec![10, 20])) |
| .await |
| .unwrap(); |
| let messages1 = table_write.prepare_commit().await.unwrap(); |
| assert_eq!(messages1.len(), 1); |
| assert_eq!(messages1[0].new_files[0].row_count, 2); |
| |
| // Second write + prepare_commit (reuse) |
| table_write |
| .write_arrow_batch(&make_batch(vec![3, 4, 5], vec![30, 40, 50])) |
| .await |
| .unwrap(); |
| let messages2 = table_write.prepare_commit().await.unwrap(); |
| assert_eq!(messages2.len(), 1); |
| assert_eq!(messages2[0].new_files[0].row_count, 3); |
| |
| // Empty prepare_commit |
| let messages3 = table_write.prepare_commit().await.unwrap(); |
| assert!(messages3.is_empty()); |
| } |
| |
| #[tokio::test] |
| async fn test_postpone_write_file_naming_and_kv_format() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_postpone_kv"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let table = test_postpone_pk_table(&file_io, table_path); |
| let mut table_write = TableWrite::new(&table, "my-commit-user".to_string()).unwrap(); |
| |
| let batch = make_batch(vec![3, 1, 2], vec![30, 10, 20]); |
| table_write.write_arrow_batch(&batch).await.unwrap(); |
| |
| let messages = table_write.prepare_commit().await.unwrap(); |
| let file = &messages[0].new_files[0]; |
| |
| // Verify postpone file naming: data-u-{commitUser}-s-{writeId}-w-{uuid}-{index}.parquet |
| assert!( |
| file.file_name.starts_with("data-u-my-commit-user-s-"), |
| "Expected postpone file prefix, got: {}", |
| file.file_name |
| ); |
| assert!( |
| file.file_name.contains("-w-"), |
| "Expected -w- in file name, got: {}", |
| file.file_name |
| ); |
| |
| // Verify KV format: read the parquet file and check physical columns |
| let bucket_dir = format!("{table_path}/bucket-postpone"); |
| let file_path = format!("{bucket_dir}/{}", file.file_name); |
| let input = file_io.new_input(&file_path).unwrap(); |
| let data = input.read().await.unwrap(); |
| let reader = |
| parquet::arrow::arrow_reader::ParquetRecordBatchReader::try_new(data, 1024).unwrap(); |
| let schema = reader.schema(); |
| // Physical schema: [_SEQUENCE_NUMBER, _VALUE_KIND, id, value] |
| assert_eq!(schema.fields().len(), 4); |
| assert_eq!(schema.field(0).name(), "_SEQUENCE_NUMBER"); |
| assert_eq!(schema.field(1).name(), "_VALUE_KIND"); |
| assert_eq!(schema.field(2).name(), "id"); |
| assert_eq!(schema.field(3).name(), "value"); |
| |
| // Data should be in arrival order (not sorted by PK): 3, 1, 2 |
| let batches: Vec<RecordBatch> = reader.into_iter().map(|r| r.unwrap()).collect(); |
| let ids: Vec<i32> = batches |
| .iter() |
| .flat_map(|b| { |
| b.column(2) |
| .as_any() |
| .downcast_ref::<Int32Array>() |
| .unwrap() |
| .values() |
| .iter() |
| .copied() |
| }) |
| .collect(); |
| assert_eq!( |
| ids, |
| vec![3, 1, 2], |
| "Postpone mode should preserve arrival order" |
| ); |
| |
| // Empty key stats for postpone mode |
| assert_eq!(file.min_key, EMPTY_SERIALIZED_ROW.clone()); |
| assert_eq!(file.max_key, EMPTY_SERIALIZED_ROW.clone()); |
| } |
| |
| // ----------------------------------------------------------------------- |
| // Cross-partition update tests (dynamic bucket, PK not including partition) |
| // ----------------------------------------------------------------------- |
| |
| fn test_cross_partition_schema() -> TableSchema { |
| // PK is only "id", partition is "pt" — PK does NOT include partition field |
| let schema = Schema::builder() |
| .column("pt", DataType::VarChar(VarCharType::string_type())) |
| .column("id", DataType::Int(IntType::new())) |
| .column("value", DataType::Int(IntType::new())) |
| .primary_key(["id"]) |
| .partition_keys(["pt"]) |
| .build() |
| .unwrap(); |
| TableSchema::new(0, &schema) |
| } |
| |
| fn test_cross_partition_table(file_io: &FileIO, table_path: &str) -> Table { |
| Table::new( |
| file_io.clone(), |
| Identifier::new("default", "test_cross_partition"), |
| table_path.to_string(), |
| test_cross_partition_schema(), |
| None, |
| ) |
| } |
| |
| #[tokio::test] |
| async fn test_cross_partition_detection() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_cross_detect"; |
| |
| // Cross-partition: PK=id, partition=pt, bucket=-1 |
| let table = test_cross_partition_table(&file_io, table_path); |
| let tw = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| assert!(matches!( |
| tw.bucket_assigner, |
| BucketAssignerEnum::CrossPartition(_) |
| )); |
| |
| // Non-cross-partition: PK includes partition field |
| let schema = Schema::builder() |
| .column("pt", DataType::VarChar(VarCharType::string_type())) |
| .column("id", DataType::Int(IntType::new())) |
| .column("value", DataType::Int(IntType::new())) |
| .primary_key(["pt", "id"]) |
| .partition_keys(["pt"]) |
| .build() |
| .unwrap(); |
| let table2 = Table::new( |
| file_io.clone(), |
| Identifier::new("default", "test"), |
| table_path.to_string(), |
| TableSchema::new(0, &schema), |
| None, |
| ); |
| let tw2 = TableWrite::new(&table2, "test-user".to_string()).unwrap(); |
| assert!(matches!( |
| tw2.bucket_assigner, |
| BucketAssignerEnum::Dynamic(_) |
| )); |
| } |
| |
| #[test] |
| fn test_rejects_cross_partition_partial_update() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_cross_partial_update"; |
| let schema = Schema::builder() |
| .column("pt", DataType::VarChar(VarCharType::string_type())) |
| .column("id", DataType::Int(IntType::new())) |
| .column("value", DataType::Int(IntType::new())) |
| .primary_key(["id"]) |
| .partition_keys(["pt"]) |
| .option("merge-engine", "partial-update") |
| .build() |
| .unwrap(); |
| let table = Table::new( |
| file_io, |
| Identifier::new("default", "test_cross_partial_update"), |
| table_path.to_string(), |
| TableSchema::new(0, &schema), |
| None, |
| ); |
| |
| let err = match TableWrite::new(&table, "test-user".to_string()) { |
| Ok(_) => panic!("cross-partition partial-update should be rejected"), |
| Err(err) => err, |
| }; |
| |
| assert!(matches!( |
| err, |
| crate::Error::Unsupported { message } |
| if message.contains("cross-partition update") |
| )); |
| } |
| |
| #[test] |
| fn test_rejects_cross_partition_aggregation() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_cross_aggregation"; |
| let schema = Schema::builder() |
| .column("pt", DataType::VarChar(VarCharType::string_type())) |
| .column("id", DataType::Int(IntType::new())) |
| .column("value", DataType::Int(IntType::new())) |
| .primary_key(["id"]) |
| .partition_keys(["pt"]) |
| .option("merge-engine", "aggregation") |
| .build() |
| .unwrap(); |
| let table = Table::new( |
| file_io, |
| Identifier::new("default", "test_cross_aggregation"), |
| table_path.to_string(), |
| TableSchema::new(0, &schema), |
| None, |
| ); |
| |
| let err = match TableWrite::new(&table, "test-user".to_string()) { |
| Ok(_) => panic!("cross-partition aggregation should be rejected"), |
| Err(err) => err, |
| }; |
| |
| assert!(matches!( |
| err, |
| crate::Error::Unsupported { message } |
| if message.contains("merge-engine=aggregation") |
| && message.contains("cross-partition update") |
| )); |
| } |
| |
| #[tokio::test] |
| async fn test_cross_partition_write_same_partition() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_cross_same_pt"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let table = test_cross_partition_table(&file_io, table_path); |
| let mut table_write = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| |
| // All rows in same partition — no cross-partition deletes |
| let batch = |
| make_partitioned_batch_3col(vec!["a", "a", "a"], vec![1, 2, 3], vec![10, 20, 30]); |
| table_write.write_arrow_batch(&batch).await.unwrap(); |
| |
| let messages = table_write.prepare_commit().await.unwrap(); |
| let total_data_files: usize = messages.iter().map(|m| m.new_files.len()).sum(); |
| assert!(total_data_files >= 1); |
| |
| let total_rows: i64 = messages |
| .iter() |
| .flat_map(|m| &m.new_files) |
| .map(|f| f.row_count) |
| .sum(); |
| assert_eq!(total_rows, 3); |
| } |
| |
| #[tokio::test] |
| async fn test_cross_partition_write_generates_delete() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_cross_delete"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let table = test_cross_partition_table(&file_io, table_path); |
| |
| // First commit: id=1 in partition "a" |
| let mut tw1 = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| let batch1 = make_partitioned_batch_3col(vec!["a"], vec![1], vec![10]); |
| tw1.write_arrow_batch(&batch1).await.unwrap(); |
| let msgs1 = tw1.prepare_commit().await.unwrap(); |
| let commit = TableCommit::new(table.clone(), "test-user".to_string()); |
| commit.commit(msgs1).await.unwrap(); |
| |
| // Second commit: id=1 moves to partition "b" |
| let mut tw2 = TableWrite::new(&table, "test-user".to_string()).unwrap(); |
| let batch2 = make_partitioned_batch_3col(vec!["b"], vec![1], vec![20]); |
| tw2.write_arrow_batch(&batch2).await.unwrap(); |
| let msgs2 = tw2.prepare_commit().await.unwrap(); |
| |
| // Should have messages for both partition "a" (delete) and partition "b" (add) |
| assert!( |
| msgs2.len() >= 2, |
| "Expected messages for both old and new partition, got {}", |
| msgs2.len() |
| ); |
| |
| // The total data files should include both the DELETE and ADD records |
| let total_rows: i64 = msgs2 |
| .iter() |
| .flat_map(|m| &m.new_files) |
| .map(|f| f.row_count) |
| .sum(); |
| // 1 ADD in partition "b" + 1 DELETE in partition "a" |
| assert_eq!(total_rows, 2); |
| |
| commit.commit(msgs2).await.unwrap(); |
| |
| // Read back — should only see id=1 with value=20 in partition "b" |
| let rb = table.new_read_builder(); |
| let scan = rb.new_scan(); |
| let plan = scan.plan().await.unwrap(); |
| let read = rb.new_read().unwrap(); |
| let batches: Vec<RecordBatch> = |
| futures::TryStreamExt::try_collect(read.to_arrow(plan.splits()).unwrap()) |
| .await |
| .unwrap(); |
| |
| let ids: Vec<i32> = batches |
| .iter() |
| .flat_map(|b| { |
| b.column(1) |
| .as_any() |
| .downcast_ref::<Int32Array>() |
| .unwrap() |
| .values() |
| .iter() |
| .copied() |
| }) |
| .collect(); |
| let values: Vec<i32> = batches |
| .iter() |
| .flat_map(|b| { |
| b.column(2) |
| .as_any() |
| .downcast_ref::<Int32Array>() |
| .unwrap() |
| .values() |
| .iter() |
| .copied() |
| }) |
| .collect(); |
| |
| assert_eq!(ids, vec![1]); |
| assert_eq!(values, vec![20]); |
| } |
| |
| fn tiny_split_pk_table(file_io: &FileIO, table_path: &str) -> Table { |
| Table::new( |
| file_io.clone(), |
| Identifier::new("default", "test_tiny_split_pk"), |
| table_path.to_string(), |
| pk_changelog_schema(&[ |
| ("source.split.target-size", "1b"), |
| ("source.split.open-file-cost", "1b"), |
| ]), |
| None, |
| ) |
| } |
| |
| async fn commit_one_batch(table: &Table, ids: Vec<i32>, values: Vec<i32>) { |
| let mut tw = TableWrite::new(table, "test-user".to_string()).unwrap(); |
| tw.write_arrow_batch(&make_batch(ids, values)) |
| .await |
| .unwrap(); |
| let msgs = tw.prepare_commit().await.unwrap(); |
| TableCommit::new(table.clone(), "test-user".to_string()) |
| .commit(msgs) |
| .await |
| .unwrap(); |
| } |
| |
| async fn read_id_value_rows(table: &Table) -> Vec<(i32, i32)> { |
| let rb = table.new_read_builder(); |
| let plan = rb.new_scan().plan().await.unwrap(); |
| let read = rb.new_read().unwrap(); |
| let batches: Vec<RecordBatch> = |
| futures::TryStreamExt::try_collect(read.to_arrow(plan.splits()).unwrap()) |
| .await |
| .unwrap(); |
| let mut rows: Vec<(i32, i32)> = batches |
| .iter() |
| .flat_map(|b| { |
| let ids = b.column(0).as_any().downcast_ref::<Int32Array>().unwrap(); |
| let values = b.column(1).as_any().downcast_ref::<Int32Array>().unwrap(); |
| (0..b.num_rows()).map(|i| (ids.value(i), values.value(i))) |
| }) |
| .collect(); |
| rows.sort_unstable(); |
| rows |
| } |
| |
| /// Regression test: three commits of the same primary key produce three |
| /// key-overlapping files. Even with a 1-byte split target they must stay |
| /// in one split so the sort-merge reader merges them, instead of each |
| /// split emitting its own (stale) version of the row. |
| #[tokio::test] |
| async fn test_pk_plan_keeps_overlapping_files_in_one_split_under_tiny_target() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_tiny_split_overlap"; |
| setup_dirs(&file_io, table_path).await; |
| let table = tiny_split_pk_table(&file_io, table_path); |
| |
| commit_one_batch(&table, vec![1], vec![10]).await; |
| commit_one_batch(&table, vec![1], vec![20]).await; |
| commit_one_batch(&table, vec![1], vec![30]).await; |
| |
| let plan = table.new_read_builder().new_scan().plan().await.unwrap(); |
| assert_eq!( |
| plan.splits().len(), |
| 1, |
| "overlapping files must share a split" |
| ); |
| assert_eq!(plan.splits()[0].data_files().len(), 3); |
| |
| assert_eq!(read_id_value_rows(&table).await, vec![(1, 30)]); |
| } |
| |
| /// Files with disjoint key ranges may still be distributed across splits |
| /// for parallelism when the split target is small. |
| #[tokio::test] |
| async fn test_pk_plan_separates_disjoint_files_under_tiny_target() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_tiny_split_disjoint"; |
| setup_dirs(&file_io, table_path).await; |
| let table = tiny_split_pk_table(&file_io, table_path); |
| |
| commit_one_batch(&table, vec![1], vec![10]).await; |
| commit_one_batch(&table, vec![2], vec![20]).await; |
| commit_one_batch(&table, vec![3], vec![30]).await; |
| |
| let plan = table.new_read_builder().new_scan().plan().await.unwrap(); |
| assert_eq!( |
| plan.splits().len(), |
| 3, |
| "disjoint files keep split parallelism" |
| ); |
| |
| assert_eq!( |
| read_id_value_rows(&table).await, |
| vec![(1, 10), (2, 20), (3, 30)] |
| ); |
| } |
| |
| /// Append-only tables have no primary keys (and empty min/max keys); they |
| /// must keep using plain file-level bin packing instead of degrading to a |
| /// single all-files section. |
| #[tokio::test] |
| async fn test_append_table_plan_uses_file_level_bin_pack_under_tiny_target() { |
| let file_io = test_file_io(); |
| let table_path = "memory:/test_tiny_split_append"; |
| setup_dirs(&file_io, table_path).await; |
| |
| let schema = Schema::builder() |
| .column("id", DataType::Int(IntType::new())) |
| .column("value", DataType::Int(IntType::new())) |
| .option("bucket", "1") |
| .option("bucket-key", "id") |
| .option("source.split.target-size", "1b") |
| .option("source.split.open-file-cost", "1b") |
| .build() |
| .unwrap(); |
| let table = Table::new( |
| file_io.clone(), |
| Identifier::new("default", "test_tiny_split_append"), |
| table_path.to_string(), |
| TableSchema::new(0, &schema), |
| None, |
| ); |
| |
| commit_one_batch(&table, vec![1], vec![10]).await; |
| commit_one_batch(&table, vec![2], vec![20]).await; |
| commit_one_batch(&table, vec![3], vec![30]).await; |
| |
| let plan = table.new_read_builder().new_scan().plan().await.unwrap(); |
| assert_eq!( |
| plan.splits().len(), |
| 3, |
| "append tables keep file-level bin pack" |
| ); |
| } |
| } |