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]