u70b3 commented on code in PR #2752:
URL: https://github.com/apache/iceberg-rust/pull/2752#discussion_r4130540645


##########
crates/iceberg/src/cow_rewrite/mod.rs:
##########
@@ -0,0 +1,1445 @@
+// 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 schema rows were read in
+//! (the planned snapshot's schema) and must preserve each source file's
+//! partition values; this primitive does not repartition rewritten rows.

Review Comment:
   Done, documented the source partition spec and unsorted output.



##########
crates/iceberg/src/cow_rewrite/mod.rs:
##########
@@ -0,0 +1,1445 @@
+// 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 schema rows were read in
+//! (the planned snapshot's schema) and must preserve each source file's
+//! partition values; this primitive does not repartition rewritten rows.
+//!
+//! The result carries data files only. A commit adapter consuming these file
+//! lists may remove position delete files and deletion vectors only when they
+//! exclusively reference removed data files. Equality delete files must be
+//! retained while they can still apply to other live data files by sequence
+//! number.
+//!
+//! ```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<CowBatchRewrite> 
{
+//!         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 written to replacement files.
+    ///
+    /// Rows the rewriter emitted for files that turned out unchanged are not
+    /// counted, so this always matches the row counts of `added_data_files`.
+    pub output_rows: u64,
+    /// Number of input batches that changed, including batches the rewriter
+    /// dropped entirely (`output: None`) even if it did not flag them.
+    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<DataFile>,
+    /// New data files that should be added by the commit action.
+    pub added_data_files: Vec<DataFile>,
+    /// Candidate files that were read and left unchanged.
+    ///
+    /// Files whose visible rows were all removed by delete files are NOT
+    /// included here: they read as zero rows and are reported in
+    /// `removed_data_files` with no replacement. A commit adapter may remove
+    /// position deletes and deletion vectors that exclusively reference these
+    /// removed files, but must retain equality deletes that can still apply to
+    /// other live files.
+    pub unchanged_data_files: Vec<DataFile>,
+    /// 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<i64>,
+    batch_size: Option<usize>,
+    case_sensitive: bool,
+    rewriter: Option<Arc<dyn CowBatchRewriter>>,
+}
+
+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 the case sensitivity used to bind the planning predicate.
+    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<dyn CowBatchRewriter>) -> 
Self {
+        self.rewriter = Some(rewriter);
+        self
+    }
+
+    /// Plans, reads, rewrites, and writes replacement data files.

Review Comment:
   Done, documented the whole-file memory bound on rewrite().



