JingsongLi commented on code in PR #966: URL: https://github.com/apache/paimon-rust/pull/966#discussion_r4172624443
########## 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] Use the data root for automatic expiration cleanup This shared helper ignores the supported data-file.path-directory option. I created a real SQL table with that option=data and snapshot min/max=1, inserted row1, then overwrote it with row2. The overwrite succeeds, main reads only row2, and automatic expiration removes snapshot1 metadata, but its original `<table>/data/bucket-0/*.parquet` remains. This helper instead attempts `<table>/bucket-0`, so ordinary production commits silently leak obsolete data/sidecars/bucket-local indexes. Resolve those paths with Table::data_file_location(), retain the metadata root for manifests/global indexes, and test actual automatic expiration on a configured data directory. Same shared defect as #968/#967; this probe requires no explicit expiration call. -- 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]
