linliu-code opened a new issue, #20104:
URL: https://github.com/apache/hudi/issues/20104

   ## Bug Description
   
   `CACHE TABLE` works on a Hudi table — queries do get an `InMemoryRelation` — 
but every subsequent
   write strands the previously cached dataset instead of recaching it. Entries 
accumulate one per
   write, `UNCACHE TABLE` reclaims none of them, and `spark.catalog.isCached` 
reports `false` the
   whole time.
   
   A plain Parquet table under the identical cycle does not do this, which is 
what localises the
   defect to Hudi.
   
   ### Measured, same session, same SQL, same harness
   
   Each iteration runs `CACHE TABLE`, then `INSERT INTO ... part='p1'` — the 
same partition that is
   already cached, so partition pruning cannot mask the behaviour. The count is
   `CacheManager.cachedData.size`, read reflectively.
   
   | iteration | Hudi: after cache → after write | Parquet: after cache → after 
write |
   | --- | --- | --- |
   | 1 | 1 → 1 | 1 → 1 |
   | 2 | 2 → 2 | 1 → 1 |
   | 3 | 3 → 3 | 1 → 1 |
   | 4 | 4 → 4 | 1 → 1 |
   | 5 | 5 → 5 | 1 → 1 |
   | **final** | **5** | **1** |
   | **after `UNCACHE TABLE`** | **5** (reclaims nothing) | **0** |
   
   Correctness is not affected on the SQL write path: after each insert the row 
count is right and
   the plan no longer uses the `InMemoryRelation`, so queries re-read rather 
than serving stale data.
   The problem is that the old materialised dataset is never released.
   
   ## Root cause
   
   `HoodieFileIndex` is a case class, so Scala derives `equals`/`hashCode` from 
**all** constructor
   parameters:
   
       case class HoodieFileIndex(spark: SparkSession,
                                  metaClient: HoodieTableMetaClient,
                                  schemaSpec: Option[StructType],
                                  options: Map[String, String],
                                  @transient fileStatusCache: FileStatusCache = 
NoopCache,
                                  includeLogFiles: Boolean = false,
                                  shouldEmbedFileSlices: Boolean = false)
   
   `@transient` affects serialization, not equality, so `fileStatusCache` is 
part of `equals`. It has
   no `equals` of its own, so it compares by identity — and Spark hands out a 
**new instance on every
   call**:
   
       FileStatusCache.getOrCreate(session) called twice in one session:
         same_instance = false
         equals        = false
         class         = 
org.apache.spark.sql.execution.datasources.SharedInMemoryCache$$anon$3
   
   `HoodieHadoopFsRelationFactory` calls 
`FileStatusCache.getOrCreate(sparkSession)` when it builds
   the relation, so a relation rebuilt by `refreshTable` carries a different 
`FileStatusCache`
   object, and therefore a `HoodieFileIndex` that is not `equals` to the cached 
one.
   
   Comparing the index from the table's analyzed plan before and after an 
`INSERT`, field by field:
   
       same_instance=false
       equals=false
       spark_same=true                   OK
       metaClient_equals=true            OK   (equals is basePath + tableType, 
value-based)
       schemaSpec_equals=true            OK
       options_equals=true               OK   (both set-diffs empty)
       fileStatusCache_same=false        <-- the only divergence
       includeLogFiles_equals=true       OK
       shouldEmbedFileSlices_equals=true OK
   
   `CacheManager.recacheByPlan` matches by plan, so it cannot find the old 
entry, leaves it
   materialised, and the next `CACHE TABLE` adds another. The same mismatch 
explains why
   `spark.catalog.isCached` returns `false`: it builds a fresh plan to look up, 
which also fails to
   match.
   
   This is a cache handle leaking into value equality. Two indexes with the 
same base path, schema,
   options and flags are the same table; which `FileStatusCache` instance they 
happen to hold is a
   performance detail.
   
   ## Proposed fix
   
   Exclude `fileStatusCache` from equality by giving `HoodieFileIndex` an 
explicit
   `equals`/`hashCode`/`canEqual` over the remaining parameters. That keeps the 
public constructor
   signature unchanged.
   
   Moving `fileStatusCache` into a second parameter list would achieve the same 
thing idiomatically,
   since Scala only derives `equals` from the first list — but 
`HoodieFileIndex` is public API and
   that would break positional construction for downstream callers, so it seems 
