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


##########
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:
   This documents the snapshot-schema choice as a deliberate divergence from 
Java's `RewriteDataFiles`, which is exactly right. Two sibling divergences are 
worth the same explicit note: replacements are written unsorted regardless of 
table sort order (`sort_order_id` stays 0), and they retain the source file's 
partition spec rather than repartitioning to the current default. Both are 
legal and readable everywhere — just worth a doc line each so the divergence is 
intentional on the record rather than surprising someone who diffs this against 
the Java action. Non-blocking.



##########
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:
   The lazy-writer buffer inside this call holds every `output: Some` batch in 
`prefix` until the first changed batch — so a file that never changes, or whose 
first change sits near EOF, materializes its entire decoded contents in memory 
before dropping them. The internal comment already calls this out and defers a 
size-capped fallback, which is the right v1 scope. I'd just lift that memory 
ceiling up into this public `rewrite()` doc so someone compacting multi-GB 
files knows the bound before they hit an OOM. Not blocking.



##########
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:
   Partition metadata is taken verbatim from the source file here, so a 
rewriter that mutates a partition-source column would produce replacement rows 
whose values disagree with the recorded partition — silent wrong results under 
partition filters, with nothing to catch it. The doc contract says rewriters 
must preserve partition values, but it's unenforced. A debug-assert that the 
output still maps to this key (or a note that the commit adapter is responsible 
for validating) would keep this a footgun rather than a corruption path. 
Follow-up, not a gate.



-- 
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