yihua opened a new pull request, #20088:
URL: https://github.com/apache/hudi/pull/20088

   ### Describe the issue this Pull Request addresses
   
   closes #20084
   part of #20064
   
   Stacked on #20069, review that first. It does not use the broadcast API from 
#20069, so it can also be rebased onto master directly.
   
   Spark serializes the upsert shuffle partitioner into every task on both 
sides of the shuffle, and the write and record creation functions into every 
task of their stages. `UpsertPartitioner` (and the bucket index partitioners) 
held the `HoodieTable`, the write config and the whole `WorkloadProfile`, 
although `getPartition` only needs the bucket maps, and looked up small files 
with one Spark task per partition path. The commit action executors captured by 
the write function kept their input records, so each write task shipped the 
lineage in front of the shuffle. `HoodieCreateRecordUtils` captured the 
DataFrame, the write config and two schemas that each task re-parsed.
   
   ### Summary and Changelog
   
   - Driver-only state in `SparkHoodiePartitioner`, `UpsertPartitioner` and the 
bucket index partitioners is `transient`; `getPartition` reads a partition's 
insert count from a small map.
   - Small files are looked up on the driver's file system view, only for 
partitions with inserts, after loading those partitions in bulk; the insert 
overwrite partitioner skips the lookup.
   - `RemoteHoodieTableFileSystemView#loadPartitions` splits the partition list 
over several requests so each URL stays below the timeline server's header 
limit (a few hundred partitions failed with 414 URI Too Long).
   - The MOR commit executor keeps the small file ids instead of the 
partitioner; input record fields of the Spark commit action executors are 
`transient`.
   - `HoodieCreateRecordUtils` builds the key generator properties, delete 
context and flags on the driver, and the schema without meta fields once per 
task.
   - Tests: `TestSparkWriteStagePayload` and `TestHoodieCreateRecordUtils` 
check the serialized partitioner, write stage and record creation tasks (and 
routing after deserialization); 
`TestRemoteHoodieTableFileSystemView#testLoadManyPartitions` and 
`TestUpsertPartitioner#testSmallFilesOfManyPartitionsFromTimelineServer` load 
1000 and 500 partitions through the timeline server. The payload and remote 
view tests fail on master.
   
   For a small upsert (Java-serialized): partitioner 86.9 KB to 1.4 KB, shuffle 
map stage payload 96.2 KB to 11.0 KB, record creation function about 98 KB to 
23 KB.
   
   ### Impact
   
   Performance only: less per-task deserialization in upsert shuffle map, write 
and record creation tasks, one fewer Spark job per upsert with inserts, and one 
schema rebuild less per record for prepped writes. No config or format change.
   
   Release note: the protected `profile` and `table` of 
`SparkHoodiePartitioner` and `config`, `smallFiles` and `bucketInfoMap` of 
`UpsertPartitioner` are now driver only (null in Spark tasks). A custom 
partitioner that reads them in `getPartition` must copy what it needs into its 
own fields. In-tree partitioners do not read them in tasks.
   
   ### Risk Level
   
   low. Covered by the new serialization tests (including routing equivalence 
on a deserialized partitioner and MOR small file routing through real tasks) 
and the existing COW and MOR write, datasource and MERGE INTO / UPDATE / DELETE 
suites. The small-file lookup now reads the driver's view and timeline instead 
of a copy in each task.
   
   ### Documentation Update
   
   none
   
   ### Contributor's checklist
   
   - [ ] Read through [contributor's 
guide](https://hudi.apache.org/contribute/how-to-contribute)
   - [ ] Enough context is provided in the sections above
   - [ ] Adequate tests were added if applicable
   


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