heye1005 commented on PR #7715:
URL: https://github.com/apache/paimon/pull/7715#issuecomment-5999393027

   > **[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.
   
   Thanks for suggests. I removed the spark filters that dropped missing 
markers and moved the check into the final execution path, so cache loss and 
recomputation won’t bypass it. I also added a regression test for this case.


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