laskoviymishka commented on code in PR #3300:
URL: https://github.com/apache/iceberg-rust/pull/3300#discussion_r4186811702
##########
crates/iceberg/src/transaction/snapshot.rs:
##########
@@ -125,23 +132,30 @@ impl<'a> SnapshotProducer<'a> {
commit_uuid: Uuid,
snapshot_properties: HashMap<String, String>,
added_data_files: Vec<DataFile>,
+ deleted_data_files: Vec<DataFile>,
) -> Self {
Self {
table,
snapshot_id: Self::generate_unique_snapshot_id(table),
commit_uuid,
snapshot_properties,
- added_data_files,
- manifest_counter: (0..),
+ // Collapsed to the first file per path, so a manifest never
references a file twice.
+ added_data_files: added_data_files
+ .into_iter()
+ .unique_by(|f| f.file_path.clone())
Review Comment:
`dedupe_added_files` used to borrow (`HashSet<&str>`); `unique_by(|f|
f.file_path.clone())` now heap-allocates a `String` per file, and since this
moved into the producer every `FastAppendAction` commit pays it too. A `retain`
with a seen-set only allocates on first-seen paths rather than every file — or
build the keep-set against `&files` with a borrowed `HashSet<&str>` and stay
allocation-free. Minor, but it's a regression on the append hot path.
##########
crates/iceberg/src/transaction/snapshot.rs:
##########
@@ -396,6 +418,19 @@ impl<'a> SnapshotProducer<'a> {
);
}
+ for data_file in &self.deleted_data_files {
Review Comment:
These `deleted-*` and `total-*` numbers come from the `DataFile`s the caller
passed, but `existing_manifest` matches only on `file_path` — so a caller with
the right path but a wrong `record_count`/size/partition writes a snapshot
whose summary and totals disagree with its own manifests, and the new
`update_totals` guard mostly just hides the negative case. Java builds the
removed-file summary from the matched `ManifestEntry` (the table's own copy),
not caller input. The rewrite already loads the real entry per match — I'd
thread those into `summary()` and compute from them, or at least validate the
caller's fields against the entry.
##########
crates/iceberg/src/transaction/overwrite.rs:
##########
@@ -0,0 +1,534 @@
+// 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, invalid_data};
+use crate::spec::{DataFile, ManifestEntry, ManifestFile, Operation};
+use crate::table::Table;
+use crate::transaction::snapshot::{
+ DefaultManifestProcess, SnapshotProduceOperation, SnapshotProducer,
+};
+use crate::transaction::{ActionCommit, TransactionAction};
+
+/// 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 operation = OverwriteOperation {
+ has_added_data_files: !self.added_data_files.is_empty(),
+ deleted_file_paths: self
+ .deleted_data_files
+ .iter()
+ .map(|f| f.file_path.clone())
+ .collect(),
+ };
+ snapshot_producer
+ .commit(operation, DefaultManifestProcess)
+ .await
+ }
+}
+
+struct OverwriteOperation {
+ has_added_data_files: bool,
+ deleted_file_paths: HashSet<String>,
+}
+
+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,
+ // Also a properties-only commit, the workaround from #1548.
+ _ => Operation::Append,
+ }
+ }
+
+ // Only the listed files are replaced, so the table totals carry over.
+ fn truncate_full_table(&self) -> bool {
+ false
+ }
+
+ async fn delete_entries(
+ &self,
+ _snapshot_produce: &SnapshotProducer<'_>,
+ ) -> Result<Vec<ManifestEntry>> {
+ Ok(vec![])
+ }
+
+ async fn existing_manifest(
+ &self,
+ snapshot_produce: &SnapshotProducer<'_>,
+ ) -> Result<Vec<ManifestFile>> {
+ let table = snapshot_produce.table;
+ let mut manifests = vec![];
+ let mut matched: HashSet<&str> = HashSet::new();
+
+ if let Some(snapshot) = table.metadata().current_snapshot() {
+ let manifest_list =
table.manifest_list_reader(snapshot).load().await?;
+ for manifest_file in manifest_list.entries() {
Review Comment:
`delete_data_files` takes any `DataFile` and nothing checks the content
type, and this loop walks every manifest including the `Deletes` ones — so a
path that matches a position/equality delete or a DV gets tombstoned in a
delete manifest, silently un-applying the deletes it was carrying. Java keeps
the two apart (`deleteFile(DataFile)` vs `deleteFile(DeleteFile)`, separate
filter managers), so I'd reject non-`Data` content in `delete_data_files` and
skip `ManifestContentType::Deletes` manifests here.
##########
crates/iceberg/src/transaction/overwrite.rs:
##########
@@ -0,0 +1,534 @@
+// 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, invalid_data};
+use crate::spec::{DataFile, ManifestEntry, ManifestFile, Operation};
+use crate::table::Table;
+use crate::transaction::snapshot::{
+ DefaultManifestProcess, SnapshotProduceOperation, SnapshotProducer,
+};
+use crate::transaction::{ActionCommit, TransactionAction};
+
+/// 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 operation = OverwriteOperation {
+ has_added_data_files: !self.added_data_files.is_empty(),
+ deleted_file_paths: self
+ .deleted_data_files
+ .iter()
+ .map(|f| f.file_path.clone())
+ .collect(),
+ };
+ snapshot_producer
+ .commit(operation, DefaultManifestProcess)
+ .await
+ }
+}
+
+struct OverwriteOperation {
+ has_added_data_files: bool,
+ deleted_file_paths: HashSet<String>,
+}
+
+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,
+ // Also a properties-only commit, the workaround from #1548.
+ _ => Operation::Append,
+ }
+ }
+
+ // Only the listed files are replaced, so the table totals carry over.
+ fn truncate_full_table(&self) -> bool {
+ false
+ }
+
+ async fn delete_entries(
+ &self,
+ _snapshot_produce: &SnapshotProducer<'_>,
+ ) -> Result<Vec<ManifestEntry>> {
+ Ok(vec![])
+ }
+
+ async fn existing_manifest(
+ &self,
+ snapshot_produce: &SnapshotProducer<'_>,
+ ) -> Result<Vec<ManifestFile>> {
+ let table = snapshot_produce.table;
+ let mut manifests = vec![];
+ let mut matched: HashSet<&str> = HashSet::new();
+
+ if let Some(snapshot) = table.metadata().current_snapshot() {
+ let manifest_list =
table.manifest_list_reader(snapshot).load().await?;
+ for manifest_file in manifest_list.entries() {
+ // Delete-only manifests record which files were removed and
must survive
+ // until `expire_snapshots` cleans them up (see #2148).
+ if !manifest_file.has_added_files()
+ && !manifest_file.has_existing_files()
+ && !manifest_file.has_deleted_files()
+ {
+ continue;
+ }
+ if self.deleted_file_paths.is_empty() {
+ manifests.push(manifest_file.clone());
+ continue;
+ }
+
+ let manifest =
table.manifest_reader().read(manifest_file).await?;
+ let deletes: Vec<&str> = manifest
+ .entries()
+ .iter()
+ .filter(|entry| entry.is_alive())
+ .filter_map(|entry|
self.deleted_file_paths.get(entry.file_path()))
+ .map(String::as_str)
+ .collect();
+ if deletes.is_empty() {
+ manifests.push(manifest_file.clone());
+ continue;
+ }
+ matched.extend(deletes);
+
+ let mut writer = snapshot_produce.new_manifest_writer(
+ manifest_file.content,
+ manifest.metadata().schema.clone(),
+ Arc::new(manifest.metadata().partition_spec.clone()),
+ )?;
+ // Like Java, only live entries are carried into the rewrite.
+ for entry in manifest.entries().iter().filter(|entry|
entry.is_alive()) {
Review Comment:
This is the `first_row_id` path I flagged last round, and I think my call
was wrong — worth saying directly. The per-entry `data_file.first_row_id`
survives the rewrite: the reader materializes it for live entries,
`add_existing_entry` carries it through, and `_serde` writes it back, so
surviving files keep their explicit ranges. The fresh manifest-level
`first_row_id` only fills null ids, which a rewrite of existing files never
has. So I'm withdrawing the corruption concern — it's spec- and Java-consistent.
What's still open from last round is a test. The whole thing rests on that
reader→serde chain and nothing in the suite asserts row lineage, so I'd want
one that appends A/B/C on a V3 table, overwrite-deletes B, then checks A and C
keep their original `first_row_id` after re-read and `next_row_id` advanced as
expected.
--
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]