From ef4f8482c481878c437db6eae5ab437e0ca61b2b Mon Sep 17 00:00:00 2001 From: kid Date: Fri, 17 Jul 2026 14:45:09 +0800 Subject: [PATCH] feat: add core COW rewrite primitive Add a core copy-on-write rewrite primitive that plans candidate data files, applies a caller-provided RecordBatch rewriter, writes replacement data files, and returns removed/added/unchanged file sets for future overwrite-style commit actions. - Plan candidates through the existing scan path so manifest pruning and delete-file application stay consistent with normal reads; the row predicate is cleared only for the full-file rewrite input. - Stream rewritten batches to the replacement writer: buffer only the prefix before the first changed batch, then open the writer lazily; unchanged files open no writer, fully-deleted files produce no replacement. - Bind the replacement writer and partition key to the planned snapshot's schema so rewrites survive schema evolution. - Build the replacement parquet writer via ParquetWriterBuilder::from_table_properties so replacement files honor the table's write.parquet.* properties. - Keep the planner internal; CowRewriteBuilder is the only public entry and CowRewriteFile fields sit behind accessors. Cover delete-style, update-style, no-op, full-file delete, no-match delete, delete-file planning, writer edge cases, schema evolution, and multi-source replacement path uniqueness. Part of #2269. Co-Authored-By: Claude --- crates/iceberg/public-api.txt | 52 + crates/iceberg/src/cow_rewrite/mod.rs | 1088 ++++++++++++++++++++ crates/iceberg/src/cow_rewrite/plan.rs | 153 +++ crates/iceberg/src/cow_rewrite/rewriter.rs | 42 + crates/iceberg/src/cow_rewrite/writer.rs | 334 ++++++ crates/iceberg/src/lib.rs | 1 + crates/iceberg/src/scan/context.rs | 12 + crates/iceberg/src/scan/mod.rs | 54 +- 8 files changed, 1722 insertions(+), 14 deletions(-) create mode 100644 crates/iceberg/src/cow_rewrite/mod.rs create mode 100644 crates/iceberg/src/cow_rewrite/plan.rs create mode 100644 crates/iceberg/src/cow_rewrite/rewriter.rs create mode 100644 crates/iceberg/src/cow_rewrite/writer.rs diff --git a/crates/iceberg/public-api.txt b/crates/iceberg/public-api.txt index 8a62277295..dced1f204b 100644 --- a/crates/iceberg/public-api.txt +++ b/crates/iceberg/public-api.txt @@ -170,6 +170,58 @@ impl serde_core::ser::Serialize for iceberg::compression::CompressionCodec pub fn iceberg::compression::CompressionCodec::serialize(&self, serializer: S) -> core::result::Result<::Ok, ::Error> impl<'de> serde_core::de::Deserialize<'de> for iceberg::compression::CompressionCodec pub fn iceberg::compression::CompressionCodec::deserialize>(deserializer: D) -> core::result::Result::Error> +pub mod iceberg::cow_rewrite +pub struct iceberg::cow_rewrite::CowBatchRewrite +pub iceberg::cow_rewrite::CowBatchRewrite::changed: bool +pub iceberg::cow_rewrite::CowBatchRewrite::output: core::option::Option +pub struct iceberg::cow_rewrite::CowRewriteBuilder<'a> +impl<'a> iceberg::cow_rewrite::CowRewriteBuilder<'a> +pub fn iceberg::cow_rewrite::CowRewriteBuilder<'a>::new(table: &'a iceberg::table::Table) -> Self +pub async fn iceberg::cow_rewrite::CowRewriteBuilder<'a>::rewrite(self) -> iceberg::Result +pub fn iceberg::cow_rewrite::CowRewriteBuilder<'a>::with_batch_size(self, batch_size: usize) -> Self +pub fn iceberg::cow_rewrite::CowRewriteBuilder<'a>::with_case_sensitive(self, case_sensitive: bool) -> Self +pub fn iceberg::cow_rewrite::CowRewriteBuilder<'a>::with_predicate(self, predicate: iceberg::expr::Predicate) -> Self +pub fn iceberg::cow_rewrite::CowRewriteBuilder<'a>::with_rewriter(self, rewriter: alloc::sync::Arc) -> Self +pub fn iceberg::cow_rewrite::CowRewriteBuilder<'a>::with_snapshot_id(self, snapshot_id: i64) -> Self +pub struct iceberg::cow_rewrite::CowRewriteFile +impl iceberg::cow_rewrite::CowRewriteFile +pub fn iceberg::cow_rewrite::CowRewriteFile::old_data_file(&self) -> &iceberg::spec::DataFile +pub fn iceberg::cow_rewrite::CowRewriteFile::scan_task(&self) -> &iceberg::scan::FileScanTask +impl core::clone::Clone for iceberg::cow_rewrite::CowRewriteFile +pub fn iceberg::cow_rewrite::CowRewriteFile::clone(&self) -> iceberg::cow_rewrite::CowRewriteFile +impl core::fmt::Debug for iceberg::cow_rewrite::CowRewriteFile +pub fn iceberg::cow_rewrite::CowRewriteFile::fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct iceberg::cow_rewrite::CowRewriteResult +pub iceberg::cow_rewrite::CowRewriteResult::added_data_files: alloc::vec::Vec +pub iceberg::cow_rewrite::CowRewriteResult::removed_data_files: alloc::vec::Vec +pub iceberg::cow_rewrite::CowRewriteResult::stats: iceberg::cow_rewrite::CowRewriteStats +pub iceberg::cow_rewrite::CowRewriteResult::unchanged_data_files: alloc::vec::Vec +impl iceberg::cow_rewrite::CowRewriteResult +pub fn iceberg::cow_rewrite::CowRewriteResult::has_changes(&self) -> bool +impl core::default::Default for iceberg::cow_rewrite::CowRewriteResult +pub fn iceberg::cow_rewrite::CowRewriteResult::default() -> iceberg::cow_rewrite::CowRewriteResult +impl core::fmt::Debug for iceberg::cow_rewrite::CowRewriteResult +pub fn iceberg::cow_rewrite::CowRewriteResult::fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct iceberg::cow_rewrite::CowRewriteStats +pub iceberg::cow_rewrite::CowRewriteStats::candidate_files: usize +pub iceberg::cow_rewrite::CowRewriteStats::changed_batches: u64 +pub iceberg::cow_rewrite::CowRewriteStats::input_rows: u64 +pub iceberg::cow_rewrite::CowRewriteStats::output_rows: u64 +pub iceberg::cow_rewrite::CowRewriteStats::rewritten_files: usize +pub iceberg::cow_rewrite::CowRewriteStats::unchanged_files: usize +impl core::clone::Clone for iceberg::cow_rewrite::CowRewriteStats +pub fn iceberg::cow_rewrite::CowRewriteStats::clone(&self) -> iceberg::cow_rewrite::CowRewriteStats +impl core::cmp::Eq for iceberg::cow_rewrite::CowRewriteStats +impl core::cmp::PartialEq for iceberg::cow_rewrite::CowRewriteStats +pub fn iceberg::cow_rewrite::CowRewriteStats::eq(&self, other: &iceberg::cow_rewrite::CowRewriteStats) -> bool +impl core::default::Default for iceberg::cow_rewrite::CowRewriteStats +pub fn iceberg::cow_rewrite::CowRewriteStats::default() -> iceberg::cow_rewrite::CowRewriteStats +impl core::fmt::Debug for iceberg::cow_rewrite::CowRewriteStats +pub fn iceberg::cow_rewrite::CowRewriteStats::fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result +impl core::marker::Copy for iceberg::cow_rewrite::CowRewriteStats +impl core::marker::StructuralPartialEq for iceberg::cow_rewrite::CowRewriteStats +pub trait iceberg::cow_rewrite::CowBatchRewriter: core::marker::Send + core::marker::Sync +pub fn iceberg::cow_rewrite::CowBatchRewriter::rewrite_batch(&self, batch: arrow_array::record_batch::RecordBatch) -> iceberg::Result pub mod iceberg::encryption pub mod iceberg::encryption::kms pub struct iceberg::encryption::kms::GeneratedKey diff --git a/crates/iceberg/src/cow_rewrite/mod.rs b/crates/iceberg/src/cow_rewrite/mod.rs new file mode 100644 index 0000000000..66210b510f --- /dev/null +++ b/crates/iceberg/src/cow_rewrite/mod.rs @@ -0,0 +1,1088 @@ +// 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. + +//! Copy-on-write rewrite primitives. +//! +//! This module plans candidate data files, reads their visible rows, applies a +//! caller-provided batch rewriter, and writes replacement data files. It returns +//! old and new file sets that can be committed by an overwrite-style transaction +//! action. +//! +//! The primitive does not parse SQL and does not commit metadata by itself. +//! Rewriters must emit batches compatible with the table schema and must +//! preserve each source file's partition values; this primitive does not +//! repartition rewritten rows. +//! +//! ```rust,no_run +//! # use std::sync::Arc; +//! # use arrow_array::RecordBatch; +//! # use iceberg::cow_rewrite::{CowBatchRewrite, CowBatchRewriter, CowRewriteBuilder}; +//! # use iceberg::table::Table; +//! # use iceberg::Result; +//! struct KeepAll; +//! +//! impl CowBatchRewriter for KeepAll { +//! fn rewrite_batch(&self, batch: RecordBatch) -> Result { +//! Ok(CowBatchRewrite { +//! output: Some(batch), +//! changed: false, +//! }) +//! } +//! } +//! +//! # async fn example(table: &Table) -> Result<()> { +//! let result = CowRewriteBuilder::new(table) +//! .with_rewriter(Arc::new(KeepAll)) +//! .rewrite() +//! .await?; +//! +//! assert!(!result.has_changes()); +//! # Ok(()) +//! # } +//! ``` + +mod plan; +mod rewriter; +pub(crate) mod writer; + +use std::sync::Arc; + +use arrow_array::RecordBatch; +use futures::TryStreamExt; +pub use plan::CowRewriteFile; +pub use rewriter::{CowBatchRewrite, CowBatchRewriter}; + +use crate::expr::Predicate; +use crate::scan::FileScanTaskStream; +use crate::spec::{DataFile, PartitionKey}; +use crate::table::Table; +use crate::{Error, ErrorKind, Result}; + +/// Counters produced by a copy-on-write rewrite. +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub struct CowRewriteStats { + /// Number of candidate files selected by planning. + pub candidate_files: usize, + /// Number of old files that have replacement output or are fully removed. + pub rewritten_files: usize, + /// Number of candidate files that did not change after row rewriting. + pub unchanged_files: usize, + /// Visible input row count read from candidate files. + pub input_rows: u64, + /// Output row count emitted by the batch rewriter. + pub output_rows: u64, + /// Number of input batches where the rewriter reported changes. + pub changed_batches: u64, +} + +/// Result of a copy-on-write rewrite operation. +#[derive(Debug, Default)] +pub struct CowRewriteResult { + /// Old data files that should be removed by the commit action. + pub removed_data_files: Vec, + /// New data files that should be added by the commit action. + pub added_data_files: Vec, + /// Candidate files that were read and left unchanged. + pub unchanged_data_files: Vec, + /// Rewrite counters. + pub stats: CowRewriteStats, +} + +impl CowRewriteResult { + /// Returns true if the rewrite produced any table changes. + pub fn has_changes(&self) -> bool { + !self.removed_data_files.is_empty() || !self.added_data_files.is_empty() + } +} + +/// Builder for orchestrating copy-on-write data file rewrites. +pub struct CowRewriteBuilder<'a> { + table: &'a Table, + predicate: Predicate, + snapshot_id: Option, + batch_size: Option, + case_sensitive: bool, + rewriter: Option>, +} + +impl<'a> CowRewriteBuilder<'a> { + /// Creates a copy-on-write rewrite builder for `table`. + pub fn new(table: &'a Table) -> Self { + Self { + table, + predicate: Predicate::AlwaysTrue, + snapshot_id: None, + batch_size: None, + case_sensitive: true, + rewriter: None, + } + } + + /// Sets the row predicate used to plan candidate files. + pub fn with_predicate(mut self, predicate: Predicate) -> Self { + self.predicate = predicate; + self + } + + /// Sets the snapshot id used to plan candidate files. + pub fn with_snapshot_id(mut self, snapshot_id: i64) -> Self { + self.snapshot_id = Some(snapshot_id); + self + } + + /// Sets the Arrow reader batch size. + pub fn with_batch_size(mut self, batch_size: usize) -> Self { + self.batch_size = Some(batch_size); + self + } + + /// Sets predicate binding case sensitivity for planning and reading. + pub fn with_case_sensitive(mut self, case_sensitive: bool) -> Self { + self.case_sensitive = case_sensitive; + self + } + + /// Sets the record batch rewriter. + pub fn with_rewriter(mut self, rewriter: Arc) -> Self { + self.rewriter = Some(rewriter); + self + } + + /// Plans, reads, rewrites, and writes replacement data files. + pub async fn rewrite(self) -> Result { + let rewriter = self.rewriter.ok_or_else(|| { + Error::new( + ErrorKind::PreconditionFailed, + "COW rewrite requires a batch rewriter", + ) + })?; + let files = plan::plan_cow_rewrite_files( + self.table, + Some(self.predicate), + self.snapshot_id, + self.case_sensitive, + ) + .await?; + + let mut result = CowRewriteResult { + stats: CowRewriteStats { + candidate_files: files.len(), + ..CowRewriteStats::default() + }, + ..CowRewriteResult::default() + }; + + for file in files { + // Schema the rows are read in (the planned snapshot's schema). The + // replacement files must be written with this schema so that batches + // remain compatible when the table's current schema has evolved past + // the snapshot the source files belong to. + let write_schema = file.scan_task.schema.clone(); + + // Batches produced before the first changed batch. They are buffered + // rather than written immediately because the primitive must not + // emit a replacement file for a source file that turns out to be + // unchanged. Once a changed batch is observed the buffered prefix is + // flushed to the writer and all subsequent batches stream straight + // through, so the in-memory footprint is bounded by the rows that + // precede the first change instead of the entire source file. + let mut prefix: Vec = Vec::new(); + let mut file_changed = false; + let mut writer: Option> = None; + + let mut task = file.scan_task.clone(); + task.predicate = None; + task.case_sensitive = self.case_sensitive; + let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream; + + let mut reader_builder = self.table.reader_builder(); + if let Some(batch_size) = self.batch_size { + reader_builder = reader_builder.with_batch_size(batch_size); + } + + let mut batches = reader_builder.build().read(tasks)?.stream(); + while let Some(batch) = batches.try_next().await? { + result.stats.input_rows += batch.num_rows() as u64; + + let rewrite = rewriter.rewrite_batch(batch)?; + if rewrite.changed { + file_changed = true; + result.stats.changed_batches += 1; + } + + if let Some(output) = rewrite.output { + result.stats.output_rows += output.num_rows() as u64; + + if file_changed { + if writer.is_none() { + let partition_key = source_partition_key( + self.table, + &file.old_data_file, + &write_schema, + )?; + writer = Some( + writer::build_replacement_writer( + self.table, + write_schema.clone(), + Some(partition_key), + ) + .await?, + ); + } + let writer = writer.as_mut().expect("writer just built"); + for prefix_batch in prefix.drain(..) { + writer.write(prefix_batch).await?; + } + writer.write(output).await?; + } else { + prefix.push(output); + } + } + } + + if file_changed { + result.stats.rewritten_files += 1; + result.removed_data_files.push(file.old_data_file.clone()); + + if let Some(mut writer) = writer { + let added_data_files = writer.close().await?; + result.added_data_files.extend(added_data_files); + } + // If `writer` is `None`, the source file was fully deleted + // (every batch dropped to `output: None`), so no replacement + // file is written. + } else { + result.stats.unchanged_files += 1; + result.unchanged_data_files.push(file.old_data_file); + // `prefix` is dropped here; no replacement file was written. + } + } + + Ok(result) + } +} + +fn source_partition_key( + table: &Table, + data_file: &DataFile, + schema: &crate::spec::SchemaRef, +) -> Result { + let spec = table + .metadata() + .partition_spec_by_id(data_file.partition_spec_id) + .ok_or_else(|| { + Error::new( + ErrorKind::DataInvalid, + format!( + "Missing partition spec {} for COW rewrite source file", + data_file.partition_spec_id + ), + ) + })? + .as_ref() + .clone(); + spec.partition_type(schema).map_err(|err| { + Error::new( + ErrorKind::DataInvalid, + format!( + "Cannot bind partition spec {} to the planned snapshot schema for COW rewrite", + data_file.partition_spec_id + ), + ) + .with_source(err) + })?; + + Ok(PartitionKey::new( + spec, + schema.clone(), + data_file.partition().clone(), + )) +} + +#[cfg(test)] +mod tests { + use std::collections::{HashMap, HashSet}; + use std::sync::Arc; + + use arrow_array::{Array, ArrayRef, BooleanArray, Int32Array, RecordBatch}; + use arrow_schema::{DataType, Field, Schema as ArrowSchema}; + use futures::TryStreamExt; + use parquet::arrow::PARQUET_FIELD_ID_META_KEY; + use tempfile::TempDir; + + use crate::cow_rewrite::{CowBatchRewrite, CowBatchRewriter, CowRewriteBuilder}; + use crate::io::LocalFsStorageFactory; + use crate::memory::{MEMORY_CATALOG_WAREHOUSE, MemoryCatalogBuilder}; + use crate::scan::{FileScanTask, FileScanTaskStream}; + use crate::spec::{DataFile, NestedField, PrimitiveType, Schema, TableProperties, Type}; + use crate::table::Table; + use crate::transaction::{AddColumn, ApplyTransactionAction, Transaction}; + use crate::{Catalog, CatalogBuilder, Error, ErrorKind, NamespaceIdent, Result, TableCreation}; + + struct KeepAll; + + impl CowBatchRewriter for KeepAll { + fn rewrite_batch(&self, batch: RecordBatch) -> Result { + Ok(CowBatchRewrite { + output: Some(batch), + changed: false, + }) + } + } + + struct DeleteEvenIds; + + impl CowBatchRewriter for DeleteEvenIds { + fn rewrite_batch(&self, batch: RecordBatch) -> Result { + let ids = batch + .column_by_name("id") + .ok_or_else(|| Error::new(ErrorKind::DataInvalid, "missing id column"))? + .as_any() + .downcast_ref::() + .ok_or_else(|| Error::new(ErrorKind::DataInvalid, "id must be Int32"))?; + + let keep = + BooleanArray::from_iter((0..ids.len()).map(|row| Some(ids.value(row) % 2 != 0))); + let filtered = arrow_select::filter::filter_record_batch(&batch, &keep) + .map_err(|err| Error::new(ErrorKind::Unexpected, err.to_string()))?; + + Ok(CowBatchRewrite { + changed: filtered.num_rows() != batch.num_rows(), + output: (filtered.num_rows() > 0).then_some(filtered), + }) + } + } + + struct IncrementValueForEvenIds; + + impl CowBatchRewriter for IncrementValueForEvenIds { + fn rewrite_batch(&self, batch: RecordBatch) -> Result { + let ids = batch + .column_by_name("id") + .ok_or_else(|| Error::new(ErrorKind::DataInvalid, "missing id column"))? + .as_any() + .downcast_ref::() + .ok_or_else(|| Error::new(ErrorKind::DataInvalid, "id must be Int32"))?; + let values = batch + .column_by_name("value") + .ok_or_else(|| Error::new(ErrorKind::DataInvalid, "missing value column"))? + .as_any() + .downcast_ref::() + .ok_or_else(|| Error::new(ErrorKind::DataInvalid, "value must be Int32"))?; + + let mut changed = false; + let updated_values = Int32Array::from_iter((0..values.len()).map(|row| { + let value = values.value(row); + if ids.value(row) % 2 == 0 { + changed = true; + Some(value + 10) + } else { + Some(value) + } + })); + let output = RecordBatch::try_new(batch.schema(), vec![ + batch.column(0).clone(), + Arc::new(updated_values), + ]) + .map_err(|err| Error::new(ErrorKind::Unexpected, err.to_string()))?; + + Ok(CowBatchRewrite { + output: Some(output), + changed, + }) + } + } + + struct CowRewriteFixture { + _temp_dir: TempDir, + table: Table, + } + + async fn test_table_with_ids(ids: Vec) -> Result { + let temp_dir = TempDir::new().unwrap(); + let warehouse = format!("file://{}", temp_dir.path().join("warehouse").display()); + let catalog = MemoryCatalogBuilder::default() + .with_storage_factory(Arc::new(LocalFsStorageFactory)) + .load( + "memory", + HashMap::from([(MEMORY_CATALOG_WAREHOUSE.to_string(), warehouse)]), + ) + .await?; + let namespace = NamespaceIdent::new("ns".to_string()); + catalog.create_namespace(&namespace, HashMap::new()).await?; + + let schema = Schema::builder() + .with_fields(vec![ + NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(), + ]) + .build()?; + let table = catalog + .create_table( + &namespace, + TableCreation::builder() + .name("cow_rewrite_fixture".to_string()) + .schema(schema) + .build(), + ) + .await?; + + let arrow_schema = Arc::new(ArrowSchema::new(vec![ + Field::new("id", DataType::Int32, false).with_metadata(HashMap::from([( + PARQUET_FIELD_ID_META_KEY.to_string(), + "1".to_string(), + )])), + ])); + let batch = RecordBatch::try_new(arrow_schema, vec![ + Arc::new(Int32Array::from(ids)) as ArrayRef + ])?; + let data_files = super::writer::write_replacement_batches( + &table, + table.metadata().current_schema().clone(), + None, + futures::stream::iter(vec![Ok(batch)]), + ) + .await?; + + let tx = Transaction::new(&table); + let tx = tx.fast_append().add_data_files(data_files).apply(tx)?; + let table = tx.commit(&catalog).await?; + + Ok(CowRewriteFixture { + _temp_dir: temp_dir, + table, + }) + } + + async fn test_table_with_id_batches(batches: Vec>) -> Result { + let temp_dir = TempDir::new().unwrap(); + let warehouse = format!("file://{}", temp_dir.path().join("warehouse").display()); + let catalog = MemoryCatalogBuilder::default() + .with_storage_factory(Arc::new(LocalFsStorageFactory)) + .load( + "memory", + HashMap::from([(MEMORY_CATALOG_WAREHOUSE.to_string(), warehouse)]), + ) + .await?; + let namespace = NamespaceIdent::new("ns".to_string()); + catalog.create_namespace(&namespace, HashMap::new()).await?; + + let schema = Schema::builder() + .with_fields(vec![ + NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(), + ]) + .build()?; + let table = catalog + .create_table( + &namespace, + TableCreation::builder() + .name("cow_rewrite_fixture".to_string()) + .schema(schema) + .properties(HashMap::from([( + TableProperties::PROPERTY_WRITE_TARGET_FILE_SIZE_BYTES.to_string(), + "1".to_string(), + )])) + .build(), + ) + .await?; + + let arrow_schema = Arc::new(ArrowSchema::new(vec![ + Field::new("id", DataType::Int32, false).with_metadata(HashMap::from([( + PARQUET_FIELD_ID_META_KEY.to_string(), + "1".to_string(), + )])), + ])); + let input = batches.into_iter().map(|ids| { + Ok(RecordBatch::try_new(arrow_schema.clone(), vec![ + Arc::new(Int32Array::from(ids)) as ArrayRef, + ])?) + }); + let data_files = super::writer::write_replacement_batches( + &table, + table.metadata().current_schema().clone(), + None, + futures::stream::iter(input), + ) + .await?; + + let tx = Transaction::new(&table); + let tx = tx.fast_append().add_data_files(data_files).apply(tx)?; + let table = tx.commit(&catalog).await?; + + Ok(CowRewriteFixture { + _temp_dir: temp_dir, + table, + }) + } + + async fn test_table_with_id_value_rows(rows: Vec<(i32, i32)>) -> Result { + let temp_dir = TempDir::new().unwrap(); + let warehouse = format!("file://{}", temp_dir.path().join("warehouse").display()); + let catalog = MemoryCatalogBuilder::default() + .with_storage_factory(Arc::new(LocalFsStorageFactory)) + .load( + "memory", + HashMap::from([(MEMORY_CATALOG_WAREHOUSE.to_string(), warehouse)]), + ) + .await?; + let namespace = NamespaceIdent::new("ns".to_string()); + catalog.create_namespace(&namespace, HashMap::new()).await?; + + let schema = Schema::builder() + .with_fields(vec![ + NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(), + NestedField::required(2, "value", Type::Primitive(PrimitiveType::Int)).into(), + ]) + .build()?; + let table = catalog + .create_table( + &namespace, + TableCreation::builder() + .name("cow_rewrite_fixture".to_string()) + .schema(schema) + .build(), + ) + .await?; + + let arrow_schema = Arc::new(ArrowSchema::new(vec![ + Field::new("id", DataType::Int32, false).with_metadata(HashMap::from([( + PARQUET_FIELD_ID_META_KEY.to_string(), + "1".to_string(), + )])), + Field::new("value", DataType::Int32, false).with_metadata(HashMap::from([( + PARQUET_FIELD_ID_META_KEY.to_string(), + "2".to_string(), + )])), + ])); + let ids = rows.iter().map(|(id, _)| *id).collect::>(); + let values = rows.iter().map(|(_, value)| *value).collect::>(); + let batch = RecordBatch::try_new(arrow_schema, vec![ + Arc::new(Int32Array::from(ids)) as ArrayRef, + Arc::new(Int32Array::from(values)) as ArrayRef, + ])?; + let data_files = super::writer::write_replacement_batches( + &table, + table.metadata().current_schema().clone(), + None, + futures::stream::iter(vec![Ok(batch)]), + ) + .await?; + + let tx = Transaction::new(&table); + let tx = tx.fast_append().add_data_files(data_files).apply(tx)?; + let table = tx.commit(&catalog).await?; + + Ok(CowRewriteFixture { + _temp_dir: temp_dir, + table, + }) + } + + async fn read_ids(table: &Table, files: &[DataFile]) -> Result> { + let schema = table.metadata().current_schema().clone(); + let project_field_ids = schema + .as_struct() + .fields() + .iter() + .map(|field| field.id) + .collect::>(); + let tasks = files + .iter() + .map(|data_file| { + Ok(FileScanTask::builder() + .with_file_size_in_bytes(data_file.file_size_in_bytes()) + .with_start(0) + .with_length(data_file.file_size_in_bytes()) + .with_record_count(Some(data_file.record_count())) + .with_data_file_path(data_file.file_path().to_string()) + .with_data_file_format(data_file.file_format()) + .with_schema(schema.clone()) + .with_project_field_ids(project_field_ids.clone()) + .with_case_sensitive(true) + .build()) + }) + .collect::>(); + let task_stream = Box::pin(futures::stream::iter(tasks)) as FileScanTaskStream; + let batches = table + .reader_builder() + .build() + .read(task_stream)? + .stream() + .try_collect::>() + .await?; + + let mut ids = Vec::new(); + for batch in batches { + let column = batch + .column_by_name("id") + .ok_or_else(|| Error::new(ErrorKind::DataInvalid, "missing id column"))? + .as_any() + .downcast_ref::() + .ok_or_else(|| Error::new(ErrorKind::DataInvalid, "id must be Int32"))?; + ids.extend((0..column.len()).map(|row| column.value(row))); + } + ids.sort_unstable(); + Ok(ids) + } + + async fn read_ids_and_values(table: &Table, files: &[DataFile]) -> Result> { + let schema = table.metadata().current_schema().clone(); + let project_field_ids = schema + .as_struct() + .fields() + .iter() + .map(|field| field.id) + .collect::>(); + let tasks = files + .iter() + .map(|data_file| { + Ok(FileScanTask::builder() + .with_file_size_in_bytes(data_file.file_size_in_bytes()) + .with_start(0) + .with_length(data_file.file_size_in_bytes()) + .with_record_count(Some(data_file.record_count())) + .with_data_file_path(data_file.file_path().to_string()) + .with_data_file_format(data_file.file_format()) + .with_schema(schema.clone()) + .with_project_field_ids(project_field_ids.clone()) + .with_case_sensitive(true) + .build()) + }) + .collect::>(); + let task_stream = Box::pin(futures::stream::iter(tasks)) as FileScanTaskStream; + let batches = table + .reader_builder() + .build() + .read(task_stream)? + .stream() + .try_collect::>() + .await?; + + let mut rows = Vec::new(); + for batch in batches { + let ids = batch + .column_by_name("id") + .ok_or_else(|| Error::new(ErrorKind::DataInvalid, "missing id column"))? + .as_any() + .downcast_ref::() + .ok_or_else(|| Error::new(ErrorKind::DataInvalid, "id must be Int32"))?; + let values = batch + .column_by_name("value") + .ok_or_else(|| Error::new(ErrorKind::DataInvalid, "missing value column"))? + .as_any() + .downcast_ref::() + .ok_or_else(|| Error::new(ErrorKind::DataInvalid, "value must be Int32"))?; + rows.extend((0..ids.len()).map(|row| (ids.value(row), values.value(row)))); + } + rows.sort_unstable_by_key(|(id, _)| *id); + Ok(rows) + } + + #[tokio::test] + async fn cow_rewrite_keep_all_produces_no_changes() -> Result<()> { + let fixture = test_table_with_ids(vec![1, 2, 3]).await?; + + let result = CowRewriteBuilder::new(&fixture.table) + .with_predicate(crate::expr::Predicate::AlwaysTrue) + .with_rewriter(Arc::new(KeepAll)) + .rewrite() + .await?; + + assert!(!result.has_changes()); + assert_eq!(result.removed_data_files.len(), 0); + assert_eq!(result.added_data_files.len(), 0); + assert_eq!(result.unchanged_data_files.len(), 1); + assert_eq!(result.stats.candidate_files, 1); + assert_eq!(result.stats.unchanged_files, 1); + assert_eq!(result.stats.input_rows, 3); + assert_eq!(result.stats.output_rows, 3); + + Ok(()) + } + + #[tokio::test] + async fn cow_rewrite_requires_rewriter() -> Result<()> { + let fixture = test_table_with_ids(vec![1]).await?; + + let err = CowRewriteBuilder::new(&fixture.table) + .rewrite() + .await + .expect_err("missing rewriter should fail"); + + assert_eq!(err.kind(), ErrorKind::PreconditionFailed); + + Ok(()) + } + + #[tokio::test] + async fn cow_rewrite_delete_rows_removes_old_file_and_adds_replacement() -> Result<()> { + let fixture = test_table_with_ids(vec![1, 2, 3, 4]).await?; + + let result = CowRewriteBuilder::new(&fixture.table) + .with_predicate(crate::expr::Predicate::AlwaysTrue) + .with_rewriter(Arc::new(DeleteEvenIds)) + .rewrite() + .await?; + + assert!(result.has_changes()); + assert_eq!(result.removed_data_files.len(), 1); + assert_eq!(result.added_data_files.len(), 1); + assert_eq!(result.stats.input_rows, 4); + assert_eq!(result.stats.output_rows, 2); + + let ids = read_ids(&fixture.table, &result.added_data_files).await?; + assert_eq!(ids, vec![1, 3]); + + Ok(()) + } + + #[tokio::test] + async fn cow_rewrite_update_rows_rewrites_file_with_updated_values() -> Result<()> { + let fixture = + test_table_with_id_value_rows(vec![(1, 10), (2, 20), (3, 30), (4, 40)]).await?; + + let result = CowRewriteBuilder::new(&fixture.table) + .with_predicate(crate::expr::Predicate::AlwaysTrue) + .with_rewriter(Arc::new(IncrementValueForEvenIds)) + .rewrite() + .await?; + + assert_eq!(result.removed_data_files.len(), 1); + assert_eq!(result.added_data_files.len(), 1); + assert_eq!(result.stats.input_rows, 4); + assert_eq!(result.stats.output_rows, 4); + + let rows = read_ids_and_values(&fixture.table, &result.added_data_files).await?; + assert_eq!(rows, vec![(1, 10), (2, 30), (3, 30), (4, 50)]); + + Ok(()) + } + + #[tokio::test] + async fn cow_rewrite_full_file_delete_removes_old_file_without_replacement() -> Result<()> { + let fixture = test_table_with_ids(vec![2, 4]).await?; + + let result = CowRewriteBuilder::new(&fixture.table) + .with_predicate(crate::expr::Predicate::AlwaysTrue) + .with_rewriter(Arc::new(DeleteEvenIds)) + .rewrite() + .await?; + + assert!(result.has_changes()); + assert_eq!(result.removed_data_files.len(), 1); + assert_eq!(result.added_data_files.len(), 0); + assert_eq!(result.stats.input_rows, 2); + assert_eq!(result.stats.output_rows, 0); + + Ok(()) + } + + #[tokio::test] + async fn cow_rewrite_delete_no_matching_rows_keeps_old_file() -> Result<()> { + let fixture = test_table_with_ids(vec![1, 3]).await?; + + let result = CowRewriteBuilder::new(&fixture.table) + .with_predicate(crate::expr::Predicate::AlwaysTrue) + .with_rewriter(Arc::new(DeleteEvenIds)) + .rewrite() + .await?; + + assert!(!result.has_changes()); + assert_eq!(result.removed_data_files.len(), 0); + assert_eq!(result.added_data_files.len(), 0); + assert_eq!(result.unchanged_data_files.len(), 1); + assert_eq!(result.stats.input_rows, 2); + assert_eq!(result.stats.output_rows, 2); + + Ok(()) + } + + #[tokio::test] + async fn cow_rewrite_uses_unique_replacement_paths_for_multiple_source_files() -> Result<()> { + let fixture = test_table_with_id_batches(vec![vec![1, 2], vec![3, 4]]).await?; + + let result = CowRewriteBuilder::new(&fixture.table) + .with_predicate(crate::expr::Predicate::AlwaysTrue) + .with_rewriter(Arc::new(DeleteEvenIds)) + .rewrite() + .await?; + + assert_eq!(result.stats.candidate_files, 2); + assert_eq!(result.removed_data_files.len(), 2); + assert_eq!(result.added_data_files.len(), 2); + + let added_paths = result + .added_data_files + .iter() + .map(|file| file.file_path().to_string()) + .collect::>(); + let removed_paths = result + .removed_data_files + .iter() + .map(|file| file.file_path().to_string()) + .collect::>(); + + assert_eq!(added_paths.len(), result.added_data_files.len()); + assert!(added_paths.is_disjoint(&removed_paths)); + + let ids = read_ids(&fixture.table, &result.added_data_files).await?; + assert_eq!(ids, vec![1, 3]); + + Ok(()) + } + + #[test] + fn cow_batch_rewriter_is_object_safe() { + let _rewriter: Arc = Arc::new(KeepAll); + } + + #[test] + fn cow_rewrite_result_reports_no_changes() { + let result = crate::cow_rewrite::CowRewriteResult { + removed_data_files: vec![], + added_data_files: vec![], + unchanged_data_files: vec![], + stats: crate::cow_rewrite::CowRewriteStats { + candidate_files: 0, + rewritten_files: 0, + unchanged_files: 0, + input_rows: 0, + output_rows: 0, + changed_batches: 0, + }, + }; + + assert!(!result.has_changes()); + assert_eq!(result.stats.candidate_files, 0); + } + + // --- Reproducer test added during code-review verification. --- + + /// `DeleteIfModThree` keeps rows whose `id` is not divisible by 3. + struct DeleteIfModThree; + + impl CowBatchRewriter for DeleteIfModThree { + fn rewrite_batch(&self, batch: RecordBatch) -> Result { + let ids = batch + .column_by_name("id") + .ok_or_else(|| Error::new(ErrorKind::DataInvalid, "missing id column"))? + .as_any() + .downcast_ref::() + .ok_or_else(|| Error::new(ErrorKind::DataInvalid, "id must be Int32"))?; + + let keep = + BooleanArray::from_iter((0..ids.len()).map(|row| Some(ids.value(row) % 3 != 0))); + let filtered = arrow_select::filter::filter_record_batch(&batch, &keep) + .map_err(|err| Error::new(ErrorKind::Unexpected, err.to_string()))?; + + Ok(CowBatchRewrite { + changed: filtered.num_rows() != batch.num_rows(), + output: (filtered.num_rows() > 0).then_some(filtered), + }) + } + } + + /// Issue #2 / #3: after adding an optional column (`value`) to the schema, + /// a COW rewrite of the pre-existing data must not fail. The replacement + /// writer must use the schema the batches were read in (the snapshot schema), + /// not the table's evolved current schema, otherwise the parquet writer + /// rejects the input batch whose column set does not include `value`. + #[tokio::test] + async fn cow_rewrite_after_optional_column_add() -> Result<()> { + let temp_dir = TempDir::new().unwrap(); + let warehouse = format!("file://{}", temp_dir.path().join("warehouse").display()); + let catalog: Arc = Arc::new( + MemoryCatalogBuilder::default() + .with_storage_factory(Arc::new(LocalFsStorageFactory)) + .load( + "memory", + HashMap::from([(MEMORY_CATALOG_WAREHOUSE.to_string(), warehouse)]), + ) + .await?, + ); + let namespace = NamespaceIdent::new("ns".to_string()); + catalog.create_namespace(&namespace, HashMap::new()).await?; + + let schema = Schema::builder() + .with_fields(vec![ + NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(), + ]) + .build()?; + let table = catalog + .create_table( + &namespace, + TableCreation::builder() + .name("evolved".to_string()) + .schema(schema) + .build(), + ) + .await?; + + // Write the original data file with the {id}-only schema. + let arrow_schema = Arc::new(ArrowSchema::new(vec![ + Field::new("id", DataType::Int32, false).with_metadata(HashMap::from([( + PARQUET_FIELD_ID_META_KEY.to_string(), + "1".to_string(), + )])), + ])); + let batch = RecordBatch::try_new(arrow_schema, vec![Arc::new(Int32Array::from(vec![ + 1, 2, 3, 4, + ])) as ArrayRef])?; + let data_files = super::writer::write_replacement_batches( + &table, + table.metadata().current_schema().clone(), + None, + futures::stream::iter(vec![Ok(batch)]), + ) + .await?; + let tx = Transaction::new(&table); + let tx = tx.fast_append().add_data_files(data_files).apply(tx)?; + let table = tx.commit(&*catalog).await?; + + // Evolve the schema: add an optional `value` column. + let tx = Transaction::new(&table); + let tx = tx + .update_schema() + .add_column(AddColumn::optional( + "value", + Type::Primitive(PrimitiveType::Int), + )) + .apply(tx)?; + let table = tx.commit(&*catalog).await?; + + // The current schema now has {id, value} but the current snapshot's + // schema is still the original {id} schema. + let current_snapshot = table.metadata().current_snapshot().unwrap(); + assert_ne!( + table.metadata().current_schema_id(), + current_snapshot.schema_id().unwrap() + ); + + // Before the fix this fails when the parquet writer rejects the {id}-only + // batch while configured with the {id, value} current schema. + let result = CowRewriteBuilder::new(&table) + .with_predicate(crate::expr::Predicate::AlwaysTrue) + .with_rewriter(Arc::new(DeleteIfModThree)) + .rewrite() + .await?; + + // id 3 is divisible by 3 and is dropped; remaining ids are 1, 2, 4. + assert!(result.has_changes()); + assert_eq!(result.removed_data_files.len(), 1); + assert_eq!(result.added_data_files.len(), 1); + + let ids = read_ids(&table, &result.added_data_files).await?; + assert_eq!(ids, vec![1, 2, 4]); + + Ok(()) + } + + /// Builds a two-file table whose files have disjoint id ranges so predicate + /// metrics pruning can exclude one file entirely. + async fn two_file_table_disjoint_ids() -> Result { + let temp_dir = TempDir::new().unwrap(); + let warehouse = format!("file://{}", temp_dir.path().join("warehouse").display()); + let catalog = MemoryCatalogBuilder::default() + .with_storage_factory(Arc::new(LocalFsStorageFactory)) + .load( + "memory", + HashMap::from([(MEMORY_CATALOG_WAREHOUSE.to_string(), warehouse)]), + ) + .await?; + let namespace = NamespaceIdent::new("ns".to_string()); + catalog.create_namespace(&namespace, HashMap::new()).await?; + + let schema = Schema::builder() + .with_fields(vec![ + NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(), + ]) + .build()?; + let table = catalog + .create_table( + &namespace, + TableCreation::builder() + .name("disjoint".to_string()) + .schema(schema) + .properties(HashMap::from([( + TableProperties::PROPERTY_WRITE_TARGET_FILE_SIZE_BYTES.to_string(), + "1".to_string(), + )])) + .build(), + ) + .await?; + + let arrow_schema = Arc::new(ArrowSchema::new(vec![ + Field::new("id", DataType::Int32, false).with_metadata(HashMap::from([( + PARQUET_FIELD_ID_META_KEY.to_string(), + "1".to_string(), + )])), + ])); + // First batch ids 1..=4, second batch ids 100..=103. Force a split via + // tiny target file size so each batch lands in its own data file. + let input = vec![vec![1, 2, 3, 4], vec![100, 101, 102, 103]] + .into_iter() + .map(|ids| { + Ok(RecordBatch::try_new(arrow_schema.clone(), vec![ + Arc::new(Int32Array::from(ids)) as ArrayRef, + ])?) + }); + let data_files = super::writer::write_replacement_batches( + &table, + table.metadata().current_schema().clone(), + None, + futures::stream::iter(input), + ) + .await?; + + let tx = Transaction::new(&table); + let tx = tx.fast_append().add_data_files(data_files).apply(tx)?; + let table = tx.commit(&catalog).await?; + + Ok(CowRewriteFixture { + _temp_dir: temp_dir, + table, + }) + } + + /// Predicate-based planning must exclude files whose metrics cannot match + /// the predicate. With `id > 50`, the file holding ids 1..=4 must never be + /// planned as a candidate. + #[tokio::test] + async fn cow_rewrite_predicate_prunes_non_matching_files() -> Result<()> { + let fixture = two_file_table_disjoint_ids().await?; + + let predicate = crate::expr::Reference::new("id").greater_than(crate::spec::Datum::int(50)); + + let result = CowRewriteBuilder::new(&fixture.table) + .with_predicate(predicate) + .with_rewriter(Arc::new(KeepAll)) + .rewrite() + .await?; + + // Only the high-id file should be a candidate. + assert_eq!(result.stats.candidate_files, 1); + // KeepAll does not change anything, so the single candidate is unchanged. + assert_eq!(result.removed_data_files.len(), 0); + assert_eq!(result.added_data_files.len(), 0); + assert_eq!(result.unchanged_data_files.len(), 1); + + // The surviving candidate must be the high-id file (ids 100..=103). + let ids = read_ids(&fixture.table, &result.unchanged_data_files).await?; + assert_eq!(ids, vec![100, 101, 102, 103]); + + Ok(()) + } +} diff --git a/crates/iceberg/src/cow_rewrite/plan.rs b/crates/iceberg/src/cow_rewrite/plan.rs new file mode 100644 index 0000000000..1ea11998ca --- /dev/null +++ b/crates/iceberg/src/cow_rewrite/plan.rs @@ -0,0 +1,153 @@ +// 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. + +use futures::TryStreamExt; + +use crate::scan::FileScanTask; +use crate::spec::DataFile; + +/// A data file selected for COW rewrite. +#[derive(Debug, Clone)] +pub struct CowRewriteFile { + /// Original data file from the manifest entry. + pub(crate) old_data_file: DataFile, + /// Full-read task for visible rows in the original file. + pub(crate) scan_task: FileScanTask, +} + +impl CowRewriteFile { + /// Original data file from the manifest entry. + pub fn old_data_file(&self) -> &DataFile { + &self.old_data_file + } + + /// Full-read task for visible rows in the original file. + pub fn scan_task(&self) -> &FileScanTask { + &self.scan_task + } +} + +/// Plan data files that may need copy-on-write rewrite. +pub(crate) async fn plan_cow_rewrite_files( + table: &crate::table::Table, + predicate: Option, + snapshot_id: Option, + case_sensitive: bool, +) -> crate::Result> { + let mut scan_builder = table + .scan() + .select_all() + .with_case_sensitive(case_sensitive); + + if let Some(predicate) = predicate { + scan_builder = scan_builder.with_filter(predicate); + } + + if let Some(snapshot_id) = snapshot_id { + scan_builder = scan_builder.snapshot_id(snapshot_id); + } + + scan_builder + .build()? + .plan_cow_rewrite_files() + .await? + .try_collect() + .await +} + +#[cfg(test)] +mod tests { + use futures::TryStreamExt; + + use crate::Result; + use crate::expr::Predicate; + use crate::scan::tests::TableTestFixture; + + #[tokio::test] + async fn cow_planner_returns_old_file_and_full_read_task() -> Result<()> { + let mut fixture = TableTestFixture::new(); + fixture.setup_manifest_files().await; + + let mut files = + super::plan_cow_rewrite_files(&fixture.table, Some(Predicate::AlwaysTrue), None, true) + .await?; + + assert_eq!(files.len(), 2); + + files.sort_by_key(|file| file.old_data_file.file_path().to_string()); + assert_eq!( + files[0].old_data_file.file_path(), + format!("{}/1.parquet", fixture.table_location) + ); + assert_eq!( + files[1].old_data_file.file_path(), + format!("{}/3.parquet", fixture.table_location) + ); + + for file in files { + assert_eq!( + file.old_data_file.file_path(), + file.scan_task.data_file_path() + ); + assert!(file.scan_task.predicate.is_none()); + assert_eq!( + Some(file.old_data_file.record_count()), + file.scan_task.record_count + ); + assert_eq!(0, file.scan_task.start); + assert_eq!( + file.old_data_file.file_size_in_bytes(), + file.scan_task.length + ); + } + + Ok(()) + } + + #[tokio::test] + async fn cow_planner_preserves_delete_files() -> Result<()> { + let mut fixture = TableTestFixture::new(); + fixture.setup_deadlock_manifests().await; + + let scan = fixture + .table + .scan() + .select_all() + .with_concurrency_limit(1) + .build()?; + + let files = tokio::time::timeout(std::time::Duration::from_secs(5), async { + scan.plan_cow_rewrite_files() + .await? + .try_collect::>() + .await + }) + .await + .expect("COW planning should not deadlock")?; + + assert_eq!(files.len(), 10); + for file in files { + assert_eq!(file.scan_task.deletes.len(), 1); + assert_eq!( + file.scan_task.deletes[0].file_path, + format!("{}/del.parquet", fixture.table_location) + ); + } + + Ok(()) + } +} diff --git a/crates/iceberg/src/cow_rewrite/rewriter.rs b/crates/iceberg/src/cow_rewrite/rewriter.rs new file mode 100644 index 0000000000..c307284f04 --- /dev/null +++ b/crates/iceberg/src/cow_rewrite/rewriter.rs @@ -0,0 +1,42 @@ +// 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. + +use arrow_array::RecordBatch; + +use crate::Result; + +/// Result of rewriting a single record batch. +pub struct CowBatchRewrite { + /// Rewritten output batch, or `None` when the input batch is fully removed. + /// + /// Output batches must use a schema compatible with the table schema and + /// must preserve the source file's partition values. This primitive writes + /// replacements into the source file's partition and does not repartition + /// rows. + pub output: Option, + /// Whether the rewrite changed the input batch contents. + /// + /// Set this to `true` whenever `output` differs from the input batch, + /// including filtered rows, updated values, reordered rows, or `None`. + pub changed: bool, +} + +/// Rewrites record batches for copy-on-write operations. +pub trait CowBatchRewriter: Send + Sync { + /// Rewrites a record batch and reports whether it changed. + fn rewrite_batch(&self, batch: RecordBatch) -> Result; +} diff --git a/crates/iceberg/src/cow_rewrite/writer.rs b/crates/iceberg/src/cow_rewrite/writer.rs new file mode 100644 index 0000000000..0c4c0b9ad9 --- /dev/null +++ b/crates/iceberg/src/cow_rewrite/writer.rs @@ -0,0 +1,334 @@ +// 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. + +use uuid::Uuid; + +use crate::Result; +#[cfg(test)] +use crate::spec::DataFile; +use crate::spec::{DataFileFormat, PartitionKey, SchemaRef}; +use crate::table::Table; +use crate::writer::IcebergWriterBuilder; +use crate::writer::base_writer::data_file_writer::DataFileWriterBuilder; +use crate::writer::file_writer::ParquetWriterBuilder; +use crate::writer::file_writer::location_generator::{ + DefaultFileNameGenerator, DefaultLocationGenerator, +}; +use crate::writer::file_writer::rolling_writer::RollingFileWriterBuilder; + +/// Builds a boxed replacement-data-file writer. +/// +/// `write_schema` is the schema the input batches are encoded in. It must match +/// the schema the rows were read in (the planned snapshot's schema), not the +/// table's possibly-evolved current schema; otherwise the parquet writer will +/// reject batches that lack columns added after the source files were written. +/// +/// Building the writer is cheap: no physical file is opened until the first +/// batch is written, so it is safe to construct one optimistically and only +/// write to it once a source file is known to have changed. +pub(crate) async fn build_replacement_writer( + table: &Table, + write_schema: SchemaRef, + partition_key: Option, +) -> Result> { + let location_generator = DefaultLocationGenerator::new(table.metadata())?; + let file_name_generator = DefaultFileNameGenerator::new( + format!("cow-rewrite-{}", Uuid::now_v7()), + None, + DataFileFormat::Parquet, + ); + let table_props = table.metadata().table_properties()?; + let parquet_builder = ParquetWriterBuilder::from_table_properties(&table_props, write_schema); + let rolling_builder = RollingFileWriterBuilder::new( + parquet_builder, + table_props.write_target_file_size_bytes, + table.file_io().clone(), + location_generator, + file_name_generator, + ); + let data_writer_builder = DataFileWriterBuilder::new(rolling_builder); + Ok(Box::new(data_writer_builder.build(partition_key).await?)) +} + +/// Writes replacement record batches as Iceberg data files. +/// +/// This is a convenience wrapper around [`build_replacement_writer`] that +/// consumes a stream end-to-end. Empty input or all-zero-row batches produce no +/// data files. +#[cfg(test)] +pub(crate) async fn write_replacement_batches( + table: &Table, + write_schema: SchemaRef, + partition_key: Option, + mut batches: S, +) -> Result> +where + S: futures::Stream> + Unpin, +{ + use futures::TryStreamExt as _; + + let mut writer = build_replacement_writer(table, write_schema, partition_key).await?; + + let mut wrote_rows = false; + while let Some(batch) = batches.try_next().await? { + if batch.num_rows() == 0 { + continue; + } + wrote_rows = true; + writer.write(batch).await?; + } + + if !wrote_rows { + return Ok(vec![]); + } + + writer.close().await +} + +#[cfg(test)] +mod tests { + use std::collections::HashMap; + use std::sync::Arc; + + use arrow_array::{ArrayRef, Int32Array, RecordBatch}; + use arrow_schema::{DataType, Field, Schema as ArrowSchema}; + use parquet::arrow::PARQUET_FIELD_ID_META_KEY; + use tempfile::TempDir; + + use super::write_replacement_batches; + use crate::io::LocalFsStorageFactory; + use crate::memory::{MEMORY_CATALOG_WAREHOUSE, MemoryCatalogBuilder}; + use crate::spec::{ + DataContentType, NestedField, PartitionKey, PrimitiveType, Schema, Struct, Transform, Type, + }; + use crate::{Catalog, CatalogBuilder, NamespaceIdent, Result, TableCreation}; + + #[tokio::test] + async fn cow_replacement_writer_outputs_data_files() -> Result<()> { + let temp_dir = TempDir::new().unwrap(); + let warehouse = format!("file://{}", temp_dir.path().join("warehouse").display()); + let catalog = MemoryCatalogBuilder::default() + .with_storage_factory(Arc::new(LocalFsStorageFactory)) + .load( + "memory", + HashMap::from([(MEMORY_CATALOG_WAREHOUSE.to_string(), warehouse)]), + ) + .await?; + let namespace = NamespaceIdent::new("ns".to_string()); + catalog.create_namespace(&namespace, HashMap::new()).await?; + + let schema = Schema::builder() + .with_fields(vec![ + NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(), + ]) + .build()?; + let table = catalog + .create_table( + &namespace, + TableCreation::builder() + .name("cow_writer".to_string()) + .schema(schema) + .build(), + ) + .await?; + + let arrow_schema = Arc::new(ArrowSchema::new(vec![ + Field::new("id", DataType::Int32, false).with_metadata(HashMap::from([( + PARQUET_FIELD_ID_META_KEY.to_string(), + "1".to_string(), + )])), + ])); + let batch = RecordBatch::try_new(arrow_schema, vec![Arc::new(Int32Array::from(vec![ + 1, 2, 3, + ])) as ArrayRef])?; + + let data_files = write_replacement_batches( + &table, + table.metadata().current_schema().clone(), + None, + futures::stream::iter(vec![Ok(batch)]), + ) + .await?; + + assert_eq!(data_files.len(), 1); + assert_eq!(data_files[0].record_count(), 3); + assert_eq!(data_files[0].content_type(), DataContentType::Data); + + Ok(()) + } + + #[tokio::test] + async fn cow_replacement_writer_skips_empty_input() -> Result<()> { + let temp_dir = TempDir::new().unwrap(); + let warehouse = format!("file://{}", temp_dir.path().join("warehouse").display()); + let catalog = MemoryCatalogBuilder::default() + .with_storage_factory(Arc::new(LocalFsStorageFactory)) + .load( + "memory", + HashMap::from([(MEMORY_CATALOG_WAREHOUSE.to_string(), warehouse)]), + ) + .await?; + let namespace = NamespaceIdent::new("ns".to_string()); + catalog.create_namespace(&namespace, HashMap::new()).await?; + + let schema = Schema::builder() + .with_fields(vec![ + NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(), + ]) + .build()?; + let table = catalog + .create_table( + &namespace, + TableCreation::builder() + .name("cow_empty_writer".to_string()) + .schema(schema) + .build(), + ) + .await?; + + let data_files = write_replacement_batches( + &table, + table.metadata().current_schema().clone(), + None, + futures::stream::iter(vec![]), + ) + .await?; + + assert!(data_files.is_empty()); + + Ok(()) + } + + #[tokio::test] + async fn cow_replacement_writer_skips_zero_row_batches() -> Result<()> { + let temp_dir = TempDir::new().unwrap(); + let warehouse = format!("file://{}", temp_dir.path().join("warehouse").display()); + let catalog = MemoryCatalogBuilder::default() + .with_storage_factory(Arc::new(LocalFsStorageFactory)) + .load( + "memory", + HashMap::from([(MEMORY_CATALOG_WAREHOUSE.to_string(), warehouse)]), + ) + .await?; + let namespace = NamespaceIdent::new("ns".to_string()); + catalog.create_namespace(&namespace, HashMap::new()).await?; + + let schema = Schema::builder() + .with_fields(vec![ + NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(), + ]) + .build()?; + let table = catalog + .create_table( + &namespace, + TableCreation::builder() + .name("cow_zero_row_writer".to_string()) + .schema(schema) + .build(), + ) + .await?; + + let arrow_schema = Arc::new(ArrowSchema::new(vec![ + Field::new("id", DataType::Int32, false).with_metadata(HashMap::from([( + PARQUET_FIELD_ID_META_KEY.to_string(), + "1".to_string(), + )])), + ])); + let batch = + RecordBatch::try_new(arrow_schema, vec![ + Arc::new(Int32Array::from(Vec::::new())) as ArrayRef, + ])?; + + let data_files = write_replacement_batches( + &table, + table.metadata().current_schema().clone(), + None, + futures::stream::iter(vec![Ok(batch)]), + ) + .await?; + + assert!(data_files.is_empty()); + + Ok(()) + } + + #[tokio::test] + async fn cow_replacement_writer_preserves_partition_key() -> Result<()> { + let temp_dir = TempDir::new().unwrap(); + let warehouse = format!("file://{}", temp_dir.path().join("warehouse").display()); + let catalog = MemoryCatalogBuilder::default() + .with_storage_factory(Arc::new(LocalFsStorageFactory)) + .load( + "memory", + HashMap::from([(MEMORY_CATALOG_WAREHOUSE.to_string(), warehouse)]), + ) + .await?; + let namespace = NamespaceIdent::new("ns".to_string()); + catalog.create_namespace(&namespace, HashMap::new()).await?; + + let schema = Schema::builder() + .with_fields(vec![ + NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(), + ]) + .build()?; + let partition_spec = crate::spec::PartitionSpec::builder(Arc::new(schema.clone())) + .add_partition_field("id", "id", Transform::Identity)? + .build()?; + let table = catalog + .create_table( + &namespace, + TableCreation::builder() + .name("cow_partitioned_writer".to_string()) + .schema(schema) + .partition_spec(partition_spec.clone().into_unbound()) + .build(), + ) + .await?; + + let arrow_schema = Arc::new(ArrowSchema::new(vec![ + Field::new("id", DataType::Int32, false).with_metadata(HashMap::from([( + PARQUET_FIELD_ID_META_KEY.to_string(), + "1".to_string(), + )])), + ])); + let batch = RecordBatch::try_new(arrow_schema, vec![Arc::new(Int32Array::from(vec![ + 1, 1, 1, + ])) as ArrayRef])?; + let partition_key = PartitionKey::new( + table.metadata().default_partition_spec().as_ref().clone(), + table.metadata().current_schema().clone(), + Struct::from_iter([Some(crate::spec::Literal::int(1))]), + ); + + let data_files = write_replacement_batches( + &table, + table.metadata().current_schema().clone(), + Some(partition_key), + futures::stream::iter(vec![Ok(batch)]), + ) + .await?; + + assert_eq!(data_files.len(), 1); + assert_eq!(data_files[0].partition_spec_id, 0); + assert_eq!( + data_files[0].partition(), + &Struct::from_iter([Some(crate::spec::Literal::int(1))]) + ); + + Ok(()) + } +} diff --git a/crates/iceberg/src/lib.rs b/crates/iceberg/src/lib.rs index 301992d15e..8cb2afb257 100644 --- a/crates/iceberg/src/lib.rs +++ b/crates/iceberg/src/lib.rs @@ -79,6 +79,7 @@ pub mod table; mod avro; pub mod cache; pub mod compression; +pub mod cow_rewrite; pub mod io; pub mod spec; diff --git a/crates/iceberg/src/scan/context.rs b/crates/iceberg/src/scan/context.rs index 176bcba459..cdb9f0cfb6 100644 --- a/crates/iceberg/src/scan/context.rs +++ b/crates/iceberg/src/scan/context.rs @@ -150,6 +150,18 @@ impl ManifestEntryContext { .with_key_metadata(self.manifest_entry.data_file.key_metadata().map(Box::from)) .build()) } + + /// Consume this `ManifestEntryContext`, returning a COW rewrite file candidate. + pub(crate) async fn into_cow_rewrite_file(self) -> Result { + let old_data_file = self.manifest_entry.data_file().clone(); + let mut scan_task = self.into_file_scan_task().await?; + scan_task.predicate = None; + + Ok(crate::cow_rewrite::CowRewriteFile { + old_data_file, + scan_task, + }) + } } /// PlanContext wraps a [`SnapshotRef`] alongside all the other diff --git a/crates/iceberg/src/scan/mod.rs b/crates/iceberg/src/scan/mod.rs index a1c85c7117..3e07f6ae19 100644 --- a/crates/iceberg/src/scan/mod.rs +++ b/crates/iceberg/src/scan/mod.rs @@ -23,6 +23,7 @@ mod context; use context::*; mod task; +use std::future::Future; use std::sync::Arc; use arrow_array::RecordBatch; @@ -379,6 +380,25 @@ pub struct TableScan { impl TableScan { /// Returns a stream of [`FileScanTask`]s. pub async fn plan_files(&self) -> Result { + self.plan_data_files(|ctx| ctx.into_file_scan_task()).await + } + + pub(crate) async fn plan_cow_rewrite_files( + &self, + ) -> Result>> { + self.plan_data_files(|ctx| ctx.into_cow_rewrite_file()) + .await + } + + async fn plan_data_files( + &self, + build_result: F, + ) -> Result>> + where + T: Send + 'static, + F: Fn(ManifestEntryContext) -> Fut + Copy + Send + Sync + 'static, + Fut: Future> + Send + 'static, + { let Some(plan_context) = self.plan_context.as_ref() else { return Ok(Box::pin(futures::stream::empty())); }; @@ -392,8 +412,8 @@ impl TableScan { let (manifest_entry_delete_ctx_tx, manifest_entry_delete_ctx_rx) = channel(concurrency_limit_manifest_files); - // used to stream the results back to the caller - let (file_scan_task_tx, file_scan_task_rx) = channel(concurrency_limit_manifest_entries); + // used to stream the planned data file results back to the caller + let (planned_file_tx, planned_file_rx) = channel(concurrency_limit_manifest_entries); let (delete_file_idx, delete_file_tx) = DeleteFileIndex::new(self.runtime.clone()); @@ -409,9 +429,9 @@ impl TableScan { manifest_entry_delete_ctx_tx, )?; - let mut channel_for_manifest_error = file_scan_task_tx.clone(); - let mut channel_for_data_manifest_entry_error = file_scan_task_tx.clone(); - let mut channel_for_delete_manifest_entry_error = file_scan_task_tx.clone(); + let mut channel_for_manifest_error = planned_file_tx.clone(); + let mut channel_for_data_manifest_entry_error = planned_file_tx.clone(); + let mut channel_for_delete_manifest_entry_error = planned_file_tx.clone(); let rt = self.runtime.clone(); @@ -468,7 +488,7 @@ impl TableScan { let rt_inner = rt.clone(); rt.cpu().spawn(async move { let result = manifest_entry_data_ctx_rx - .map(|me_ctx| Ok((me_ctx, file_scan_task_tx.clone()))) + .map(|me_ctx| Ok((me_ctx, planned_file_tx.clone()))) .try_for_each_concurrent( concurrency_limit_manifest_entries, |(manifest_entry_context, tx)| { @@ -480,6 +500,7 @@ impl TableScan { Self::process_data_manifest_entry( manifest_entry_context, tx, + build_result, ) .await }) @@ -495,7 +516,7 @@ impl TableScan { }); } - Ok(file_scan_task_rx.boxed()) + Ok(planned_file_rx.boxed()) } /// Returns an [`ArrowRecordBatchStream`]. @@ -526,10 +547,15 @@ impl TableScan { self.plan_context.as_ref().map(|x| &x.snapshot) } - async fn process_data_manifest_entry( + async fn process_data_manifest_entry( manifest_entry_context: ManifestEntryContext, - mut file_scan_task_tx: Sender>, - ) -> Result<()> { + mut planned_file_tx: Sender>, + build_result: F, + ) -> Result<()> + where + F: Fn(ManifestEntryContext) -> Fut + Send + Sync + 'static, + Fut: Future> + Send + 'static, + { // skip processing this manifest entry if it has been marked as deleted if !manifest_entry_context.manifest_entry.is_alive() { return Ok(()); @@ -574,10 +600,10 @@ impl TableScan { } // congratulations! the manifest entry has made its way through the - // entire plan without getting filtered out. Create a corresponding - // FileScanTask and push it to the result stream - file_scan_task_tx - .send(Ok(manifest_entry_context.into_file_scan_task().await?)) + // entire plan without getting filtered out. Build the planned file before + // sending so delete-file lookup preserves the original scan timing. + planned_file_tx + .send(Ok(build_result(manifest_entry_context).await?)) .await?; Ok(())