JingsongLi commented on PR #7715: URL: https://github.com/apache/paimon/pull/7715#issuecomment-5997238836
**[P1] Preserve the abort condition when Spark recomputes cached references** There is still a data-loss path in the Spark implementation at head `c2cb7f0535cafa04ec6bc4d1b1822abecee8813b`. The two `isEmpty` checks validate the references during the initial reads, but the cached datasets do not provide a durable snapshot. If an executor loses cached partitions, Spark can recompute them from the original lineage and surviving shuffle output ([Spark persistence documentation](https://spark.apache.org/docs/3.5.8/rdd-programming-guide.html#rdd-persistence)). A concurrent commit can merge the old manifests while retaining their data files, and snapshot expiration can then remove those old manifests. On recomputation, the reader correctly emits a missing-manifest marker, but [line 177](https://github.com/apache/paimon/blob/c2cb7f0535cafa04ec6bc4d1b1822abecee8813b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/procedure/SparkOrphanFilesClean.scala#L177) filters that marker out instead of aborting. [Line 181](https://github.com/apache/paimon/blob/c2cb7f0535cafa04ec6bc4d1b1822abecee8813b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/procedure/SparkOrphanFilesClean.scala#L181) has the equivalent issue for manifest-list markers. The final `left_anti` join can therefore delete old data files that the current snapshot still references. I reproduced this with Spark 3.5.8, real Paimon commits and expiration, and simulated cache loss: 1. Create an append table with `bucket=1`, `bucket-key=id`, `manifest.merge-min-count=1`, and both snapshot retention limits set to `1`. 2. Insert two batches, set their data-file modification times to two days ago, and use the default one-day cleanup cutoff. 3. Call `doOrphanClean()`; both missing-manifest prechecks pass. Obtain the returned deletion dataset and initialize its execution plan. 4. Insert a third batch. Expiration removes obsolete merged manifests, while the latest snapshot still contains all three rows. 5. Simulate loss of the cached RDD contents while keeping the SQL cache registrations: ```scala spark.sparkContext.getPersistentRDDs.values.foreach(_.unpersist(blocking = true)) ``` 6. Execute `deleted.collect()`. It reports `(2, 932)`, and both old data files still referenced by the latest snapshot are deleted. This is a remaining gap in the intended concurrent-expiration fix; the baseline also has a similar risk. Please make a missing live manifest/manifest-list fail the execution even during recomputation, or preserve a global abort barrier in the final deletion query. A reliable checkpoint of the validated references is another option. A regression test covering commit + expiration + cache loss would protect this path. -- 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]
