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]

Reply via email to