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]

Reply via email to