yihua opened a new issue, #20064: URL: https://github.com/apache/hudi/issues/20064
Hudi 1.x Spark reads ship the table meta client (with a full copy of the Hadoop configuration and the loaded timeline) to every task, and do per-task work that 0.14 did not: parsing the default configuration, and reading the timeline or `.hoodie` on executors for older table versions. Similar per-task payloads exist on the write paths, table services, metadata table reads and writes, and in the Trino, Flink and Hive integrations. This grows the task binary and the per-task CPU, which adds up on scans with many small tasks. On a local micro-benchmark (one file per task), the planned changes bring the Spark read task binary from 147 KB to 28 KB and per-task deserialization CPU from 2.4 ms to 0.3 ms, back to 0.14.1 levels (0.29 ms). The work is split into the PRs below, opened in stages. - [ ] perf(spark): ship file group reader scan state once per executor - [ ] refactor(core): read file groups from a table state instead of a meta client - [ ] perf(spark,flink): resolve schema-on-read file schemas from a history loaded once per query - [ ] perf(spark): stop per-file Hadoop configuration parsing and copies on the parquet read path - [ ] test(spark): guard against executor .hoodie access and task payload growth on the read path - [ ] perf(trino): read file groups on workers from table state shipped on the handle - [ ] perf(flink,hive): read splits with the table state captured at planning - [ ] feat(core): add an engine-agnostic broadcast to HoodieEngineContext - [ ] perf(index): keep tables and timeline reloads out of index lookup tasks - [ ] perf(metadata): stop shipping the table metadata to metadata table lookup tasks - [ ] perf(metadata): keep table state and per-task metadata I/O out of metadata index record generation - [ ] perf(spark): trim the upsert partitioner, write stage and record creation task payloads - [ ] perf(core): keep the table out of table service tasks and fix the bootstrap listing configuration - [ ] perf(client): serialize nested write configs as differences from the write config props - [ ] fix(cdc): fix CDC reads of MERGE_ON_READ tables of table version 6 - [ ] perf(flink): initialize the table once per instant on write tasks -- 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]
