linliu-code opened a new issue, #20113:
URL: https://github.com/apache/hudi/issues/20113
## Bug Description
A `CACHE TABLE`'d Hudi table is not invalidated when the table is written
through the DataFrame
writer, so queries keep returning pre-write results with no indication
anything is wrong. A
Parquet table in the same situation is invalidated correctly.
Unlike #20104, which is about memory, this one returns **wrong answers**.
### Measured
Same session, same shape of write - straight to the table's path, bypassing
the catalog:
| | rows before write | rows after write | stale? |
| --- | --- | --- | --- |
| Hudi (`df.write.format("hudi")...save(path)`) | 1 | **1** | **yes** |
| Parquet (`df.write.parquet(path)`) | 1 | 2 | no |
`REFRESH TABLE` recovers both (each then reports 2).
`spark.catalog.isCached` still reports
`true` for the Hudi table while it is serving the stale result.
The row written targets a partition that is already cached, so partition
pruning cannot account
for the difference.
## Root cause
Spark's own file-source writer invalidates **by path**, which sidesteps plan
matching entirely:
InsertIntoHadoopFsRelationCommand.scala:207
sparkSession.sharedState.cacheManager.recacheByPath(sparkSession,
outputPath, fs)
CacheManager.recacheByPath(spark, resourcePath, fs)
recacheByCondition(spark, _.plan.exists(lookupAndRefresh(_, fs,
qualifiedPath)))
-> re-caches every entry whose plan reads that path
Hudi's write path never calls it. `DefaultSource.createRelation` goes to
`HoodieSparkSqlWriter`,
which does only this:
// HoodieSparkSqlWriter, at the end of a successful write
if (metaSyncEnabled) {
getHiveTableNames(hoodieConfig).foreach(name => {
val syncDb = hoodieConfig.getStringOrDefault(HIVE_DATABASE)
val qualifiedTableName = String.join(".", syncDb, name)
if (spark.catalog.databaseExists(syncDb) &&
spark.catalog.tableExists(qualifiedTableName)) {
spark.catalog.refreshTable(qualifiedTableName)
}
})
}
So the invalidation is conditional on meta sync being enabled and on the
synced name resolving in
the catalog. With meta sync off - the default for a plain `save(path)` -
nothing is invalidated
at all, and the cached table silently serves stale rows.
This also means a table registered in the catalog under a name that does not
match the meta-sync
configuration is missed even when meta sync is on.
## Suggested fix
Call `recacheByPath` on the written base path, unconditionally, alongside
the existing
`refreshTable` call. That matches what Spark's own writer does, needs no
knowledge of whether
the table is catalog-registered or meta-synced, and covers the case where
the same data is
cached under a different name.
## How to reproduce
In one Spark session:
CREATE TABLE t (id int, name string, price double, part string)
USING hudi PARTITIONED BY (part)
TBLPROPERTIES (primaryKey = 'id', type = 'cow') LOCATION '<path>';
INSERT INTO t VALUES (1,'a',cast(10.0 as double),'p1');
CACHE TABLE t;
SELECT count(*) FROM t; -- 1
then, from the same session, write to the path with the DataFrame writer:
spark.sql("select 2 as id, 'b' as name, cast(20.0 as double) as price,
'p1' as part")
.write.format("hudi")
.option("hoodie.table.name", "t")
.option("hoodie.datasource.write.recordkey.field", "id")
.option("hoodie.datasource.write.partitionpath.field", "part")
.option("hoodie.datasource.write.operation", "insert")
.mode("append").save("<path>")
SELECT count(*) FROM t; -- still 1; expected 2
REFRESH TABLE t;
SELECT count(*) FROM t; -- 2
Repeating the same steps with `USING parquet` and
`df.write.mode("append").partitionBy("part").parquet(path)`
returns 2 without the refresh.
## Wider consequence
Spark's `CacheManager` is session-local, so a write from a *different*
process - a separate
ingest job, Deltastreamer, Flink - cannot invalidate a cache held by a
reading session whatever
Hudi does here. That is a broader problem and not what this issue proposes
to fix; it is worth
noting only because the in-session case described above is the part that is
fixable today, and
it is currently broken where Parquet's is not.
## Related
- #20104 - the same feature area: a `CACHE TABLE` entry is stranded rather
than recached when a
SQL write does force a rebuild. Different mechanism (plan equality), and
its fix does not
address this one.
- #16359 - a query-only session not seeing later commits, rooted in the same
reused-relation
behaviour.
## Environment
- Hudi `master` (1.3.0-SNAPSHOT)
- Spark 3.5, Scala 2.12, JDK 11
- COPY_ON_WRITE, partitioned, catalog-registered, meta sync not enabled
## What was not measured
- MERGE_ON_READ
- Whether an `UPSERT` or `DELETE` through the DataFrame writer behaves the
same as the `INSERT`
tested here
- Whether `recacheByPath` alone is sufficient when the base path is written
through a symlink or
a differently-qualified URI
--
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]