zhuxiangyi commented on code in PR #965: URL: https://github.com/apache/paimon-rust/pull/965#discussion_r4178227748
########## crates/paimon/src/table/snapshot_deletion.rs: ########## @@ -0,0 +1,504 @@ +// 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, Review Comment: Fixed in 5e9cbad. I rebased onto main 8c3527b and reproduced the E0432 first. #1016 removed the legacy deletion-vector location fallback from reads altogether, not only that helper, so expiration now resolves index files only through `committed_index_file_path`, the layout reads use today. The round-1 legacy-location test is dropped, because that layout is no longer readable on main. The merge with current main then compiles. The full core suite, the DataFusion procedure tests and clippy with `-D warnings` pass on top of 8c3527b. ########## crates/integrations/datafusion/src/procedures.rs: ########## @@ -1214,6 +1224,82 @@ async fn proc_list_policies( ) } +/// `CALL sys.expire_snapshots`, with the arguments of Java's Flink and Spark +/// `ExpireSnapshotsProcedure`. `options` are dynamic table options applied for +/// this call only (for example `snapshot.time-retained`). +async fn proc_expire_snapshots( + ctx: &SessionContext, + catalog: &Arc<dyn Catalog>, + catalog_name: &str, + args: &HashMap<String, String>, +) -> DFResult<DataFrame> { + let mut table = get_table(catalog, catalog_name, args).await?; + if let Some(options) = args.get("options") { + table = table.copy_with_options(parse_key_value_options(options)?); + } + let mut expire = table.new_expire_snapshots(); + if let Some(value) = optional_i32_arg(args, "retain_max")? { + expire.with_retain_max(value); + } + if let Some(value) = optional_i32_arg(args, "retain_min")? { + expire.with_retain_min(value); + } + if let Some(value) = optional_i32_arg(args, "max_deletes")? { + expire.with_max_deletes(value); + } + if let Some(older_than) = args.get("older_than").filter(|v| !v.trim().is_empty()) { + expire.with_older_than_millis(parse_older_than(older_than)?); + } + let deleted = expire.execute().await.map_err(to_datafusion_error)?; + let deleted = i32::try_from(deleted).unwrap_or(i32::MAX); + + let schema = Arc::new(Schema::new(vec![Field::new( + "deleted_snapshots_count", + ArrowDataType::Int32, + false, + )])); + let batch = RecordBatch::try_new(schema, vec![Arc::new(Int32Array::from(vec![deleted]))])?; + ctx.read_batch(batch) +} + +fn optional_i32_arg(args: &HashMap<String, String>, name: &str) -> DFResult<Option<i32>> { + args.get(name) + .map(|value| { + value.trim().parse::<i32>().map_err(|_| { + DataFusionError::Plan(format!("Invalid integer for '{name}': '{value}'")) + }) + }) + .transpose() +} + +/// `older_than` as epoch milliseconds, or as a timestamp such as +/// `2024-01-01 12:00:00` in the session's local time zone, which is how Java +/// (`DateTimeUtils.parseTimestampData` with the default time zone) reads it. +fn parse_older_than(value: &str) -> DFResult<i64> { + let value = value.trim(); + if let Ok(millis) = value.parse::<i64>() { + return Ok(millis); + } + let naive = ["%Y-%m-%d %H:%M:%S%.f", "%Y-%m-%dT%H:%M:%S%.f"] + .iter() + .find_map(|format| NaiveDateTime::parse_from_str(value, format).ok()) + .or_else(|| { + NaiveDate::parse_from_str(value, "%Y-%m-%d") + .ok() + .and_then(|date| date.and_hms_opt(0, 0, 0)) + }) + .ok_or_else(|| DataFusionError::Plan(format!("Invalid older_than timestamp: '{value}'")))?; + Local Review Comment: Fixed in 22d9e7e. `older_than` now goes through `spec::parse_local_timestamp_millis`. That is main's Java-compatible `scan.timestamp` parser (jiff's compatible disambiguation, as in `LocalDateTime.atZone`), generalized to take an option name. A time in a gap moves forward, and one in a fold takes the earlier instant. Epoch milliseconds are still accepted. `test_parse_older_than_follows_java_across_dst` runs the parser in a child process pinned to `TZ=America/New_York`, like the existing `scan.timestamp` test. `2024-03-10 02:30:00` gives 1710055800000 (03:30 EDT), and `2024-11-03 01:30:00` gives the earlier 1730611800000. On the old code it fails with your error message. -- 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]
