laskoviymishka commented on code in PR #3253:
URL: https://github.com/apache/iceberg-rust/pull/3253#discussion_r4153177113
##########
crates/iceberg/src/transaction/snapshot.rs:
##########
@@ -82,9 +83,24 @@ pub(crate) trait SnapshotProduceOperation: Send + Sync {
/// - **Overwrite operations**: May exclude manifests for partitions being
overwritten
/// - **Delete operations**: May exclude manifests for partitions being
deleted
fn existing_manifest(
- &self,
+ &mut self,
snapshot_produce: &SnapshotProducer<'_>,
) -> impl Future<Output = Result<Vec<ManifestFile>>> + Send;
+
+ /// Returns the data files this operation actually removed, each with the
schema and
+ /// partition spec of the manifest that recorded it.
+ ///
+ /// Only populated once [`Self::existing_manifest`] has run, so the
snapshot summary counts
+ /// what was really removed rather than what the caller asked to remove.
+ fn removed_data_files(&self) -> &[(DataFile, SchemaRef, PartitionSpecRef)]
{
+ &[]
+ }
+
+ /// Returns whether this operation replaces the whole table, dropping the
previous totals
+ /// from the snapshot summary. A partial overwrite must return `false`.
+ fn truncate_full_table(&self) -> bool {
Review Comment:
Flipping this from `operation() == Operation::Overwrite` to an opt-in flag
is the right move for partial overwrites — this PR's own
`RemoveExistingFileOperation` depends on the `false` default to report actual
removed files instead of truncating. My worry is what it does to the ops that
come after it: a full-table overwrite (replace-partitions, overwrite-by-filter
with no predicate) returns `Overwrite` too, and if it forgets to override this
it silently accumulates totals instead of resetting them — no error, and no
repair path once the baseline is wrong.
Since this is step 1 and the follow-ons are the exact callers at risk, I'd
bind the invariant to something harder than a doc comment: at minimum a
`debug_assert!`/warn in `summary()` when `operation() == Overwrite &&
!truncate_full_table()`, or better, tie truncation to a dedicated
operation/path so the type carries it. I'd also add one e2e test that drives
`truncate_full_table() == true` through `commit()` — right now only the
internal `truncate_table_summary` fn is covered, so the wiring here is untested.
##########
crates/iceberg/src/transaction/snapshot.rs:
##########
@@ -82,9 +83,24 @@ pub(crate) trait SnapshotProduceOperation: Send + Sync {
/// - **Overwrite operations**: May exclude manifests for partitions being
overwritten
/// - **Delete operations**: May exclude manifests for partitions being
deleted
fn existing_manifest(
- &self,
+ &mut self,
snapshot_produce: &SnapshotProducer<'_>,
) -> impl Future<Output = Result<Vec<ManifestFile>>> + Send;
+
+ /// Returns the data files this operation actually removed, each with the
schema and
+ /// partition spec of the manifest that recorded it.
+ ///
+ /// Only populated once [`Self::existing_manifest`] has run, so the
snapshot summary counts
+ /// what was really removed rather than what the caller asked to remove.
+ fn removed_data_files(&self) -> &[(DataFile, SchemaRef, PartitionSpecRef)]
{
Review Comment:
`&[(DataFile, SchemaRef, PartitionSpecRef)]` is a positional three-tuple
whose roles you have to infer, and the COW series is going to stack callers on
top of it. A small named struct (`RemovedFile { data_file, schema,
partition_spec }`) costs nothing now and saves a coordinated touch at every
call site the first time this needs a fourth field.
##########
crates/iceberg/src/transaction/snapshot.rs:
##########
@@ -403,6 +433,10 @@ impl<'a> SnapshotProducer<'a> {
);
}
+ for (data_file, schema, partition_spec) in
snapshot_produce_operation.removed_data_files() {
Review Comment:
This loop builds the summary's `deleted-*` counts purely from
`removed_data_files()`, while the manifest's `Deleted` entries come from
whatever `existing_manifest()` passed to `add_delete_entry()`. Nothing binds
the two, so an operation that marks entries deleted in the manifest without
reporting them here (or the reverse) commits manifests and a summary that
disagree — and `total-data-files` drifts from there on every subsequent commit.
`test_summary_counts_files_the_operation_reports_as_removed` is the case in
point: `RemoveOneFileOperation` reports a removal but writes no `Deleted`
entry, so the summary claims `deleted-data-files=1` over manifests with zero
deleted entries, and it passes. I'd either sum `deleted_files_count` across the
produced data manifests and `debug_assert!` it matches
`removed_data_files().len()` in `produce_manifests`, or have
`existing_manifest` return the removed set so the two can't diverge by
construction.
##########
crates/iceberg/src/transaction/snapshot.rs:
##########
@@ -308,7 +328,9 @@ impl<'a> SnapshotProducer<'a> {
// Write manifest file for added data files and return the ManifestFile
for ManifestList.
async fn write_added_manifest(&mut self) -> Result<ManifestFile> {
- let added_data_files = std::mem::take(&mut self.added_data_files);
+ // Cloned rather than taken: the summary is built after the manifests,
so the added
+ // files must still be here.
+ let added_data_files = self.added_data_files.clone();
Review Comment:
The ordering fix is right, but cloning the whole `added_data_files` vec here
is a real cost on the commit hot path — each `DataFile` carries several
`HashMap`s, so a large append allocates a full duplicate where `mem::take` was
free. Since the only reason to keep them is so `summary()` can read them
afterward, I'd pass `&[DataFile]` into `write_added_manifest` (or iterate
borrowing) instead of cloning the vec.
##########
crates/iceberg/src/transaction/snapshot.rs:
##########
@@ -234,22 +250,26 @@ impl<'a> SnapshotProducer<'a> {
snapshot_id
}
- fn new_manifest_writer(&mut self, content: ManifestContentType) ->
Result<ManifestWriter> {
- let new_manifest_path = format!(
+ /// Returns the path for the next manifest file of this commit.
+ fn new_manifest_path(&self) -> Result<String> {
+ Ok(format!(
"{}/{}-m{}.{}",
self.table.metadata().metadata_location()?,
self.commit_uuid,
- self.manifest_counter.next().unwrap(),
+ self.manifest_counter.fetch_add(1, Ordering::Relaxed),
DataFileFormat::Avro
- );
- let output_file = self.table.file_io().new_output(new_manifest_path)?;
- let partition_spec = self
- .table
- .metadata()
- .default_partition_spec()
- .as_ref()
- .clone();
- let schema = self.table.metadata().current_schema().clone();
+ ))
+ }
+
+ /// Returns a writer for the next manifest file of this commit, routed
through the table's
+ /// encryption manager when one is configured.
+ pub(crate) fn new_manifest_writer(
+ &self,
+ content: ManifestContentType,
+ schema: SchemaRef,
+ partition_spec: PartitionSpec,
Review Comment:
`schema` comes in as `SchemaRef` but `partition_spec` is taken by value, so
every caller holding an `Arc<PartitionSpec>` has to `.as_ref().clone()` a full
spec to hand it over — `write_added_manifest` does exactly that just below, and
`removed_data_files()` already standardizes on `PartitionSpecRef`. I'd take
`partition_spec: PartitionSpecRef` here and clone once internally if the
builder needs an owned value, so the call sites stop deep-cloning and the Arc
convention stays consistent.
##########
crates/iceberg/src/spec/snapshot_summary.rs:
##########
@@ -532,7 +532,7 @@ fn update_totals(
return;
};
- let new_total = previous_total + added - removed;
+ let new_total = (previous_total + added).saturating_sub(removed);
Review Comment:
The `saturating_sub` fixes the subtraction underflow, but `previous_total +
added` on the left is still plain `u64` addition — a corrupt or adversarial
previous summary with `total-data-files` near `u64::MAX` panics in debug and
wraps to a small number in release. `saturating_add` closes it: `let new_total
= previous_total.saturating_add(added).saturating_sub(removed);`.
Separately, saturating to `0` diverges from Java:
`SnapshotProducer.updateTotal` *omits* the `total-*` property entirely when the
value would go negative, and skips recomputation when the key is absent — so
writing `0` makes a Java reader treat it as a real basis and carry the
undercount forward. I'd match Java and omit the property on underflow, and flip
`test_update_totals_saturates_...` to assert the key is absent rather than
`"0"`.
##########
crates/iceberg/src/transaction/snapshot.rs:
##########
@@ -116,7 +132,7 @@ pub(crate) struct SnapshotProducer<'a> {
// A counter used to generate unique manifest file names.
// It starts from 0 and increments for each new manifest file.
// Note: This counter is limited to the range of (0..u64::MAX).
- manifest_counter: RangeFrom<u64>,
+ manifest_counter: AtomicU64,
Review Comment:
The carried-over comment just above says the counter is "limited to the
range of (0..u64::MAX)", but `fetch_add` wraps at `u64::MAX` rather than
stopping — so the note now reads as a guarantee the type doesn't give (a wrap
hands out a duplicate manifest path). While we're here, a word on why it's
`AtomicU64` would help: it's interior mutability so `new_manifest_writer` can
take `&self`, not cross-thread sharing — which is why `Relaxed` is fine.
--
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]