voonhous commented on issue #20064:
URL: https://github.com/apache/hudi/issues/20064#issuecomment-5979735678
### Before: every Spark task carries the file format and a meta client
```mermaid
flowchart TB
subgraph DRV["Driver"]
FI["HoodieFileIndex<br/>plans the scan with the relation's meta client"]
FF["File format builds a second meta client<br/>with a full Hadoop conf"]
CL["Per-file read closure<br/>captures the file format + meta client"]
FI --> FF --> CL
end
subgraph EX["Executor: repeated in every task"]
T1["Java-deserialize file format,<br/>meta client, Hadoop conf"]
T2["Re-parse schemas"]
T3["Write read options into<br/>the live table config"]
T4["Load timeline<br/>log block commit check, table version before 8"]
T5["Load schema history<br/>schema-on-read"]
T6["Read file group"]
T1 --> T2 --> T3 --> T6
T4 --> T6
T5 --> T6
end
HD[(".hoodie on storage")]
CL -- "shipped with every task" --> T1
HD -. "read per task" .-> T4
HD -. "read per task" .-> T5
```
Problems: large task binary and per-task deserialization CPU; `.hoodie` I/O
on executors; tasks of one query can see different timelines; a read option
named like a table config key can change merge results, depending on which file
a task read first.
### After: scan state built once on the driver, shipped once per executor
```mermaid
flowchart TB
subgraph DRV["Driver: once per scan"]
MC["Relation's meta client<br/>timeline the scan was planned against"]
ST["Scan state #20077<br/>schemas + reader props<br/>read options merged
into a copy of table props"]
TS["FileGroupReaderTableState #20078<br/>table config, base
path,<br/>committed instants, schema history"]
BC["Spark broadcast"]
FN["HoodieFileGroupReaderFunction #20077<br/>holds broadcast handles
only"]
MC --> ST
MC --> TS
ST --> BC
TS --> BC
end
subgraph EX["Executor"]
DS["Deserialize scan state<br/>once per executor"]
subgraph TASKS["Each task"]
P["Copy per-file props"]
R["Read file group<br/>no .hoodie I/O"]
P --> R
end
DS -- "shared by concurrent tasks" --> P
end
BC -- "once per executor" --> DS
FN -- "small closure per task" --> P
```
Result: one consistent timeline snapshot per scan, no executor `.hoodie`
reads, and read options never mutate shared table config. #20086 (Trino) and
#20087 (Flink, Hive) apply the same pattern: capture the table state where the
read is planned and ship it on the handle or split. #20112 adds tests that fail
if a task touches `.hoodie` or the task binary grows past 24 KB.
--
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]