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


##########
crates/integrations/datafusion/src/procedures.rs:
##########
@@ -1214,6 +1226,141 @@ 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)
+}
+
+/// `CALL sys.remove_orphan_files`, with the arguments of Java's Flink and
+/// Spark `RemoveOrphanFilesProcedure`. Only one table per call and the local
+/// mode are supported.
+async fn proc_remove_orphan_files(
+    ctx: &SessionContext,
+    catalog: &Arc<dyn Catalog>,
+    catalog_name: &str,
+    args: &HashMap<String, String>,
+) -> DFResult<DataFrame> {
+    if require_arg(args, "table")?.trim().ends_with(".*") {
+        return Err(DataFusionError::NotImplemented(
+            "remove_orphan_files on all tables of a database ('db.*') is not 
supported yet"
+                .to_string(),
+        ));
+    }
+    if let Some(mode) = args.get("mode") {
+        if !mode.trim().eq_ignore_ascii_case("local") {
+            return Err(DataFusionError::NotImplemented(format!(
+                "remove_orphan_files only supports mode => 'local', got 
'{mode}'"
+            )));
+        }
+    }
+    let table = get_table(catalog, catalog_name, args).await?;
+    let mut clean = table.new_remove_orphan_files();
+    if let Some(older_than) = args.get("older_than").filter(|v| 
!v.trim().is_empty()) {
+        clean.with_older_than_millis(parse_older_than(older_than)?);
+    }
+    if let Some(dry_run) = args.get("dry_run") {
+        let dry_run = dry_run.trim().parse::<bool>().map_err(|_| {
+            DataFusionError::Plan(format!("Invalid boolean for 'dry_run': 
'{dry_run}'"))
+        })?;
+        clean.with_dry_run(dry_run);
+    }
+    if let Some(parallelism) = optional_i32_arg(args, "parallelism")? {
+        if parallelism < 1 {
+            return Err(DataFusionError::Plan(format!(
+                "parallelism must be at least 1, got {parallelism}"
+            )));
+        }
+        clean.with_parallelism(parallelism as usize);
+    }
+    let result = clean.execute().await.map_err(to_datafusion_error)?;
+
+    let schema = Arc::new(Schema::new(vec![
+        Field::new("deletedFileCount", ArrowDataType::Int64, false),
+        Field::new("deletedFileTotalLenInBytes", ArrowDataType::Int64, false),
+    ]));
+    let batch = RecordBatch::try_new(
+        schema,
+        vec![
+            Arc::new(Int64Array::from(vec![result.deleted_file_count as i64])),
+            Arc::new(Int64Array::from(vec![
+                result.deleted_file_total_bytes as i64,
+            ])),
+        ],
+    )?;
+    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:
   [P2] Match Java's forward resolution of DST gap timestamps
   
   With TZ=America/New_York, CALL sys.expire_snapshots(table => 'test_db.t1', 
older_than => '2024-03-10 02:30:00', retain_min => 1) fails with 'does not 
exist in the local time zone'. Java ProcedureUtils uses 
DateTimeUtils.parseTimestampData with the default time zone; 
LocalDateTime.atZone resolves this gap forward to 03:30 EDT (1710055800000). 
This rejects a valid Java procedure input instead of expiring normally. Use the 
same gap-forward / fold-earlier resolution already implemented for timestamp 
options, and cover this input in an isolated time-zone test.



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