zhuxiangyi commented on code in PR #965:
URL: https://github.com/apache/paimon-rust/pull/965#discussion_r4173318672


##########
crates/paimon/src/table/expire_snapshots.rs:
##########
@@ -0,0 +1,473 @@
+// 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.
+
+//! Snapshot expiration.
+//!
+//! Reference: Java `ExpireSnapshotsImpl`, configured like the Flink and Spark
+//! `expire_snapshots` procedures (`ProcedureUtils.fillInSnapshotOptions`).
+//!
+//! Expiring snapshots `[earliest, end)` removes, in order:
+//! 1. data files deleted by the delta manifests of `(earliest, end]` that no
+//!    tag still reads;
+//! 2. changelog files added by `[earliest, end)`;
+//! 3. manifest lists, manifests, index manifests, index files, statistics and
+//!    reassign plans of `[earliest, end)` that neither `end` nor a tag in the
+//!    range still references;
+//! 4. the snapshot files themselves, then the EARLIEST hint moves to `end`.
+//!
+//! Snapshot files go last, so an interrupted run leaves every remaining
+//! snapshot readable and a later run finishes the job.
+
+use crate::spec::{CoreOptions, Snapshot};
+use crate::table::snapshot_deletion::{read_long_lived_changelogs, DataFileKey, 
SnapshotDeletion};
+use crate::table::{BranchManager, Table};
+use crate::{Error, Result};
+use futures::{stream, StreamExt, TryStreamExt};
+use std::collections::hash_map::Entry;
+use std::collections::{HashMap, HashSet};
+use std::time::{SystemTime, UNIX_EPOCH};
+
+const SNAPSHOT_READ_CONCURRENCY: usize = 16;
+
+/// Retention rules for one expiration run. Java `ExpireConfig` (snapshot 
part).
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+struct ExpireConfig {
+    retain_max: i32,
+    retain_min: i32,
+    /// Snapshots whose successor was committed at or after this time are kept.
+    older_than_millis: i64,
+    max_deletes: i32,
+}
+
+/// Expire old snapshots of a table and delete the files only they reference.
+///
+/// Defaults come from the table options `snapshot.num-retained.min`,
+/// `snapshot.num-retained.max`, `snapshot.time-retained` and
+/// `snapshot.expire.limit`; each can be overridden per run, like the arguments
+/// of Java's `expire_snapshots` procedure.
+pub struct ExpireSnapshots<'a> {
+    table: &'a Table,
+    retain_max: Option<i32>,
+    retain_min: Option<i32>,
+    older_than_millis: Option<i64>,
+    max_deletes: Option<i32>,
+    current_time_millis: Option<i64>,
+}
+
+impl<'a> ExpireSnapshots<'a> {
+    pub(crate) fn new(table: &'a Table) -> Self {
+        Self {
+            table,
+            retain_max: None,
+            retain_min: None,
+            older_than_millis: None,
+            max_deletes: None,
+            current_time_millis: None,
+        }
+    }
+
+    /// Maximum number of completed snapshots to retain.
+    pub fn with_retain_max(&mut self, retain_max: i32) -> &mut Self {
+        self.retain_max = Some(retain_max);
+        self
+    }
+
+    /// Minimum number of completed snapshots to retain.
+    pub fn with_retain_min(&mut self, retain_min: i32) -> &mut Self {
+        self.retain_min = Some(retain_min);
+        self
+    }
+
+    /// Expire snapshots older than this epoch-millisecond timestamp instead of
+    /// `snapshot.time-retained`.
+    pub fn with_older_than_millis(&mut self, older_than_millis: i64) -> &mut 
Self {
+        self.older_than_millis = Some(older_than_millis);
+        self
+    }
+
+    /// Maximum number of snapshots to expire in this run.
+    pub fn with_max_deletes(&mut self, max_deletes: i32) -> &mut Self {
+        self.max_deletes = Some(max_deletes);
+        self
+    }
+
+    #[cfg(test)]
+    pub(crate) fn with_current_time_millis(&mut self, current_time_millis: 
i64) -> &mut Self {
+        self.current_time_millis = Some(current_time_millis);
+        self
+    }
+
+    fn config(&self) -> Result<ExpireConfig> {
+        let core = CoreOptions::new(self.table.schema().options());
+        let retain_max = match self.retain_max {
+            Some(value) => positive("retain_max", value)?,
+            None => core.snapshot_num_retained_max()?,
+        };
+        let retain_min = match self.retain_min {
+            Some(value) => positive("retain_min", value)?,
+            None => core.snapshot_num_retained_min()?,
+        };
+        if retain_max < retain_min {
+            return Err(Error::DataInvalid {
+                message: format!(
+                    "retainMax ({retain_max}) must not be less than retainMin 
({retain_min})."
+                ),
+                source: None,
+            });
+        }
+        let max_deletes = match self.max_deletes {
+            Some(value) => positive("max_deletes", value)?,
+            None => core.snapshot_expire_limit()?,
+        };
+        let older_than_millis = match self.older_than_millis {
+            Some(value) => value,
+            None => {
+                let now = 
self.current_time_millis.unwrap_or_else(current_time_millis);
+                let retained = 
i64::try_from(core.snapshot_time_retained_ms()?).unwrap_or(i64::MAX);
+                now.saturating_sub(retained)
+            }
+        };
+        Ok(ExpireConfig {
+            retain_max,
+            retain_min,
+            older_than_millis,
+            max_deletes,
+        })
+    }
+
+    /// Expire snapshots and return how many snapshot files were removed.
+    pub async fn execute(&self) -> Result<usize> {
+        self.table.ensure_not_branch_reference_for_write()?;
+        let config = self.config()?;
+
+        let snapshot_manager = self.table.snapshot_manager();
+        let Some(latest) = snapshot_manager.get_latest_snapshot().await? else {
+            return Ok(0);
+        };
+        let latest = latest.id();
+        let Some(earliest) = snapshot_manager.earliest_snapshot_id().await? 
else {
+            return Ok(0);
+        };
+
+        // The oldest snapshot `snapshot.num-retained.max` lets us keep.
+        let min = (latest - i64::from(config.retain_max) + 1).max(earliest);
+        // `snapshot.num-retained.min` protects the newest snapshots.
+        let mut max_exclusive = latest - i64::from(config.retain_min) + 1;
+        // A snapshot a consumer still reads from cannot be deleted.
+        let consumers = self.table.consumer_manager().list_all().await?;
+        if let Some(min_next) = consumers.iter().map(|(_, next)| *next).min() {
+            max_exclusive = max_exclusive.min(min_next);
+        }
+        // `snapshot.expire.limit` bounds one run.
+        max_exclusive = 
max_exclusive.min(earliest.saturating_add(i64::from(config.max_deletes)));
+
+        for id in min..max_exclusive {
+            // A snapshot expires only once its successor has also outlived the
+            // retention time.
+            if let Some(next) = snapshot_manager.try_get_snapshot(id + 
1).await? {
+                if config.older_than_millis <= next.time_millis() as i64 {
+                    return self.expire_until(earliest, id).await;
+                }
+            }
+        }
+        self.expire_until(earliest, max_exclusive).await
+    }
+
+    async fn expire_until(&self, earliest: i64, end_exclusive: i64) -> 
Result<usize> {
+        let snapshot_manager = self.table.snapshot_manager();
+        if end_exclusive <= earliest {
+            // Nothing expires; record the earliest snapshot so later runs find
+            // it without listing the directory.
+            if !snapshot_manager.earliest_hint_exists().await? {
+                snapshot_manager.write_earliest_hint(earliest).await?;
+            }
+            return Ok(0);
+        }
+
+        // Futures are built eagerly: a borrowing closure inside the stream
+        // would make callers' futures lose `Send`.
+        let reads = (earliest..=end_exclusive)
+            .map(|id| snapshot_manager.try_get_snapshot(id))
+            .collect::<Vec<_>>();
+        let snapshots_including_end = stream::iter(reads)
+            .buffered(SNAPSHOT_READ_CONCURRENCY)
+            .try_collect::<Vec<_>>()
+            .await?
+            .into_iter()
+            .flatten()
+            .collect::<Vec<_>>();
+        let Some(first) = snapshots_including_end.first() else {
+            return Ok(0);
+        };
+        let begin_inclusive = first.id();
+        let snapshots_excluding_end = snapshots_including_end
+            .iter()
+            .filter(|snapshot| snapshot.id() != end_exclusive)
+            .collect::<Vec<_>>();
+
+        let tagged = self.tagged_snapshots().await?;
+        let deletion = SnapshotDeletion::new(self.table)?;
+        // Branches and long-lived changelogs share files with this history.
+        // Read everything they reference before deleting anything; if that is
+        // impossible, the run stops with nothing changed.
+        let external = self.external_owners(&deletion).await?;
+
+        // Data files deleted by a snapshot are unused from that snapshot on,
+        // so the range is (begin, end].
+        let data_snapshots = snapshots_including_end
+            .iter()
+            .filter(|snapshot| snapshot.id() != begin_inclusive)
+            .collect::<Vec<_>>();
+        self.clean_data_files(&deletion, &data_snapshots, &tagged, 
&external.data_files)
+            .await;
+
+        let mut changelog_files = Vec::new();
+        for snapshot in &snapshots_excluding_end {
+            let still_read = snapshot
+                .changelog_manifest_list()
+                .is_some_and(|list| external.changelog_lists.contains(list));
+            if !still_read {
+                
changelog_files.extend(deletion.changelog_files(snapshot).await);
+            }
+        }
+        deletion.delete_quietly(changelog_files).await;
+
+        let last = snapshots_including_end
+            .last()
+            .expect("checked non-empty above");
+        if last.id() != end_exclusive {
+            // The end snapshot is gone, so nothing protects the manifests it
+            // would share; stop before deleting any of them.
+            return Ok(0);
+        }
+        let mut skipping_snapshots = find_skipping_tags(&tagged, 
begin_inclusive, end_exclusive);
+        skipping_snapshots.push(last);
+        let external_snapshots = external.snapshots.iter().collect::<Vec<_>>();
+        let skipping = match deletion
+            .manifest_skipping_set(&skipping_snapshots, false)
+            .await
+        {
+            Ok(mut skipping) => {
+                skipping.extend(
+                    deletion
+                        .manifest_skipping_set(&external_snapshots, true)

Review Comment:
   Fixed in f7ba2d3. The skipping set of the external owners reads their index 
and changelog manifests. It is now built right after the owners are read, 
before any data or changelog file is deleted. A run that cannot read one of 
them now fails with nothing changed.
   
   - `test_unreadable_branch_index_manifest_changes_nothing` follows your 
probe: build a global index, branch from the indexed snapshot, drop main's 
index, overwrite twice, then corrupt only the branch-owned index manifest. The 
call fails, and every snapshot and data file remains.
   - `test_unreadable_external_changelog_manifest_changes_nothing` covers an 
unreadable changelog manifest of a long-lived changelog.
   
   Both tests fail if the set is built after the deletions again.
   



##########
crates/paimon/src/table/snapshot_deletion.rs:
##########
@@ -0,0 +1,500 @@
+// 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.
+
+//! Planning and deleting the files of expired snapshots.
+//!
+//! Reference: Java `FileDeletionBase` and `SnapshotDeletion`.
+//!
+//! Every read failure makes the plan delete less, never more: an unreadable
+//! delta manifest cancels that snapshot's data-file deletion, and a skipping
+//! set that cannot be built cancels manifest deletion. A file left behind is
+//! an orphan that orphan-file cleanup can remove later; a referenced file
+//! deleted by mistake is data loss.
+
+use crate::io::FileIO;
+use crate::spec::{
+    bucket_path, BinaryRow, CoreOptions, FileKind, IndexManifest, Manifest, 
ManifestEntry,
+    ManifestFileMeta, ManifestList, PartitionComputer, Snapshot,
+};
+use crate::table::index_file_path::{
+    committed_index_file_path, resolve_legacy_deletion_vector_entries,
+};
+use crate::table::{SnapshotManager, Table};
+use crate::Result;
+use futures::{stream, StreamExt, TryStreamExt};
+use std::collections::hash_map::Entry;
+use std::collections::{HashMap, HashSet};
+
+const FILE_OPERATION_CONCURRENCY: usize = 16;
+const STATISTICS_DIR: &str = "statistics";
+/// Java `SerializationAssignment.PLAN_FILE_PROPERTY`.
+const REASSIGN_PLAN_FILE_PROPERTY: &str = "row-id-reassign.plan";
+/// Java `SerializationAssignment.REASSIGN_SNAPSHOT_ID`.
+const REASSIGN_SNAPSHOT_ID_PROPERTY: &str = "reassign-snapshot-id";
+
+/// A data file within a bucket: `(partition, bucket, file name)`. File names
+/// are unique within a bucket, so this is how Java's `containsDataFile`
+/// matches entries across snapshots.
+pub(crate) type DataFileKey = (Vec<u8>, i32, String);
+
+/// Data files that a snapshot's delta manifests delete, with the physical
+/// paths (data file plus extra files) to remove for each.
+#[derive(Debug, Default)]
+pub(crate) struct DataFileDeletionPlan {
+    candidates: Vec<(DataFileKey, Vec<String>)>,
+}
+
+impl DataFileDeletionPlan {
+    /// Paths of candidates that none of `retained` protects.
+    pub(crate) fn paths_to_delete(&self, retained: &[&HashSet<DataFileKey>]) 
-> Vec<String> {
+        self.candidates
+            .iter()
+            .filter(|(key, _)| !retained.iter().any(|retained| 
retained.contains(key)))
+            .flat_map(|(_, paths)| paths.iter().cloned())
+            .collect()
+    }
+}
+
+pub(crate) struct SnapshotDeletion {
+    table: Table,
+    file_io: FileIO,
+    table_location: String,
+    snapshot_manager: SnapshotManager,
+    partition_computer: Option<PartitionComputer>,
+    index_file_in_data_file_dir: bool,
+}
+
+impl SnapshotDeletion {
+    pub(crate) fn new(table: &Table) -> Result<Self> {
+        let schema = table.schema();
+        let core_options = CoreOptions::new(schema.options());
+        let partition_computer = if schema.partition_keys().is_empty() {
+            None
+        } else {
+            Some(PartitionComputer::new(
+                schema.partition_keys(),
+                schema.fields(),
+                core_options.partition_default_name(),
+                core_options.legacy_partition_name(),
+            )?)
+        };
+        Ok(Self {
+            table: table.clone(),
+            file_io: table.file_io().clone(),
+            table_location: table.location().trim_end_matches('/').to_string(),
+            snapshot_manager: table.snapshot_manager(),
+            partition_computer,
+            index_file_in_data_file_dir: 
core_options.index_file_in_data_file_dir(),
+        })
+    }
+
+    fn bucket_path(&self, partition: &[u8], bucket: i32) -> Result<String> {
+        let partition = match self.partition_computer {
+            // An unpartitioned table never decodes its (empty) partition blob.
+            None => BinaryRow::new(0),
+            Some(_) => BinaryRow::from_serialized_bytes(partition)?,
+        };
+        bucket_path(
+            &self.table_location,

Review Comment:
   Fixed in f7ba2d3. Bucket directories now come from 
`Table::data_file_location()`, as for writers, readers and commit cleanup. That 
covers data files, their sidecars, and bucket-local index files. Manifests, 
statistics and global index files stay under the table location.
   
   Covered in the core tests for a relative (`data`) and an absolute 
`data-file.path-directory`, and end to end through `CALL sys.expire_snapshots` 
in `test_expire_snapshots_procedure_with_data_directory`. Each test fails if 
the table location is used again.
   



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

Reply via email to