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


##########
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 the shared expiration code of #965 (f7ba2d3). The external owners' 
skipping set, including their index and changelog manifests, is now built 
before any deletion, so an unreadable owner fails the run with nothing changed. 
The regression tests are described in the #965 thread.
   



##########
crates/paimon/src/table/orphan_files_clean.rs:
##########
@@ -0,0 +1,531 @@
+// 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.
+
+//! Orphan file cleanup.
+//!
+//! Reference: Java `OrphanFilesClean` and `LocalOrphanFilesClean`.
+//!
+//! A file is an orphan when it sits in a directory Paimon writes to (the
+//! manifest, index and statistics directories, bucket directories and data
+//! file external paths), is older than `older_than`, and no snapshot, tag or
+//! long-lived changelog of any branch references it. Temporary files left in
+//! snapshot and changelog directories by interrupted commits are removed too.
+//!
+//! `older_than` defaults to one day ago and must be in the past, so files that
+//! an in-flight write has created but not yet committed are never candidates.
+//! Unlike Java, which reads a missing manifest as empty, a missing manifest of
+//! a snapshot that still exists aborts the run: its data files would otherwise
+//! look unreferenced and be deleted.
+
+use crate::io::{FileIO, FileStatus};
+use crate::spec::{IndexManifest, Manifest, ManifestList, Snapshot};
+use crate::table::snapshot_deletion::reassign_plan_file;
+use crate::table::{BranchManager, SnapshotManager, Table};
+use crate::{Error, Result};
+use futures::{stream, StreamExt, TryStreamExt};
+use std::collections::{BTreeMap, HashSet};
+use std::time::{SystemTime, UNIX_EPOCH};
+
+const DEFAULT_OLDER_THAN_MS: i64 = 24 * 60 * 60 * 1000;
+const FILE_OPERATION_CONCURRENCY: usize = 16;
+const MAIN_BRANCH: &str = "main";
+const SNAPSHOT_PREFIX: &str = "snapshot-";
+const CHANGELOG_PREFIX: &str = "changelog-";
+const EARLIEST: &str = "EARLIEST";
+const LATEST: &str = "LATEST";
+const BUCKET_PREFIX: &str = "bucket-";
+/// Java `ManagedBlobReferenceFile.MANAGED_BLOB_SUFFIX`: never cleaned here.
+const MANAGED_BLOB_SUFFIX: &str = ".managed.blob";
+const DATA_FILE_EXTERNAL_PATHS_OPTION: &str = "data-file.external-paths";
+
+/// Files removed (or, in a dry run, that would be removed) by a cleanup.
+#[derive(Debug, Clone, Default, PartialEq, Eq)]
+pub struct OrphanFilesCleanResult {
+    pub deleted_file_count: u64,
+    pub deleted_file_total_bytes: u64,
+    /// Full paths, sorted.
+    pub deleted_files: Vec<String>,
+}
+
+/// Remove files that no snapshot, tag or changelog of the table references.
+pub struct RemoveOrphanFiles<'a> {
+    table: &'a Table,
+    older_than_millis: Option<i64>,
+    dry_run: bool,
+    parallelism: usize,
+    current_time_millis: Option<i64>,
+}
+
+impl<'a> RemoveOrphanFiles<'a> {
+    pub(crate) fn new(table: &'a Table) -> Self {
+        Self {
+            table,
+            older_than_millis: None,
+            dry_run: false,
+            parallelism: FILE_OPERATION_CONCURRENCY,
+            current_time_millis: None,
+        }
+    }
+
+    /// Maximum number of concurrent manifest reads and file deletions.
+    pub fn with_parallelism(&mut self, parallelism: usize) -> &mut Self {
+        self.parallelism = parallelism.max(1);
+        self
+    }
+
+    /// Only files last modified before this epoch-millisecond timestamp are
+    /// candidates. Defaults to one day ago; must be in the past.
+    pub fn with_older_than_millis(&mut self, older_than_millis: i64) -> &mut 
Self {
+        self.older_than_millis = Some(older_than_millis);
+        self
+    }
+
+    /// Report the orphan files without deleting them.
+    pub fn with_dry_run(&mut self, dry_run: bool) -> &mut Self {
+        self.dry_run = dry_run;
+        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
+    }
+
+    pub async fn execute(&self) -> Result<OrphanFilesCleanResult> {
+        self.table.ensure_not_branch_reference_for_write()?;

Review Comment:
   Fixed in 2f6f156. Before listing anything, the cleanup now rejects every 
table type except `table` and `materialized-table`. Format, object, Lance and 
Iceberg tables keep live files that no snapshot references, and Java likewise 
accepts only `FileStoreTable`.
   
   - `test_format_table_is_rejected_and_its_files_kept` follows your probe: a 
Format Table's data file at `bucket-0/part-0.parquet` is older than the 
cut-off. The call returns `Unsupported`, and the file and its rows remain.
   - `test_remove_orphan_files_procedure_rejects_format_tables` checks the same 
end to end through SQL.
   



##########
crates/paimon/src/table/orphan_files_clean.rs:
##########
@@ -0,0 +1,531 @@
+// 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.
+
+//! Orphan file cleanup.
+//!
+//! Reference: Java `OrphanFilesClean` and `LocalOrphanFilesClean`.
+//!
+//! A file is an orphan when it sits in a directory Paimon writes to (the
+//! manifest, index and statistics directories, bucket directories and data
+//! file external paths), is older than `older_than`, and no snapshot, tag or
+//! long-lived changelog of any branch references it. Temporary files left in
+//! snapshot and changelog directories by interrupted commits are removed too.
+//!
+//! `older_than` defaults to one day ago and must be in the past, so files that
+//! an in-flight write has created but not yet committed are never candidates.
+//! Unlike Java, which reads a missing manifest as empty, a missing manifest of
+//! a snapshot that still exists aborts the run: its data files would otherwise
+//! look unreferenced and be deleted.
+
+use crate::io::{FileIO, FileStatus};
+use crate::spec::{IndexManifest, Manifest, ManifestList, Snapshot};
+use crate::table::snapshot_deletion::reassign_plan_file;
+use crate::table::{BranchManager, SnapshotManager, Table};
+use crate::{Error, Result};
+use futures::{stream, StreamExt, TryStreamExt};
+use std::collections::{BTreeMap, HashSet};
+use std::time::{SystemTime, UNIX_EPOCH};
+
+const DEFAULT_OLDER_THAN_MS: i64 = 24 * 60 * 60 * 1000;
+const FILE_OPERATION_CONCURRENCY: usize = 16;
+const MAIN_BRANCH: &str = "main";
+const SNAPSHOT_PREFIX: &str = "snapshot-";
+const CHANGELOG_PREFIX: &str = "changelog-";
+const EARLIEST: &str = "EARLIEST";
+const LATEST: &str = "LATEST";
+const BUCKET_PREFIX: &str = "bucket-";
+/// Java `ManagedBlobReferenceFile.MANAGED_BLOB_SUFFIX`: never cleaned here.
+const MANAGED_BLOB_SUFFIX: &str = ".managed.blob";
+const DATA_FILE_EXTERNAL_PATHS_OPTION: &str = "data-file.external-paths";
+
+/// Files removed (or, in a dry run, that would be removed) by a cleanup.
+#[derive(Debug, Clone, Default, PartialEq, Eq)]
+pub struct OrphanFilesCleanResult {
+    pub deleted_file_count: u64,
+    pub deleted_file_total_bytes: u64,
+    /// Full paths, sorted.
+    pub deleted_files: Vec<String>,
+}
+
+/// Remove files that no snapshot, tag or changelog of the table references.
+pub struct RemoveOrphanFiles<'a> {
+    table: &'a Table,
+    older_than_millis: Option<i64>,
+    dry_run: bool,
+    parallelism: usize,
+    current_time_millis: Option<i64>,
+}
+
+impl<'a> RemoveOrphanFiles<'a> {
+    pub(crate) fn new(table: &'a Table) -> Self {
+        Self {
+            table,
+            older_than_millis: None,
+            dry_run: false,
+            parallelism: FILE_OPERATION_CONCURRENCY,
+            current_time_millis: None,
+        }
+    }
+
+    /// Maximum number of concurrent manifest reads and file deletions.
+    pub fn with_parallelism(&mut self, parallelism: usize) -> &mut Self {
+        self.parallelism = parallelism.max(1);
+        self
+    }
+
+    /// Only files last modified before this epoch-millisecond timestamp are
+    /// candidates. Defaults to one day ago; must be in the past.
+    pub fn with_older_than_millis(&mut self, older_than_millis: i64) -> &mut 
Self {
+        self.older_than_millis = Some(older_than_millis);
+        self
+    }
+
+    /// Report the orphan files without deleting them.
+    pub fn with_dry_run(&mut self, dry_run: bool) -> &mut Self {
+        self.dry_run = dry_run;
+        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
+    }
+
+    pub async fn execute(&self) -> Result<OrphanFilesCleanResult> {
+        self.table.ensure_not_branch_reference_for_write()?;
+        let now = self.current_time_millis.unwrap_or_else(current_time_millis);
+        let older_than = match self.older_than_millis {
+            None => now - DEFAULT_OLDER_THAN_MS,
+            Some(value) if value < now => value,
+            Some(_) => {
+                return Err(Error::DataInvalid {
+                    message: "older_than must be earlier than now; files being 
written and not \
+                              yet referenced by a snapshot would be deleted"
+                        .to_string(),
+                    source: None,
+                })
+            }
+        };
+        let clean = Clean {
+            file_io: self.table.file_io().clone(),
+            table_location: 
self.table.location().trim_end_matches('/').to_string(),
+            older_than,
+        };
+
+        let branches = self.valid_branches().await?;
+        let mut deleted = BTreeMap::new();
+        for branch in &branches {
+            
deleted.extend(clean.non_snapshot_files(&self.branch_root(branch)).await?);
+        }
+
+        let candidates = clean
+            .candidate_files(
+                self.table.schema().partition_keys().len(),
+                self.table
+                    .schema()
+                    .options()
+                    .get(DATA_FILE_EXTERNAL_PATHS_OPTION),
+            )
+            .await?;
+        if !candidates.is_empty() {
+            let mut used = HashSet::new();
+            for branch in &branches {
+                self.collect_used_files(&clean, branch, &mut used).await?;
+            }
+            for (path, size) in candidates {
+                if !used.contains(file_name(&path)) {
+                    deleted.insert(path, size);
+                }
+            }
+        }
+
+        if !self.dry_run {
+            let deletes = deleted
+                .keys()
+                .map(|path| clean.delete_quietly(path))
+                .collect::<Vec<_>>();
+            stream::iter(deletes)
+                .buffer_unordered(self.parallelism)
+                .collect::<Vec<_>>()
+                .await;
+        }
+        Ok(OrphanFilesCleanResult {
+            deleted_file_count: deleted.len() as u64,
+            deleted_file_total_bytes: deleted.values().sum(),
+            deleted_files: deleted.into_keys().collect(),
+        })
+    }
+
+    /// Every branch plus main. A branch without a schema aborts the run, as in
+    /// Java: its files cannot be accounted for.
+    async fn valid_branches(&self) -> Result<Vec<String>> {
+        let branch_manager = BranchManager::new(
+            self.table.file_io().clone(),
+            self.table.location().to_string(),
+        );
+        let branches = branch_manager.list_all().await?;
+        let mut abnormal = Vec::new();
+        for branch in &branches {
+            let schema_manager = 
self.table.schema_manager().with_branch(branch);
+            if schema_manager.latest().await?.is_none() {
+                abnormal.push(branch.clone());
+            }
+        }
+        if !abnormal.is_empty() {
+            return Err(Error::DataInvalid {
+                message: format!(
+                    "Branches {abnormal:?} have no schemas. Orphan files 
cleaning aborted. \
+                     Please check these branches manually."
+                ),
+                source: None,
+            });
+        }
+        let mut all = branches;
+        all.push(MAIN_BRANCH.to_string());
+        Ok(all)
+    }
+
+    fn branch_root(&self, branch: &str) -> String {
+        if branch == MAIN_BRANCH {
+            self.table.location().trim_end_matches('/').to_string()
+        } else {
+            BranchManager::new(
+                self.table.file_io().clone(),
+                self.table.location().to_string(),
+            )
+            .branch_path(branch)
+        }
+    }
+
+    /// Names of every file a snapshot, tag or long-lived changelog of `branch`
+    /// references. Java `LocalOrphanFilesClean#getUsedFiles`.
+    async fn collect_used_files(
+        &self,
+        clean: &Clean,
+        branch: &str,
+        used: &mut HashSet<String>,
+    ) -> Result<()> {
+        let snapshot_manager = 
self.table.snapshot_manager().with_branch(branch);
+        let mut owners = Vec::new();
+        for id in snapshot_manager.list_all_ids().await? {
+            if let Some(snapshot) = 
snapshot_manager.try_get_snapshot(id).await? {
+                owners.push(Owner {
+                    file: Some(snapshot_manager.snapshot_path(id)),
+                    snapshot,
+                });
+            }
+        }
+        let tag_manager = if branch == MAIN_BRANCH {
+            self.table.tag_manager()
+        } else {
+            self.table.tag_manager().with_branch(branch)
+        };
+        for (_, snapshot) in tag_manager.list_all().await? {
+            owners.push(Owner {
+                file: None,
+                snapshot,
+            });
+        }
+        owners.extend(clean.changelogs(&self.branch_root(branch)).await?);
+
+        let mut manifests = HashSet::new();
+        for owner in &owners {
+            if let Some(files) = clean.metadata_files(&snapshot_manager, 
owner).await? {
+                used.extend(files.names);
+                manifests.extend(files.manifests);
+            }
+        }
+        let paths = manifests
+            .iter()
+            .map(|name| snapshot_manager.manifest_path(name))
+            .collect::<Vec<_>>();
+        let reads = paths
+            .iter()
+            .map(|path| Manifest::read(&clean.file_io, path))
+            .collect::<Vec<_>>();
+        let entries = stream::iter(reads)
+            .buffer_unordered(self.parallelism)
+            .try_collect::<Vec<_>>()
+            .await?;
+        for entry in entries.into_iter().flatten() {
+            used.insert(entry.file().file_name.clone());
+            used.extend(entry.file().extra_files.iter().cloned());
+        }
+        Ok(())
+    }
+}
+
+/// A snapshot-like object that references files, and the file that makes it
+/// live (`None` for a tag, whose manifests must always be readable).
+struct Owner {
+    file: Option<String>,
+    snapshot: Snapshot,
+}
+
+struct MetadataFiles {
+    names: Vec<String>,
+    manifests: Vec<String>,
+}
+
+struct Clean {
+    file_io: FileIO,
+    table_location: String,
+    older_than: i64,
+}
+
+impl Clean {
+    fn old_enough(&self, status: &FileStatus) -> bool {
+        // A store that reports no modification time never exposes a file to 
deletion.
+        status
+            .last_modified
+            .is_some_and(|modified| modified.timestamp_millis() < 
self.older_than)
+    }
+
+    /// Files in `dir`, or nothing when it does not exist.
+    async fn list(&self, dir: &str) -> Result<Vec<FileStatus>> {
+        if !self.file_io.exists_dir(dir).await? {
+            return Ok(Vec::new());
+        }
+        self.file_io.list_status(dir).await
+    }
+
+    /// Old non-snapshot files in the snapshot and changelog directories, such
+    /// as temporary files of interrupted commits. Java 
`cleanBranchSnapshotDir`.
+    ///
+    /// Only `snapshot-<id>` / `changelog-<id>` and the hint files are kept: a
+    /// writer's temporary file is `snapshot-<id>.tmp-<uuid>`, which shares the
+    /// prefix.
+    async fn non_snapshot_files(&self, branch_root: &str) -> 
Result<Vec<(String, u64)>> {
+        let mut files = Vec::new();
+        for (dir, prefix) in [
+            ("snapshot", SNAPSHOT_PREFIX),
+            ("changelog", CHANGELOG_PREFIX),
+        ] {
+            for status in self.list(&format!("{branch_root}/{dir}")).await? {
+                let name = file_name(&status.path);
+                if !status.is_dir
+                    && !is_numbered(name, prefix)
+                    && name != EARLIEST
+                    && name != LATEST
+                    && self.old_enough(&status)
+                {
+                    files.push((status.path, status.size));
+                }
+            }
+        }
+        Ok(files)
+    }
+
+    /// Old files in the directories Paimon writes data and metadata to, keyed
+    /// by full path. Java `getCandidateDeletingFiles`.
+    async fn candidate_files(
+        &self,
+        partition_depth: usize,
+        external_paths: Option<&String>,
+    ) -> Result<BTreeMap<String, u64>> {
+        let mut dirs = ["manifest", "index", "statistics"]
+            .map(|dir| format!("{}/{dir}", self.table_location))
+            .to_vec();
+        dirs.extend(
+            self.bucket_dirs(&self.table_location, partition_depth)

Review Comment:
   Fixed in 2f6f156. Bucket candidates are now enumerated under 
`Table::data_file_location()`, and each external root applies the same 
`data-file.path-directory` as the writer. Manifest, index and statistics 
directories stay under the table location.
   
   `test_relative_data_directory_is_scanned` and 
`test_absolute_data_directory_is_scanned` plant an old orphan next to a live 
file. Only the orphan is removed. Expiration got the matching fix in #965 
(f7ba2d3), so the two now agree with writers and readers.
   



##########
crates/paimon/src/table/orphan_files_clean.rs:
##########
@@ -0,0 +1,531 @@
+// 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.
+
+//! Orphan file cleanup.
+//!
+//! Reference: Java `OrphanFilesClean` and `LocalOrphanFilesClean`.
+//!
+//! A file is an orphan when it sits in a directory Paimon writes to (the
+//! manifest, index and statistics directories, bucket directories and data
+//! file external paths), is older than `older_than`, and no snapshot, tag or
+//! long-lived changelog of any branch references it. Temporary files left in
+//! snapshot and changelog directories by interrupted commits are removed too.
+//!
+//! `older_than` defaults to one day ago and must be in the past, so files that
+//! an in-flight write has created but not yet committed are never candidates.
+//! Unlike Java, which reads a missing manifest as empty, a missing manifest of
+//! a snapshot that still exists aborts the run: its data files would otherwise
+//! look unreferenced and be deleted.
+
+use crate::io::{FileIO, FileStatus};
+use crate::spec::{IndexManifest, Manifest, ManifestList, Snapshot};
+use crate::table::snapshot_deletion::reassign_plan_file;
+use crate::table::{BranchManager, SnapshotManager, Table};
+use crate::{Error, Result};
+use futures::{stream, StreamExt, TryStreamExt};
+use std::collections::{BTreeMap, HashSet};
+use std::time::{SystemTime, UNIX_EPOCH};
+
+const DEFAULT_OLDER_THAN_MS: i64 = 24 * 60 * 60 * 1000;
+const FILE_OPERATION_CONCURRENCY: usize = 16;
+const MAIN_BRANCH: &str = "main";
+const SNAPSHOT_PREFIX: &str = "snapshot-";
+const CHANGELOG_PREFIX: &str = "changelog-";
+const EARLIEST: &str = "EARLIEST";
+const LATEST: &str = "LATEST";
+const BUCKET_PREFIX: &str = "bucket-";
+/// Java `ManagedBlobReferenceFile.MANAGED_BLOB_SUFFIX`: never cleaned here.
+const MANAGED_BLOB_SUFFIX: &str = ".managed.blob";
+const DATA_FILE_EXTERNAL_PATHS_OPTION: &str = "data-file.external-paths";
+
+/// Files removed (or, in a dry run, that would be removed) by a cleanup.
+#[derive(Debug, Clone, Default, PartialEq, Eq)]
+pub struct OrphanFilesCleanResult {
+    pub deleted_file_count: u64,
+    pub deleted_file_total_bytes: u64,
+    /// Full paths, sorted.
+    pub deleted_files: Vec<String>,
+}
+
+/// Remove files that no snapshot, tag or changelog of the table references.
+pub struct RemoveOrphanFiles<'a> {
+    table: &'a Table,
+    older_than_millis: Option<i64>,
+    dry_run: bool,
+    parallelism: usize,
+    current_time_millis: Option<i64>,
+}
+
+impl<'a> RemoveOrphanFiles<'a> {
+    pub(crate) fn new(table: &'a Table) -> Self {
+        Self {
+            table,
+            older_than_millis: None,
+            dry_run: false,
+            parallelism: FILE_OPERATION_CONCURRENCY,
+            current_time_millis: None,
+        }
+    }
+
+    /// Maximum number of concurrent manifest reads and file deletions.
+    pub fn with_parallelism(&mut self, parallelism: usize) -> &mut Self {
+        self.parallelism = parallelism.max(1);
+        self
+    }
+
+    /// Only files last modified before this epoch-millisecond timestamp are
+    /// candidates. Defaults to one day ago; must be in the past.
+    pub fn with_older_than_millis(&mut self, older_than_millis: i64) -> &mut 
Self {
+        self.older_than_millis = Some(older_than_millis);
+        self
+    }
+
+    /// Report the orphan files without deleting them.
+    pub fn with_dry_run(&mut self, dry_run: bool) -> &mut Self {
+        self.dry_run = dry_run;
+        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
+    }
+
+    pub async fn execute(&self) -> Result<OrphanFilesCleanResult> {
+        self.table.ensure_not_branch_reference_for_write()?;
+        let now = self.current_time_millis.unwrap_or_else(current_time_millis);
+        let older_than = match self.older_than_millis {
+            None => now - DEFAULT_OLDER_THAN_MS,
+            Some(value) if value < now => value,
+            Some(_) => {
+                return Err(Error::DataInvalid {
+                    message: "older_than must be earlier than now; files being 
written and not \
+                              yet referenced by a snapshot would be deleted"
+                        .to_string(),
+                    source: None,
+                })
+            }
+        };
+        let clean = Clean {
+            file_io: self.table.file_io().clone(),
+            table_location: 
self.table.location().trim_end_matches('/').to_string(),
+            older_than,
+        };
+
+        let branches = self.valid_branches().await?;
+        let mut deleted = BTreeMap::new();
+        for branch in &branches {
+            
deleted.extend(clean.non_snapshot_files(&self.branch_root(branch)).await?);
+        }
+
+        let candidates = clean
+            .candidate_files(
+                self.table.schema().partition_keys().len(),
+                self.table
+                    .schema()
+                    .options()
+                    .get(DATA_FILE_EXTERNAL_PATHS_OPTION),
+            )
+            .await?;
+        if !candidates.is_empty() {
+            let mut used = HashSet::new();
+            for branch in &branches {
+                self.collect_used_files(&clean, branch, &mut used).await?;
+            }
+            for (path, size) in candidates {
+                if !used.contains(file_name(&path)) {
+                    deleted.insert(path, size);
+                }
+            }
+        }
+
+        if !self.dry_run {
+            let deletes = deleted
+                .keys()
+                .map(|path| clean.delete_quietly(path))
+                .collect::<Vec<_>>();
+            stream::iter(deletes)
+                .buffer_unordered(self.parallelism)
+                .collect::<Vec<_>>()
+                .await;
+        }
+        Ok(OrphanFilesCleanResult {
+            deleted_file_count: deleted.len() as u64,
+            deleted_file_total_bytes: deleted.values().sum(),
+            deleted_files: deleted.into_keys().collect(),
+        })
+    }
+
+    /// Every branch plus main. A branch without a schema aborts the run, as in
+    /// Java: its files cannot be accounted for.
+    async fn valid_branches(&self) -> Result<Vec<String>> {
+        let branch_manager = BranchManager::new(
+            self.table.file_io().clone(),
+            self.table.location().to_string(),
+        );
+        let branches = branch_manager.list_all().await?;
+        let mut abnormal = Vec::new();
+        for branch in &branches {
+            let schema_manager = 
self.table.schema_manager().with_branch(branch);
+            if schema_manager.latest().await?.is_none() {
+                abnormal.push(branch.clone());
+            }
+        }
+        if !abnormal.is_empty() {
+            return Err(Error::DataInvalid {
+                message: format!(
+                    "Branches {abnormal:?} have no schemas. Orphan files 
cleaning aborted. \
+                     Please check these branches manually."
+                ),
+                source: None,
+            });
+        }
+        let mut all = branches;
+        all.push(MAIN_BRANCH.to_string());
+        Ok(all)
+    }
+
+    fn branch_root(&self, branch: &str) -> String {
+        if branch == MAIN_BRANCH {
+            self.table.location().trim_end_matches('/').to_string()
+        } else {
+            BranchManager::new(
+                self.table.file_io().clone(),
+                self.table.location().to_string(),
+            )
+            .branch_path(branch)
+        }
+    }
+
+    /// Names of every file a snapshot, tag or long-lived changelog of `branch`
+    /// references. Java `LocalOrphanFilesClean#getUsedFiles`.
+    async fn collect_used_files(
+        &self,
+        clean: &Clean,
+        branch: &str,
+        used: &mut HashSet<String>,
+    ) -> Result<()> {
+        let snapshot_manager = 
self.table.snapshot_manager().with_branch(branch);
+        let mut owners = Vec::new();
+        for id in snapshot_manager.list_all_ids().await? {
+            if let Some(snapshot) = 
snapshot_manager.try_get_snapshot(id).await? {
+                owners.push(Owner {
+                    file: Some(snapshot_manager.snapshot_path(id)),
+                    snapshot,
+                });
+            }
+        }
+        let tag_manager = if branch == MAIN_BRANCH {
+            self.table.tag_manager()
+        } else {
+            self.table.tag_manager().with_branch(branch)
+        };
+        for (_, snapshot) in tag_manager.list_all().await? {
+            owners.push(Owner {
+                file: None,
+                snapshot,
+            });
+        }
+        owners.extend(clean.changelogs(&self.branch_root(branch)).await?);
+
+        let mut manifests = HashSet::new();
+        for owner in &owners {
+            if let Some(files) = clean.metadata_files(&snapshot_manager, 
owner).await? {
+                used.extend(files.names);
+                manifests.extend(files.manifests);
+            }
+        }
+        let paths = manifests
+            .iter()
+            .map(|name| snapshot_manager.manifest_path(name))
+            .collect::<Vec<_>>();
+        let reads = paths
+            .iter()
+            .map(|path| Manifest::read(&clean.file_io, path))
+            .collect::<Vec<_>>();
+        let entries = stream::iter(reads)
+            .buffer_unordered(self.parallelism)
+            .try_collect::<Vec<_>>()
+            .await?;
+        for entry in entries.into_iter().flatten() {
+            used.insert(entry.file().file_name.clone());
+            used.extend(entry.file().extra_files.iter().cloned());
+        }
+        Ok(())
+    }
+}
+
+/// A snapshot-like object that references files, and the file that makes it
+/// live (`None` for a tag, whose manifests must always be readable).
+struct Owner {
+    file: Option<String>,
+    snapshot: Snapshot,
+}
+
+struct MetadataFiles {
+    names: Vec<String>,
+    manifests: Vec<String>,
+}
+
+struct Clean {
+    file_io: FileIO,
+    table_location: String,
+    older_than: i64,
+}
+
+impl Clean {
+    fn old_enough(&self, status: &FileStatus) -> bool {
+        // A store that reports no modification time never exposes a file to 
deletion.
+        status
+            .last_modified
+            .is_some_and(|modified| modified.timestamp_millis() < 
self.older_than)
+    }
+
+    /// Files in `dir`, or nothing when it does not exist.
+    async fn list(&self, dir: &str) -> Result<Vec<FileStatus>> {
+        if !self.file_io.exists_dir(dir).await? {
+            return Ok(Vec::new());
+        }
+        self.file_io.list_status(dir).await
+    }
+
+    /// Old non-snapshot files in the snapshot and changelog directories, such
+    /// as temporary files of interrupted commits. Java 
`cleanBranchSnapshotDir`.
+    ///
+    /// Only `snapshot-<id>` / `changelog-<id>` and the hint files are kept: a
+    /// writer's temporary file is `snapshot-<id>.tmp-<uuid>`, which shares the
+    /// prefix.
+    async fn non_snapshot_files(&self, branch_root: &str) -> 
Result<Vec<(String, u64)>> {
+        let mut files = Vec::new();
+        for (dir, prefix) in [
+            ("snapshot", SNAPSHOT_PREFIX),
+            ("changelog", CHANGELOG_PREFIX),
+        ] {
+            for status in self.list(&format!("{branch_root}/{dir}")).await? {
+                let name = file_name(&status.path);
+                if !status.is_dir
+                    && !is_numbered(name, prefix)
+                    && name != EARLIEST
+                    && name != LATEST
+                    && self.old_enough(&status)
+                {
+                    files.push((status.path, status.size));
+                }
+            }
+        }
+        Ok(files)
+    }
+
+    /// Old files in the directories Paimon writes data and metadata to, keyed
+    /// by full path. Java `getCandidateDeletingFiles`.
+    async fn candidate_files(
+        &self,
+        partition_depth: usize,
+        external_paths: Option<&String>,
+    ) -> Result<BTreeMap<String, u64>> {
+        let mut dirs = ["manifest", "index", "statistics"]
+            .map(|dir| format!("{}/{dir}", self.table_location))
+            .to_vec();
+        dirs.extend(
+            self.bucket_dirs(&self.table_location, partition_depth)
+                .await?,
+        );
+        for external in external_paths
+            .into_iter()
+            .flat_map(|paths| paths.split(','))
+            .map(str::trim)
+            .filter(|path| !path.is_empty())
+        {
+            dirs.extend(
+                self.bucket_dirs(external.trim_end_matches('/'), 
partition_depth)
+                    .await?,
+            );
+        }
+
+        let mut candidates = BTreeMap::new();
+        for dir in dirs {
+            for status in self.list(&dir).await? {
+                if !status.is_dir

Review Comment:
   Fixed in 2f6f156. External buckets are now listed recursively. That reaches 
the four `entropy-inject` hash directories, and the cut-off, the managed-BLOB 
exclusion and the referenced-file protection still apply. 
`test_entropy_injected_external_files_are_scanned` writes through the real 
`entropy-inject` writer and plants an old orphan beside the written file. Only 
the orphan is removed, and the table still reads correctly.
   



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