JingsongLi commented on code in PR #967:
URL: https://github.com/apache/paimon-rust/pull/967#discussion_r4172614411


##########
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:
   [P2] Resolve expired data paths from the configured data root
   
   With a real Python table configured `data-file.path-directory=data`, INSERT 
writes the old Parquet under `<table>/data/bucket-0`. After INSERT OVERWRITE, 
the new Python `expire_snapshots(retain_min=1, retain_max=1)` returns 1 and 
only snapshot 2 remains, but the expired snapshot's original Parquet still 
exists. Latest rows remain correct; physical reclamation silently fails because 
this helper attempts `<table>/bucket-0` instead of the writer/read path. Use 
Table::data_file_location() for data and bucket-local sidecars/indexes, 
retaining table_location for metadata/global indexes, and add a 
configured-directory regression through the Python API. This is the same shared 
foundation defect reproduced on #968.



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