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]

Reply via email to