the worse trade.
   
   Not proposed: `spark` is also identity-compared. It is stable within a 
session so it does not
   cause this bug, and excluding it would change cross-session semantics that 
have not been examined
   here.
   
   ### Blast radius
   
   Across the repository, 29 non-test files reference `HoodieFileIndex`. None 
of them compare
   instances: a grep for `fileIndex ==`, `.equals(fileIndex)`, 
`Set[HoodieFileIndex]` and
   `Map[HoodieFileIndex, ...]` returns only the `def unapply(relation): 
Option[HoodieFileIndex]`
   extractors in the three `Spark*HoodiePruneFileSourcePartitions` rules, which 
match on type rather
   than equality. No test asserts on index equality either.
   
   So the only consumer of this equality is Spark's own plan machinery — 
`HadoopFsRelation`
   equality, plan canonicalization, and `CacheManager`. The change moves Hudi 
onto the behaviour
   Parquet already has rather than away from it.
   
   One question worth stating explicitly, since it looks alarming: if plans 
compare equal across a
   write, could a lookup then match a **stale** entry? That is exactly how 
Parquet already behaves.
   Spark's contract is that plans stay equal across a refresh and 
`recacheByPlan` re-executes to
   replace the contents — equality is table identity, freshness is 
`refreshTable`'s job.
   
   ## A separate defect, mentioned only so it is not conflated
   
   Writes through the DataFrame writer do not invalidate the cache at all, and 
that one **does**
   produce stale reads:
   
       rows before df.write...save(path) = 3
       rows after                        = 3      (the newly written row is 
invisible)
       isCached                          = true
   
   `HoodieSparkSqlWriter` only refreshes the catalog table when meta sync is 
enabled:
   
       if (metaSyncEnabled) {
         getHiveTableNames(hoodieConfig).foreach(name => {
           ...
           if (spark.catalog.databaseExists(syncDb) && 
spark.catalog.tableExists(qualifiedTableName)) {
             spark.catalog.refreshTable(qualifiedTableName)
           }
         })
       }
   
   This has a different mechanism from the equality bug above and is not fixed 
by it. Happy to file
   it separately if that is preferred.
   
   By extension, any writer outside the reading Spark session — a separate job, 
Flink, a streaming
   ingest — leaves a cached Hudi table stale with no signal, because Spark's 
`CacheManager` is
   session-local. That is arguably not Hudi's to fix in-session, but a cheap 
"has the timeline
   advanced" check that a session could poll would make it solvable.
   
   ## How to reproduce
   
   In a single Spark session:
   
       CREATE TABLE t (id int, name string, price double, part string)
         USING hudi PARTITIONED BY (part)
         TBLPROPERTIES (primaryKey = 'id', type = 'cow');
       INSERT INTO t VALUES (1,'a',cast(10.0 as double),'p1');
   
       -- repeat 5 times:
       CACHE TABLE t;
       SELECT count(*) FROM t;
       INSERT INTO t VALUES (2,'b',cast(20.0 as double),'p1');
   
   Then read `CacheManager.cachedData.size` (it is private; reflection works). 
It climbs by one per
   iteration on Hudi and stays at 1 on an otherwise identical `USING parquet` 
table. `UNCACHE TABLE`
   returns the Parquet count to 0 and leaves the Hudi count unchanged.
   
   ## Environment
   
   - Hudi: `master` (1.3.0-SNAPSHOT)
   - Spark 3.5, Scala 2.12, JDK 11
   - Table type: COPY_ON_WRITE, partitioned, catalog-registered
   
   ## Related
   
   - #16359 — in a query-only Spark session the latest visible commit is not 
updated, because Spark
     reuses the cached relation and therefore the same `HoodieFileIndex` 
instance. Same reuse
     behaviour, different consequence: that issue is about read staleness, this 
one is about the
     `CacheManager` entry being stranded when a write does force a rebuild.
   - #20057 — unbounded growth of `cachedAllInputFileSlices` in a long-lived 
session. Also rooted in
     the reused index, and overlaps with #16359's description of the same 
fields.
   
   ## What was not measured
   
   - Bytes retained per stranded entry, and whether Spark evicts them under 
memory pressure — so the
     practical severity is not established, only the accumulation
   - MERGE_ON_READ behaviour
   - Whether `spark.catalog.clearCache()` reclaims the stranded entries
   


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