JanKaul commented on code in PR #3046: URL: https://github.com/apache/iceberg-rust/pull/3046#discussion_r4192688698
########## crates/iceberg/src/transaction/rewrite.rs: ########## @@ -0,0 +1,1353 @@ +// 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. + +//! Transaction action for rewriting data files (compaction). +//! +//! [`RewriteFilesAction`] replaces a set of data files with a new set while +//! keeping the logical table contents unchanged. This is used for compaction — +//! merging many small files into fewer large ones. +//! +//! The resulting snapshot uses [`Operation::Replace`] to indicate that files +//! were reorganised without changing the data. + +use std::sync::Arc; + +use async_trait::async_trait; + +use crate::error::{Result, invalid_data}; +use crate::spec::{DataFile, ManifestContentType, Operation}; +use crate::table::Table; +use crate::transaction::merging::MergingSnapshotProducer; +use crate::transaction::{ActionCommit, TransactionAction}; + +/// A transaction action that rewrites (replaces) data files. +/// +/// This is the Rust equivalent of Java's `BaseRewriteFiles`. It uses +/// `MergingSnapshotProducer` to handle manifest filtering and creation, +/// and commits a snapshot with [`Operation::Replace`]. +/// +/// # Example +/// +/// ```ignore +/// let tx = Transaction::new(&table); +/// let action = tx.rewrite_files() +/// .delete_file(old_file_1) +/// .delete_file(old_file_2) +/// .add_file(merged_file); +/// let tx = action.apply(tx)?; +/// let table = tx.commit(&catalog).await?; +/// ``` +/// +/// # Concurrent deletes +/// +/// The action records the table's current snapshot when it is created, and +/// refuses to commit if the table gained a delete manifest after it. Java +/// narrows that to the files being replaced in +/// `validateNoNewDeletesForDataFiles`; until that exists here the check is +/// table-wide, so a delete committed to an unrelated partition while the +/// rewrite ran also fails the commit with +/// [`ErrorKind::DataInvalid`](crate::ErrorKind::DataInvalid). So do a +/// rewrite planned against a table that had no snapshot yet, and one whose +/// starting snapshot has since been expired: neither can rule out a delete. +/// Replan the rewrite against the current table and run it again. +/// +/// Deletes that were already committed when the rewrite was planned are not +/// covered by this check. Either apply them while rewriting, or keep them +/// applicable to the new file with +/// [`data_sequence_number`](RewriteFilesAction::data_sequence_number). +pub struct RewriteFilesAction { + producer: MergingSnapshotProducer, + /// The snapshot the rewrite was planned against, if the table had one. + /// Deletes newer than it fail the commit. + starting_snapshot_id: Option<i64>, +} + +impl RewriteFilesAction { + pub(crate) fn new(starting_snapshot_id: Option<i64>) -> Self { + Self { + producer: MergingSnapshotProducer::new(Operation::Replace), + starting_snapshot_id, + } + } + + /// Register a data file to be removed from the table. + /// + /// The file must exist in the current snapshot; otherwise the commit + /// will fail with a validation error. + pub fn delete_file(mut self, file: DataFile) -> Self { + self.producer.delete_data_file(file); + self + } + + /// Register a data file to be added to the table. + /// + /// Typically this is the merged output of the files being deleted. + pub fn add_file(mut self, file: DataFile) -> Self { + self.producer.add_data_file(file); + self + } + + /// Set the data sequence number recorded for every added file. + /// + /// Without this the added files inherit the sequence number of the new + /// snapshot. A compaction that must not shadow concurrently written + /// deletes sets the sequence number of the files it replaces instead. + /// V1 manifest entries carry no sequence number, so this has no effect on + /// a V1 table. + pub fn data_sequence_number(mut self, sequence_number: i64) -> Self { + self.producer.set_data_sequence_number(sequence_number); + self + } + + fn validate(&self) -> Result<()> { + if !self.producer.has_deleted_data_files() { + return Err(invalid_data!( + "Rewrite files requires at least one file to delete" + )); + } + if !self.producer.has_added_data_files() { + return Err(invalid_data!( + "Rewrite files requires at least one file to add" + )); + } Review Comment: Should this follow `BaseRewriteFiles` and allow a rewrite with no added files? Java's `validateReplacedAndAddedFiles` only requires the removal side to be non-empty ("Files to delete cannot be empty"); its other two checks forbid adds without matching removals, not the reverse. Remove-only rewrites are legitimate even at this PR's data-file-only scope: a compaction group whose remaining rows are all shadowed by deletes produces zero output files, and this check would reject that commit. It also matters for the follow-ups — `RemoveDanglingDeleteFiles` is exactly a remove-only `RewriteFiles`, and rewritten position-delete groups can shrink to nothing once dangling rows are dropped. Concretely: drop this check and keep only the `has_deleted_data_files()` requirement above (which then generalizes to "deletes data or delete files" once delete-file removal lands). Happy to send that as part of the delete-file follow-up instead if you'd rather not touch it here — but relaxing a validation after release is noisier than merging the Java-parity rule now. -- 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]
