yihua commented on PR #20078: URL: https://github.com/apache/hudi/pull/20078#issuecomment-6027948918
## Spark read benchmark of the landed changes Recording where the Spark read path stands now that this PR has landed. "landed master" below is apache/hudi master at [438c0d3](https://github.com/apache/hudi/commit/438c0d372569bb40138cc3046069ed405f75bdbe), which includes this PR together with #20077, #20079, #20080 and #20115. "master before" is 463b6159, before any of them. **Setup.** Spark 3.4.3 bundles, `local[4]`, JDK 11, 8 GB heap, metadata table off. Each table has 400 file groups x 300 rows in 10 partitions, 100 top-level columns (82 primitives plus structs, arrays of structs and maps of structs nested to depth 6), and 1,000 completed instants on the active timeline before the data commit, each with 100 write stats and the schema in its commit metadata. MOR tables add one upsert of every 10th row (one log file per file group). The session Hadoop conf carries 1,000 extra entries on top of Spark's ~1,000 defaults. Every query runs with one file slice per Spark task (400 tasks). Values are per-task medians of 10 measured iterations after 5 warmups (0.14.1: 6 after 3), averaged over 2 rounds that alternate the builds. Task binary is the serialized (lz4) RDD and dependency that Spark broadcasts per stage and every task deserializes; deserialize and executor CPU are Spark's `executorDeserializeCpuTime` and `executorCpuTime`. Incremental queries read the bulk insert commit on COW and the upsert delta commit on MOR. The narrow snapshot projects 4 columns, one of them a nested leaf. **Each version reading the table it wrote** (0.14.1: table version 6, 1.2.1: version 9, master: version 10) | metric | 0.14.1 | 1.2.1 | master before | landed master | |---|---|---|---|---| | task binary, COW snapshot (kB) | 43.9 | 351.1 | 353.5 | 47.8 | | Hudi read function in the task binary (kB) | 10.8 | 318.1 | 320.6 | 1.2 | | deserialize CPU per task, COW snapshot (ms) | 0.73 | 5.51 | 5.33 | 0.74 | | deserialize CPU per task, MOR snapshot (ms) | 0.96 | 5.56 | 5.31 | 0.78 | | executor CPU per task, COW snapshot (ms) | 43.01 | 38.03 | 36.67 | 32.71 | | executor CPU per task, COW incremental (ms) | 130.44 | 40.38 | 38.33 | 35.06 | | driver time before the scan job is submitted, COW snapshot (ms) | 2,408 | 414 | 404 | 360 | **The 0.14.1-written table (version 6) read by every version** | metric | 0.14.1 | 1.2.1 | master before | landed master | |---|---|---|---|---| | executor `.hoodie` calls per MOR query | 2,400 | 400 | 400 | 0 | | deserialize CPU per task, MOR snapshot (ms) | 0.96 | 5.48 | 5.42 | 0.81 | | executor CPU per task, MOR snapshot (ms) | 189.68 | 162.36 | 164.89 | 49.19 | | executor CPU per task, MOR incremental (ms) | 150.70 | 134.72 | 138.90 | 21.54 | | executor CPU per task, narrow MOR snapshot (ms) | 131.46 | 125.53 | 126.18 | 7.08 | | executor CPU per task, COW snapshot (ms) | 43.01 | 36.44 | 38.16 | 33.24 | **The 1.2.1-written table (version 9) read by every 1.x version** | metric | 1.2.1 | master before | landed master | |---|---|---|---| | task binary, COW / MOR snapshot (kB) | 351.1 / 353.1 | 353.5 / 355.4 | 47.8 / 49.6 | | deserialize CPU per task, COW snapshot (ms) | 5.51 | 5.55 | 0.73 | | executor CPU per task, COW snapshot (ms) | 38.03 | 37.44 | 33.34 | | executor CPU per task, MOR snapshot (ms) | 50.27 | 48.67 | 49.20 | **Takeaways** - The per-task payload growth in 1.x is gone. Task binary goes from 350-358 kB to 47-51 kB (17 kB for the narrow projection), of which the Hudi read function is 1.2 kB and the rest is Spark's own plan. Deserialize CPU per task goes from 5.2-5.8 ms to 0.7-0.8 ms, back to the 0.14.1 level. This holds for snapshot, read optimized and incremental queries on table versions 6, 9 and 10. - Executors no longer touch `.hoodie`. On the version 6 MOR table, 1.2.1 and master before listed the timeline once per task (400 calls per query); landed master makes none, and MOR executor CPU per task drops 3.4x (snapshot), 6.4x (incremental) and 18x (narrow snapshot). - The shipped state moved into per-executor broadcasts: 315-337 kB per executor (the 191 kB Hadoop conf broadcast 1.2.1 already had, a 123 kB scan state, and 20 kB of committed instants on version 6 MOR tables), shipped once instead of ~300 kB extra in every task. - Query results (sum of row hashes and row count) are identical between master before and landed master on every query and table. - The two rounds differ by up to ~15% on the smaller MOR numbers (8-60 ms), so differences of that size among 1.x builds on table versions 9 and 10 are noise; the changes above are well outside it. -- 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]
