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