[ 
https://issues.apache.org/jira/browse/SPARK-59631?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

ASF GitHub Bot updated SPARK-59631:
-----------------------------------
    Labels: pull-request-available  (was: )

> Persisted V1 time-travel relations are recached after table writes
> ------------------------------------------------------------------
>
>                 Key: SPARK-59631
>                 URL: https://issues.apache.org/jira/browse/SPARK-59631
>             Project: Spark
>          Issue Type: Bug
>          Components: SQL
>    Affects Versions: 4.1.0, 4.2.0
>            Reporter: Joel Robin
>            Priority: Major
>              Labels: pull-request-available
>
> h2. Description
> A persisted DataFrame reading a fixed table version through the V1 data 
> source path is unnecessarily recached after a later write advances the live 
> table.
> The pinned snapshot is immutable, but recaching re-executes its entire plan. 
> This can repeat expensive or non-idempotent UDFs even though the DataFrame 
> remains registered as cached.
> h2. Reproduction
> Tested with unpatched Apache Spark 4.2.0, Scala 2.13, Java 17, and Delta Lake 
> master commit {{052429d500f64f7b3d482f28092b2eb066d20cdb}}. Delta V2 mode is 
> disabled to exercise the affected V1 relation path.
> Compile and run the following standalone application with the local Spark and 
> Delta classpath:
> {code:scala}
> import java.nio.file.Files
> import java.util.UUID
> import org.apache.spark.sql.SparkSession
> import org.apache.spark.sql.functions.{col, udf}
> import org.apache.spark.storage.StorageLevel
> object DeltaTimeTravelCacheRepro {
>   def main(args: Array[String]): Unit = {
>     val root = Files.createTempDirectory("delta-time-travel-cache-repro-")
>     val table = "tt_cache_" + UUID.randomUUID().toString.replace("-", "")
>     val spark = SparkSession.builder()
>       .appName("DeltaTimeTravelCacheRepro")
>       .master("local[2]")
>       .config("spark.ui.enabled", "false")
>       .config("spark.sql.shuffle.partitions", "2")
>       .config("spark.sql.warehouse.dir", 
> root.resolve("warehouse").toUri.toString)
>       .config("spark.sql.extensions", 
> "io.delta.sql.DeltaSparkSessionExtension")
>       .config(
>         "spark.sql.catalog.spark_catalog",
>         "org.apache.spark.sql.delta.catalog.DeltaCatalog")
>       .config("spark.databricks.delta.v2.enableMode", "NONE")
>       .getOrCreate()
>     var cached: Option[org.apache.spark.sql.DataFrame] = None
>     try {
>       spark.range(10).write.format("delta").saveAsTable(table)
>       val udfCalls = spark.sparkContext.longAccumulator("udf-calls")
>       val expensiveUdf = udf { id: Long =>
>         udfCalls.add(1L)
>         s"$id-${UUID.randomUUID()}"
>       }.asNondeterministic()
>       val pinned = spark.read
>         .option("versionAsOf", 0)
>         .table(table)
>         .withColumn("token", expensiveUdf(col("id")))
>         .persist()
>       cached = Some(pinned)
>       // Materialize the cache before advancing the live table.
>       val before = pinned.orderBy("id").select("token").collect().toSeq
>       val callsBefore = udfCalls.value
>       // This creates version 1. The pinned DataFrame still reads version 0.
>       spark.range(10, 20)
>         .write
>         .format("delta")
>         .mode("append")
>         .saveAsTable(table)
>       val after = pinned.orderBy("id").select("token").collect().toSeq
>       val callsAfter = udfCalls.value
>       val pinnedCount = pinned.count()
>       val liveCount = spark.table(table).count()
>       println(s"is_cached=${pinned.storageLevel != StorageLevel.NONE}")
>       println(s"udf_calls=$callsBefore->$callsAfter")
>       println(s"tokens_same=${before == after}")
>       println(s"pinned_count=$pinnedCount")
>       println(s"live_count=$liveCount")
>       assert(callsBefore == 10)
>       assert(callsAfter == callsBefore)
>       assert(after == before)
>       assert(pinnedCount == 10)
>       assert(liveCount == 20)
>     } finally {
>       cached.foreach(_.unpersist(blocking = true))
>       spark.sql(s"DROP TABLE IF EXISTS $table")
>       spark.stop()
>     }
>   }
> }
> {code}
> The nondeterministic Scala UDF is a stand-in for the reported pandas UDF that 
> performs LLM calls. It isolates the cache invalidation behavior without 
> requiring Python or Arrow.
> h2. Actual Result
> On unpatched Spark 4.2.0:
> {noformat}
> is_cached=true
> udf_calls=10->20
> tokens_same=false
> pinned_count=10
> live_count=20
> {noformat}
> The second action recomputes the cached plan and invokes the UDF again. The 
> cache entry remains registered, which explains why the DataFrame still 
> reports as cached while its materialized buffers have been discarded.
> h2. Expected Result
> {noformat}
> is_cached=true
> udf_calls=10->10
> tokens_same=true
> pinned_count=10
> live_count=20
> {noformat}
> A write that advances the live table should refresh live-table cache entries 
> but should not invalidate a cache entry for an immutable, version-pinned 
> snapshot. Explicit operations such as {{REFRESH TABLE}}, {{refreshByPath}}, 
> or uncache should retain their existing ability to invalidate pinned 
> snapshots.



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to