linliu-code opened a new pull request, #20120:
URL: https://github.com/apache/hudi/pull/20120

   ### Describe the issue this Pull Request addresses
   
   Closes #20113
   
   A DataFrame write to a Hudi table does not invalidate Spark session cache 
entries built from that table, so a cached Hudi table keeps serving rows from 
before the commit, with no signal to the reader. A plain Parquet table under 
the identical cycle does not do this.
   
   Measured, same session, same harness, meta sync off and neither table 
registered in the catalog:
   
   |  | rows before write | rows after write | stale |
   | --- | --- | --- | --- |
   | Hudi | 3 | 3 | yes |
   | Parquet | 3 | 4 | no |
   
   `REFRESH TABLE` recovers both, so the data on storage is correct; only the 
cache is wrong.
   
   Spark's own file-based writers invalidate the cache from 
`InsertIntoHadoopFsRelationCommand`, which calls `CacheManager.recacheByPath` 
on the output path. Hudi writes go through `DefaultSource.createRelation` -> 
`HoodieSparkSqlWriter` and never reach that command. The only invalidation 
`HoodieSparkSqlWriter` performed was a catalog `refreshTable`, and only when 
`hoodie.meta.sync.enable` is on — which cannot reach an entry built from a 
DataFrame that was never registered in the catalog.
   
   ### Summary and Changelog
   
   - `HoodieSparkSqlWriter.metaSync` now calls 
`spark.catalog.refreshByPath(basePath)` after every commit. That is the public 
entry point to `CacheManager.recacheByPath`, the same call Spark makes for its 
own writers. Matching by path rather than by plan is what reaches non-catalog 
entries.
   - The call is wrapped in a `NonFatal` guard that logs a warning. The 
rationale is in the code comment and worth stating here, since swallowing an 
exception on a write path normally deserves a rethrow: the commit has already 
succeeded by this point, and `CacheManager.recacheByCondition` removes the 
matching entries from `cachedData` and clears their blocks *before* it attempts 
to rebuild them. So the invalidation has already taken effect when the rebuild 
throws — rethrowing would report a failure for a write that succeeded, while 
gaining nothing. The rebuild does throw in practice: an overwrite that replaces 
a partitioned table with a non-partitioned one leaves the cached plan holding 
the old partition schema, and re-optimizing it fails on the new layout with 
`Empty partition column value in 'partition='`. That case is covered by a test.
   - New `TestHoodieWriteCacheInvalidation` with four cases: the Hudi write 
parameterized over COPY_ON_WRITE and MERGE_ON_READ, the layout-changing 
overwrite above, and a Parquet control so a failure of the Hudi cases can be 
attributed to Hudi rather than to the harness or to Spark's caching.
   
   Verification: all four tests fail-then-pass across the change. Without the 
fix, `testCachedDataFrameSeesHudiWrite` fails `expected: <4> but was: <3>` 
while the Parquet control passes on the same base. `TestCOWDataSource` + 
`TestMORDataSource` run 234 tests with 0 failures and 0 errors on the final 
diff; the unguarded first draft of this change broke 
`TestCOWDataSource.testInferPartitionBy`, which is the case the guard and its 
test now cover.
   
   ### Impact
   
   Behavior change on the write path: a write now invalidates and rebuilds 
Spark cache entries rooted at the table base path, matching what Spark already 
does for its own file-based writers. Reads that were silently stale now see 
committed data.
   
   Cost is one `recacheByPath` per commit. With an empty session cache it is a 
filter over an empty list. With a matching entry it discards and re-registers 
the `InMemoryRelation`, which is lazy — nothing is materialized until the next 
query. This is the same per-write cost Parquet already pays.
   
   Not addressed here: a writer outside the reading Spark session — a separate 
job, Flink, a streaming ingest — still leaves a cached Hudi table stale, 
because Spark's `CacheManager` is session-local and no in-process hook exists 
for it.
   
   ### Risk Level
   
   low
   
   The added call is a no-op when nothing is cached under the table's base 
path, and it cannot fail the write. The behavior it changes is behavior Spark's 
own writers already have. The regression risk that did materialize — rebuilding 
a cached plan that is no longer valid against the table — was found by running 
the existing COW and MOR datasource suites against the change, and is handled 
by the guard plus a dedicated test.
   
   ### Documentation Update
   
   none
   
   ### Contributor's checklist
   
   - [x] Read through [contributor's 
guide](https://hudi.apache.org/contribute/how-to-contribute)
   - [x] Enough context is provided in the sections above
   - [x] Adequate tests were added if applicable
   
   🤖 Generated with [Claude Code](https://claude.com/claude-code)
   


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