##########
crates/iceberg/src/cow_rewrite/mod.rs:
##########
@@ -0,0 +1,1445 @@
+// 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 schema rows were read in
+//! (the planned snapshot's schema) and must preserve each source file's
+//! partition values; this primitive does not repartition rewritten rows.
+//!
+//! The result carries data files only. A commit adapter consuming these file
+//! lists may remove position delete files and deletion vectors only when they
+//! exclusively reference removed data files. Equality delete files must be
+//! retained while they can still apply to other live data files by sequence
+//! number.
+//!
+//! ```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<CowBatchRewrite> 
{
+//!         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 written to replacement files.
+    ///
+    /// Rows the rewriter emitted for files that turned out unchanged are not
+    /// counted, so this always matches the row counts of `added_data_files`.
+    pub output_rows: u64,
+    /// Number of input batches that changed, including batches the rewriter
+    /// dropped entirely (`output: None`) even if it did not flag them.
+    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<DataFile>,
+    /// New data files that should be added by the commit action.
+    pub added_data_files: Vec<DataFile>,
+    /// Candidate files that were read and left unchanged.
+    ///
+    /// Files whose visible rows were all removed by delete files are NOT
+    /// included here: they read as zero rows and are reported in
+    /// `removed_data_files` with no replacement. A commit adapter may remove
+    /// position deletes and deletion vectors that exclusively reference these
+    /// removed files, but must retain equality deletes that can still apply to
+    /// other live files.
+    pub unchanged_data_files: Vec<DataFile>,
+    /// 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<i64>,
+    batch_size: Option<usize>,
+    case_sensitive: bool,
+    rewriter: Option<Arc<dyn CowBatchRewriter>>,
+}
+
+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 the case sensitivity used to bind the planning predicate.
+    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<dyn CowBatchRewriter>) -> 
Self {
+        self.rewriter = Some(rewriter);
+        self
+    }
+
+    /// Plans, reads, rewrites, and writes replacement data files.
+    pub async fn rewrite(self) -> Result<CowRewriteResult> {
+        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 {
+            let CowRewriteFile {
+                old_data_file,
+                scan_task,
+            } = file;
+            // 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. This preserves the 
older
+            // schema in replacement files rather than promoting them to the
+            // table's current schema during this rewrite.
+            let write_schema = scan_task.schema_ref();
+            let has_delete_files = !scan_task.deletes().is_empty();
+
+            // 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.
+            //
+            // Worst-case footprint: for a file that never changes (or whose
+            // first change sits at its very end) the prefix holds the entire
+            // decoded source file in memory. Files are processed sequentially,
+            // so peak usage is one file at a time, but that can still be 
several
+            // GB for a compaction-sized file. A size-capped fallback that 
starts
+            // writing the replacement once the buffer crosses a threshold is
+            // left for follow-up work.
+            let mut prefix: Vec<RecordBatch> = Vec::new();
+            let mut file_changed = false;
+            let mut file_input_rows = 0_u64;
+            let mut file_output_rows = 0_u64;
+            let mut writer: Option<Box<dyn crate::writer::IcebergWriter>> = 
None;
+
+            // Planning already cleared the row predicate (see
+            // `ManifestEntryContext::into_cow_rewrite_file`), so this task
+            // reads every row of the source file.
+            let tasks = Box::pin(futures::stream::iter(vec![Ok(scan_task)])) 
as FileScanTaskStream;
+
+            // Each candidate file gets its own reader so the per-file prefix
+            // and lazy-writer semantics stay intact; the delete-file cache is
+            // therefore also per file, and equality deletes shared by several
+            // candidates are fetched once per file.
+            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;
+                file_input_rows += batch.num_rows() as u64;
+
+                let rewrite = rewriter.rewrite_batch(batch)?;
+                // `output: None` means the batch is fully removed, which is
+                // itself a change. Derive the effective flag instead of
+                // trusting every rewriter to keep `changed` consistent with
+                // `output` — otherwise a `{changed: false, output: None}`
+                // batch would silently drop its rows while leaving the file
+                // marked unchanged.
+                let changed = rewrite.changed || rewrite.output.is_none();
+                if changed {
+                    file_changed = true;
+                    result.stats.changed_batches += 1;
+                }
+
+                // A dropped batch can be the first change. Flush any kept
+                // prefix even when this batch has no output; defer opening the
+                // writer only if there are no rows to preserve yet.
+                if file_changed
+                    && writer.is_none()
+                    && (!prefix.is_empty() || rewrite.output.is_some())
+                {
+                    let partition_key =
+                        source_partition_key(self.table, &old_data_file, 
&write_schema)?;
+                    writer = Some(
+                        writer::build_replacement_writer(
+                            self.table,
+                            write_schema.clone(),
+                            Some(partition_key),
+                        )
+                        .await?,
+                    );
+                }
+
+                if let Some(writer) = writer.as_mut() {
+                    for prefix_batch in prefix.drain(..) {
+                        writer.write(prefix_batch).await?;
+                    }
+                    if let Some(output) = rewrite.output {
+                        file_output_rows += output.num_rows() as u64;
+                        writer.write(output).await?;
+                    }
+                } else if let Some(output) = rewrite.output {
+                    file_output_rows += output.num_rows() as u64;
+                    prefix.push(output);
+                }
+            }
+
+            // A candidate whose visible rows were all removed by its delete
+            // files reads as zero rows and the loop above never runs. Treat it
+            // the same as a rewriter that dropped every batch — changed with
+            // no replacement — so the file and its delete files can be
+            // compacted away instead of being pinned in the table forever.
+            let fully_removed_by_deletes =
+                !file_changed && file_input_rows == 0 && has_delete_files;
+
+            if file_changed || fully_removed_by_deletes {
+                result.stats.rewritten_files += 1;
+                result.stats.output_rows += file_output_rows;
+                result.removed_data_files.push(old_data_file);
+
+                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`, no visible rows remain after delete
+                // files and batch rewriting, so no replacement is written.
+            } else {
+                result.stats.unchanged_files += 1;
+                result.unchanged_data_files.push(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<PartitionKey> {
+    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();
+    // `PartitionKey::new` does not bind the spec to the schema, so validate 
the
+    // binding here: a spec that is incompatible with the planned snapshot
+    // schema should fail with a clear error instead of producing a bad
+    // partition path when the writer later calls `PartitionKey::to_path`.
+    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(

Review Comment:
   Done, added a debug partition guard and regression test.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to