jopdorp commented on code in PR #3300:
URL: https://github.com/apache/iceberg-rust/pull/3300#discussion_r4168726469


##########
crates/iceberg/src/spec/manifest/writer.rs:
##########
@@ -361,6 +361,19 @@ impl ManifestWriter {
         Ok(())
     }
 
+    /// Add a manifest entry as a tombstone, preserving its original sequence 
numbers and
+    /// snapshot id.
+    ///
+    /// The caller owns `snapshot_id`: set it to the deleting snapshot for a 
file this commit
+    /// removes, and leave it untouched when carrying an earlier tombstone 
forward. Use
+    /// `add_delete_entry` instead to have this manifest's snapshot id stamped 
on the entry.

Review Comment:
   `add_tombstone_entry` is gone: the rewrite now only carries live entries, 
like Java's `ManifestFilterManager`, so it uses the existing `add_delete_entry` 
/ `add_existing_entry` and both lose their `allow(dead_code)`.



##########
crates/iceberg/src/transaction/overwrite.rs:
##########
@@ -0,0 +1,912 @@
+// 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 std::collections::{HashMap, HashSet};
+use std::sync::Arc;
+
+use async_trait::async_trait;
+use uuid::Uuid;
+
+use crate::error::Result;
+use crate::spec::{
+    DataFile, Manifest, ManifestEntry, ManifestFile, Operation, 
PartitionSpecRef, SchemaRef,
+};
+use crate::table::Table;
+use crate::transaction::snapshot::{
+    DefaultManifestProcess, SnapshotProduceOperation, SnapshotProducer,
+};
+use crate::transaction::{ActionCommit, TransactionAction};
+use crate::{Error, ErrorKind};
+
+/// OverwriteAction is a transaction action for overwriting data files in the 
table.
+///
+/// Creates a snapshot with `Operation::Overwrite` semantics — adds new data 
files and
+/// optionally removes existing data files by rewriting affected manifests 
with those
+/// entries marked as `ManifestStatus::Deleted`.
+pub struct OverwriteAction {
+    check_duplicate: bool,
+    commit_uuid: Option<Uuid>,
+    snapshot_properties: HashMap<String, String>,
+    added_data_files: Vec<DataFile>,
+    deleted_data_files: Vec<DataFile>,
+}
+
+impl OverwriteAction {
+    pub(crate) fn new() -> Self {
+        Self {
+            check_duplicate: true,
+            commit_uuid: None,
+            snapshot_properties: HashMap::default(),
+            added_data_files: vec![],
+            deleted_data_files: vec![],
+        }
+    }
+
+    /// Set whether to check duplicate files.
+    pub fn with_check_duplicate(mut self, v: bool) -> Self {
+        self.check_duplicate = v;
+        self
+    }
+
+    /// Add data files to the snapshot.
+    pub fn add_data_files(mut self, data_files: impl IntoIterator<Item = 
DataFile>) -> Self {
+        self.added_data_files.extend(data_files);
+        self
+    }
+
+    /// Specify data files to be removed from the table in this overwrite.
+    pub fn delete_data_files(mut self, data_files: impl IntoIterator<Item = 
DataFile>) -> Self {
+        self.deleted_data_files.extend(data_files);
+        self
+    }
+
+    /// Set commit UUID for the snapshot.
+    pub fn set_commit_uuid(mut self, commit_uuid: Uuid) -> Self {
+        self.commit_uuid = Some(commit_uuid);
+        self
+    }
+
+    /// Set snapshot summary properties.
+    pub fn set_snapshot_properties(mut self, snapshot_properties: 
HashMap<String, String>) -> Self {
+        self.snapshot_properties = snapshot_properties;
+        self
+    }
+}
+
+#[async_trait]
+impl TransactionAction for OverwriteAction {
+    async fn commit(self: Arc<Self>, table: &Table) -> Result<ActionCommit> {
+        let snapshot_producer = SnapshotProducer::new(
+            table,
+            self.commit_uuid.unwrap_or_else(Uuid::now_v7),
+            self.snapshot_properties.clone(),
+            self.added_data_files.clone(),

Review Comment:
   Pushed it down into `SnapshotProducer::new` for both added and deleted 
files, so fast append and overwrite share it.



##########
crates/iceberg/src/transaction/overwrite.rs:
##########
@@ -0,0 +1,912 @@
+// 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 std::collections::{HashMap, HashSet};
+use std::sync::Arc;
+
+use async_trait::async_trait;
+use uuid::Uuid;
+
+use crate::error::Result;
+use crate::spec::{
+    DataFile, Manifest, ManifestEntry, ManifestFile, Operation, 
PartitionSpecRef, SchemaRef,
+};
+use crate::table::Table;
+use crate::transaction::snapshot::{
+    DefaultManifestProcess, SnapshotProduceOperation, SnapshotProducer,
+};
+use crate::transaction::{ActionCommit, TransactionAction};
+use crate::{Error, ErrorKind};
+
+/// OverwriteAction is a transaction action for overwriting data files in the 
table.
+///
+/// Creates a snapshot with `Operation::Overwrite` semantics — adds new data 
files and
+/// optionally removes existing data files by rewriting affected manifests 
with those
+/// entries marked as `ManifestStatus::Deleted`.
+pub struct OverwriteAction {
+    check_duplicate: bool,
+    commit_uuid: Option<Uuid>,
+    snapshot_properties: HashMap<String, String>,
+    added_data_files: Vec<DataFile>,
+    deleted_data_files: Vec<DataFile>,
+}
+
+impl OverwriteAction {
+    pub(crate) fn new() -> Self {
+        Self {
+            check_duplicate: true,
+            commit_uuid: None,
+            snapshot_properties: HashMap::default(),
+            added_data_files: vec![],
+            deleted_data_files: vec![],
+        }
+    }
+
+    /// Set whether to check duplicate files.
+    pub fn with_check_duplicate(mut self, v: bool) -> Self {
+        self.check_duplicate = v;
+        self
+    }
+
+    /// Add data files to the snapshot.
+    pub fn add_data_files(mut self, data_files: impl IntoIterator<Item = 
DataFile>) -> Self {
+        self.added_data_files.extend(data_files);
+        self
+    }
+
+    /// Specify data files to be removed from the table in this overwrite.
+    pub fn delete_data_files(mut self, data_files: impl IntoIterator<Item = 
DataFile>) -> Self {
+        self.deleted_data_files.extend(data_files);
+        self
+    }
+
+    /// Set commit UUID for the snapshot.
+    pub fn set_commit_uuid(mut self, commit_uuid: Uuid) -> Self {
+        self.commit_uuid = Some(commit_uuid);
+        self
+    }
+
+    /// Set snapshot summary properties.
+    pub fn set_snapshot_properties(mut self, snapshot_properties: 
HashMap<String, String>) -> Self {
+        self.snapshot_properties = snapshot_properties;
+        self
+    }
+}
+
+#[async_trait]
+impl TransactionAction for OverwriteAction {
+    async fn commit(self: Arc<Self>, table: &Table) -> Result<ActionCommit> {
+        let snapshot_producer = SnapshotProducer::new(
+            table,
+            self.commit_uuid.unwrap_or_else(Uuid::now_v7),
+            self.snapshot_properties.clone(),
+            self.added_data_files.clone(),
+            self.deleted_data_files.clone(),
+        );
+
+        snapshot_producer.validate_added_data_files()?;

Review Comment:
   Dropped the "for fast append".



##########
crates/iceberg/src/transaction/overwrite.rs:
##########
@@ -0,0 +1,912 @@
+// 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 std::collections::{HashMap, HashSet};
+use std::sync::Arc;
+
+use async_trait::async_trait;
+use uuid::Uuid;
+
+use crate::error::Result;
+use crate::spec::{
+    DataFile, Manifest, ManifestEntry, ManifestFile, Operation, 
PartitionSpecRef, SchemaRef,
+};
+use crate::table::Table;
+use crate::transaction::snapshot::{
+    DefaultManifestProcess, SnapshotProduceOperation, SnapshotProducer,
+};
+use crate::transaction::{ActionCommit, TransactionAction};
+use crate::{Error, ErrorKind};
+
+/// OverwriteAction is a transaction action for overwriting data files in the 
table.
+///
+/// Creates a snapshot with `Operation::Overwrite` semantics — adds new data 
files and
+/// optionally removes existing data files by rewriting affected manifests 
with those
+/// entries marked as `ManifestStatus::Deleted`.
+pub struct OverwriteAction {
+    check_duplicate: bool,
+    commit_uuid: Option<Uuid>,
+    snapshot_properties: HashMap<String, String>,
+    added_data_files: Vec<DataFile>,
+    deleted_data_files: Vec<DataFile>,
+}
+
+impl OverwriteAction {
+    pub(crate) fn new() -> Self {
+        Self {
+            check_duplicate: true,
+            commit_uuid: None,
+            snapshot_properties: HashMap::default(),
+            added_data_files: vec![],
+            deleted_data_files: vec![],
+        }
+    }
+
+    /// Set whether to check duplicate files.
+    pub fn with_check_duplicate(mut self, v: bool) -> Self {
+        self.check_duplicate = v;
+        self
+    }
+
+    /// Add data files to the snapshot.
+    pub fn add_data_files(mut self, data_files: impl IntoIterator<Item = 
DataFile>) -> Self {
+        self.added_data_files.extend(data_files);
+        self
+    }
+
+    /// Specify data files to be removed from the table in this overwrite.
+    pub fn delete_data_files(mut self, data_files: impl IntoIterator<Item = 
DataFile>) -> Self {
+        self.deleted_data_files.extend(data_files);
+        self
+    }
+
+    /// Set commit UUID for the snapshot.
+    pub fn set_commit_uuid(mut self, commit_uuid: Uuid) -> Self {
+        self.commit_uuid = Some(commit_uuid);
+        self
+    }
+
+    /// Set snapshot summary properties.
+    pub fn set_snapshot_properties(mut self, snapshot_properties: 
HashMap<String, String>) -> Self {
+        self.snapshot_properties = snapshot_properties;
+        self
+    }
+}
+
+#[async_trait]
+impl TransactionAction for OverwriteAction {
+    async fn commit(self: Arc<Self>, table: &Table) -> Result<ActionCommit> {
+        let snapshot_producer = SnapshotProducer::new(
+            table,
+            self.commit_uuid.unwrap_or_else(Uuid::now_v7),
+            self.snapshot_properties.clone(),
+            self.added_data_files.clone(),
+            self.deleted_data_files.clone(),
+        );
+
+        snapshot_producer.validate_added_data_files()?;
+
+        if self.check_duplicate {
+            snapshot_producer.validate_duplicate_files().await?;
+        }
+
+        let deleted_file_paths: HashSet<String> = self
+            .deleted_data_files
+            .iter()
+            .map(|f| f.file_path.clone())
+            .collect();
+
+        let snapshot_id = snapshot_producer.snapshot_id();
+        snapshot_producer
+            .commit(
+                OverwriteOperation {
+                    deleted_file_paths,
+                    has_added_data_files: !self.added_data_files.is_empty(),
+                    snapshot_id,
+                    removed_data_files: vec![],
+                },
+                DefaultManifestProcess,
+            )
+            .await
+    }
+}
+
+struct OverwriteOperation {
+    deleted_file_paths: HashSet<String>,
+    has_added_data_files: bool,
+    snapshot_id: i64,
+    removed_data_files: Vec<(DataFile, SchemaRef, PartitionSpecRef)>,
+}
+
+impl SnapshotProduceOperation for OverwriteOperation {
+    fn operation(&self) -> Operation {
+        match (
+            self.has_added_data_files,
+            !self.deleted_file_paths.is_empty(),
+        ) {
+            (true, true) => Operation::Overwrite,
+            (false, true) => Operation::Delete,
+            _ => Operation::Append,

Review Comment:
   Added a one-line comment on the `_` arm.



##########
crates/iceberg/src/transaction/overwrite.rs:
##########
@@ -0,0 +1,912 @@
+// 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 std::collections::{HashMap, HashSet};
+use std::sync::Arc;
+
+use async_trait::async_trait;
+use uuid::Uuid;
+
+use crate::error::Result;
+use crate::spec::{
+    DataFile, Manifest, ManifestEntry, ManifestFile, Operation, 
PartitionSpecRef, SchemaRef,
+};
+use crate::table::Table;
+use crate::transaction::snapshot::{
+    DefaultManifestProcess, SnapshotProduceOperation, SnapshotProducer,
+};
+use crate::transaction::{ActionCommit, TransactionAction};
+use crate::{Error, ErrorKind};
+
+/// OverwriteAction is a transaction action for overwriting data files in the 
table.
+///
+/// Creates a snapshot with `Operation::Overwrite` semantics — adds new data 
files and
+/// optionally removes existing data files by rewriting affected manifests 
with those
+/// entries marked as `ManifestStatus::Deleted`.
+pub struct OverwriteAction {
+    check_duplicate: bool,
+    commit_uuid: Option<Uuid>,
+    snapshot_properties: HashMap<String, String>,
+    added_data_files: Vec<DataFile>,
+    deleted_data_files: Vec<DataFile>,
+}
+
+impl OverwriteAction {
+    pub(crate) fn new() -> Self {
+        Self {
+            check_duplicate: true,
+            commit_uuid: None,
+            snapshot_properties: HashMap::default(),
+            added_data_files: vec![],
+            deleted_data_files: vec![],
+        }
+    }
+
+    /// Set whether to check duplicate files.
+    pub fn with_check_duplicate(mut self, v: bool) -> Self {
+        self.check_duplicate = v;
+        self
+    }
+
+    /// Add data files to the snapshot.
+    pub fn add_data_files(mut self, data_files: impl IntoIterator<Item = 
DataFile>) -> Self {
+        self.added_data_files.extend(data_files);
+        self
+    }
+
+    /// Specify data files to be removed from the table in this overwrite.
+    pub fn delete_data_files(mut self, data_files: impl IntoIterator<Item = 
DataFile>) -> Self {
+        self.deleted_data_files.extend(data_files);
+        self
+    }
+
+    /// Set commit UUID for the snapshot.
+    pub fn set_commit_uuid(mut self, commit_uuid: Uuid) -> Self {
+        self.commit_uuid = Some(commit_uuid);
+        self
+    }
+
+    /// Set snapshot summary properties.
+    pub fn set_snapshot_properties(mut self, snapshot_properties: 
HashMap<String, String>) -> Self {
+        self.snapshot_properties = snapshot_properties;
+        self
+    }
+}
+
+#[async_trait]
+impl TransactionAction for OverwriteAction {
+    async fn commit(self: Arc<Self>, table: &Table) -> Result<ActionCommit> {
+        let snapshot_producer = SnapshotProducer::new(
+            table,
+            self.commit_uuid.unwrap_or_else(Uuid::now_v7),
+            self.snapshot_properties.clone(),
+            self.added_data_files.clone(),
+            self.deleted_data_files.clone(),
+        );
+
+        snapshot_producer.validate_added_data_files()?;
+
+        if self.check_duplicate {
+            snapshot_producer.validate_duplicate_files().await?;
+        }
+
+        let deleted_file_paths: HashSet<String> = self
+            .deleted_data_files
+            .iter()
+            .map(|f| f.file_path.clone())
+            .collect();
+
+        let snapshot_id = snapshot_producer.snapshot_id();
+        snapshot_producer
+            .commit(
+                OverwriteOperation {
+                    deleted_file_paths,
+                    has_added_data_files: !self.added_data_files.is_empty(),
+                    snapshot_id,
+                    removed_data_files: vec![],
+                },
+                DefaultManifestProcess,
+            )
+            .await
+    }
+}
+
+struct OverwriteOperation {
+    deleted_file_paths: HashSet<String>,
+    has_added_data_files: bool,
+    snapshot_id: i64,
+    removed_data_files: Vec<(DataFile, SchemaRef, PartitionSpecRef)>,
+}
+
+impl SnapshotProduceOperation for OverwriteOperation {
+    fn operation(&self) -> Operation {
+        match (
+            self.has_added_data_files,
+            !self.deleted_file_paths.is_empty(),
+        ) {
+            (true, true) => Operation::Overwrite,
+            (false, true) => Operation::Delete,
+            _ => Operation::Append,
+        }
+    }
+
+    async fn delete_entries(
+        &self,
+        _snapshot_produce: &SnapshotProducer<'_>,
+    ) -> Result<Vec<ManifestEntry>> {
+        Ok(vec![])
+    }
+
+    fn removed_data_files(&self) -> &[(DataFile, SchemaRef, PartitionSpecRef)] 
{
+        &self.removed_data_files
+    }
+
+    async fn existing_manifest(
+        &mut self,
+        snapshot_produce: &SnapshotProducer<'_>,
+    ) -> Result<Vec<ManifestFile>> {
+        let Some(snapshot) = 
snapshot_produce.table.metadata().current_snapshot() else {
+            self.validate_deletes_matched()?;
+            return Ok(vec![]);
+        };
+
+        let manifest_list = snapshot_produce
+            .table
+            .manifest_list_reader(snapshot)
+            .load()
+            .await?;
+
+        if self.deleted_file_paths.is_empty() {
+            return Ok(manifest_list
+                .entries()
+                .iter()
+                .filter(|entry| {
+                    // Delete-only manifests record which files were removed 
and must survive
+                    // until `expire_snapshots` cleans them up (see #2148).
+                    entry.has_added_files()
+                        || entry.has_existing_files()
+                        || entry.has_deleted_files()
+                })
+                .cloned()
+                .collect());
+        }
+
+        let mut result = Vec::new();
+
+        for manifest_file in manifest_list.entries() {
+            if !manifest_file.has_added_files()
+                && !manifest_file.has_existing_files()
+                && !manifest_file.has_deleted_files()
+            {
+                continue;
+            }
+
+            let manifest = snapshot_produce
+                .table
+                .manifest_reader()
+                .read(manifest_file)
+                .await?;
+
+            let has_deletes = manifest.entries().iter().any(|entry| {
+                entry.is_alive() && 
self.deleted_file_paths.contains(entry.file_path())
+            });
+
+            if has_deletes {
+                let rewritten = self
+                    .rewrite_manifest(snapshot_produce, manifest_file, 
&manifest)
+                    .await?;
+                result.push(rewritten);
+            } else {
+                result.push(manifest_file.clone());
+            }
+        }
+
+        self.validate_deletes_matched()?;
+
+        Ok(result)
+    }
+}
+
+impl OverwriteOperation {
+    /// Fails when a requested delete path was not live in any manifest, 
matching the existence
+    /// check Java's `BaseOverwriteFiles` performs before committing.
+    fn validate_deletes_matched(&self) -> Result<()> {
+        let removed: HashSet<&str> = self
+            .removed_data_files
+            .iter()
+            .map(|(data_file, _, _)| data_file.file_path.as_str())
+            .collect();
+        let mut missing: Vec<&str> = self
+            .deleted_file_paths
+            .iter()
+            .map(String::as_str)
+            .filter(|path| !removed.contains(path))
+            .collect();
+
+        if missing.is_empty() {
+            return Ok(());
+        }
+
+        missing.sort_unstable();
+        Err(Error::new(
+            ErrorKind::DataInvalid,
+            format!("Missing required files to delete: {}", missing.join(", 
")),
+        ))
+    }
+
+    /// Rewrite a manifest, marking entries whose file paths are in 
`deleted_file_paths`
+    /// as `ManifestStatus::Deleted`.
+    async fn rewrite_manifest(
+        &mut self,
+        snapshot_produce: &SnapshotProducer<'_>,
+        manifest_file: &ManifestFile,
+        manifest: &Manifest,
+    ) -> Result<ManifestFile> {
+        let partition_spec: PartitionSpecRef = 
Arc::new(manifest.metadata().partition_spec.clone());
+        let mut writer = snapshot_produce.new_manifest_writer(
+            manifest_file.content,
+            manifest.metadata().schema.clone(),
+            partition_spec.as_ref().clone(),
+        )?;
+
+        for entry in manifest.entries() {
+            let entry: ManifestEntry = (**entry).clone();
+            if !entry.is_alive() {
+                // A tombstone keeps its original snapshot id so the delete 
stays attributed
+                // to the snapshot that made it.
+                writer.add_tombstone_entry(entry)?;
+            } else if self.deleted_file_paths.contains(entry.file_path()) {
+                let mut deleted = entry;
+                deleted.snapshot_id = Some(self.snapshot_id);
+                self.removed_data_files.push((
+                    deleted.data_file().clone(),
+                    manifest.metadata().schema.clone(),
+                    partition_spec.clone(),
+                ));
+                writer.add_tombstone_entry(deleted)?;
+            } else {
+                writer.add_existing_entry(entry)?;
+            }
+        }
+
+        writer.write_manifest_file().await

Review Comment:
   We looked at this before changing it, and I think the rows are safe: the 
manifest reader gives every live entry its `first_row_id` on read and the v3 
writer writes it out explicitly, so the surviving entries keep their ids in the 
rewritten manifest. In a quick v3 check, b keeps row id 1 while the rewritten 
manifest gets `first_row_id` 2. Java does the same, `ManifestFilterManager` 
writes the filtered manifest with a null `firstRowId` and `ManifestListWriter` 
assigns a fresh one. Happy to carry the original over anyway if you'd rather, 
it's two lines.



##########
crates/iceberg/src/transaction/snapshot.rs:
##########
@@ -113,26 +129,33 @@ pub(crate) struct SnapshotProducer<'a> {
     commit_uuid: Uuid,
     snapshot_properties: HashMap<String, String>,
     added_data_files: Vec<DataFile>,
+    deleted_data_files: Vec<DataFile>,

Review Comment:
   It drives the summary now (from #3253), counted like the added files as 
pyiceberg does, so it's no longer just answering `is_empty()`. The 
`removed_data_files` side is gone.



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