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]
