Repository navigation
Conversation
Give collections a capacity where their final size is known, and stop allocating a fresh structure in every iteration of per-row and per-batch loops where one can be reused, borrowed or moved instead.
d60b157 to
d156eff
Compare
sunchao
left a comment
There was a problem hiding this comment.
Reviewed full head d156eff7b60e700dbd6dd680688688bc3f07210b against supplied base c1e00e7b92b0ada119b140b016dbefdd198bbc35. The PR is not a draft. No introduced P1/P2 issues found within this review.
Reviewed all 20 PR-changed files and the three additional direct base/head differences. Commit ancestry confirms those additional differences come from base-only fix #6310, not deletions introduced by this PR. The snapshot and live reviews, conversation comments, inline comments, and review threads contained no existing feedback.
Summary
- Prior state and problem: Native loops created temporary collections, cloned strings and column vectors, and repeatedly constructed shuffle writers even when sizes or schemas were already known.
- Design approach: Presize collections, borrow or move existing data, and reuse scratch buffers and schema-dependent writers within their existing execution scope.
- Correctness: Checked surrounding code and relevant Spark implementations for split, Bloom filters, JSON array length, percentile, and generators. String hashing still skips nulls, Bloom serialization preserves V1/V2 headers and big-endian words, and flattening preserves column order. Focused split and MERGE-mask comparisons passed.
- Compatibility analysis: Compared relevant Spark sources across supported versions 3.4.3, 3.5.9, 4.0.4, 4.1.3, and 4.2.0. No configuration, shim, protobuf, or shuffle-format contract changes were introduced. The native split limit-zero difference is inherited, with Spark codegen still used by default. Comparison with
branch-1.1ate9efd9f764ee0a59b7898ff028d6985d4a7a28e1identified no new result changes attributable to these allocation edits. - Key design decisions: The cached shuffle writer is rebuilt when the batch schema changes. Dictionary encoding still creates a fresh stream writer per block. Nested-struct scratch vectors are cleared before each sibling field.
- Implementation sketch: Changes cover expression allocation, expand/explode/MERGE construction, Parquet schema matching and batch assembly, shuffle conversion, plan construction, and trace-line reading.
- Performance: The code removes identifiable allocations and copies, including per-string scalar materialization and repeated schema encoding. No benchmarks were supplied or run, so throughput improvements and workload-wide performance neutrality remain unmeasured. No reproducible P1/P2 performance regression was identified.
- Design: The changes remain local and preserve existing ownership and error-handling paths. Reuse is bounded by a call or partition, without introducing process-wide caches.
- Abstraction & complexity: Existing slices,
Vec,into_parts(), and Arrow builders are sufficient. The schema/writer pair adds limited state with an explicit invalidation condition. No abstraction issue meeting the reporting bar was identified. - Behavioral changes worth calling out: The native edits preserve intended results and serialization. Restoring the base's memory-manager tests produced three failures at the pinned head and none at the base, but ancestry establishes that this is the already-fixed race from base-only #6310. It is not a PR regression.
- Suggested improvements: No introduced P1/P2 issue requiring a change was identified.
Routed skills: review-comet-pr, review-comet-expression-pr, review-comet-memory-pr, review-comet-shuffle-pr, review-comet-ffi-pr, and review-comet-iceberg-write-pr.
Exact-head CI: 26 successful checks and 15 skipped checks, with no failures or pending checks. Required Checks, all four Linux Spark 4.1 test groups, native tests, and TPC-H/TPC-DS verification passed. The Rust log reports 2,224 passing tests plus a separate five-test invocation passing.
Local validation: Extracted base/head split functions matched across 122,850 UTF-8/LargeUTF-8 cases. MERGE Boolean masks matched for 1,025 lengths. The isolated memory-manager suite passed 6/9 tests at head and 9/9 at base, explained by the base-only fix above.
Validation limits: No full local Comet build, end-to-end Spark/Iceberg run, or benchmark was performed. Upstream Spark SQL, Iceberg, PyArrow UDF, and macOS CI suites were skipped. No project source or GitHub state was changed.
Which issue does this PR close?
There is no issue for this change. It is a sweep of the native Rust code for two allocation patterns.
Rationale for this change
Two patterns allocate more than needed across
native/:Vec::new(),HashMap::new(),String::new(), an Arrow builder) and then filled to a size that is already known, so it reallocates as it grows.Vec,Stringor writer in every iteration, where one could be reused, or the data borrowed or moved instead.This PR fixes the instances in production code where the bound is exact or the reuse is local.
What changes are included in this PR?
Allocations removed per row:
bloom_filter_agg.rs:update_batchbuilt aScalarValueper row, which allocates aStringforUtf8input. String arrays are now hashed in place.split.rs: withlimit = 0, each row's parts were collected into aVecso that trailing empty parts could be dropped. Empty parts are now held back until a non-empty part follows.Allocations removed per batch:
spark_unsafe/row.rs:process_sorted_row_partitionbuilt aShuffleBlockWriter, which pre-encodes the IPC schema, for every batch. The writer is now kept until the batch schema changes, which happens when a dictionary-encoded column falls back to a plain array.spark_unsafe/row.rs: the field-major struct conversion allocated three vectors for every nested struct field. They are now allocated once per call and cleared for each field.explode.rs:flatten_struct_colsand the column merge inbuild_batchallocated a one-elementVecper column before flattening. Both now append into one vector of known capacity.multi_partition.rs(hash-all round robin),schema_align.rsandparquet_writer.rsborrow the batch's columns, or move them out withinto_parts(), instead of copying them into a newVec.spark_bloom_filter.rs: serialization cloned the bit array and then copied it twice more. It now writes each word into one presized buffer, andSparkBitArray::data()returns a slice.expand.rsandSparkPercentileGroupsAccumulator::state, an all-true MERGE mask built without a temporaryVec<bool>inmerge_rows.rs, and noclass_nameclone per batch inJvmScalarUdfExpr.Cold paths, changed only where the bound is exact: presized maps and sets in parquet schema matching (
parquet_support.rs,schema_adapter.rs),iceberg_dictionary.rs, the builders inprocess_sorted_row_partition, explode's per-batch maps, a loop-invariantSchema::empty()hoisted inplanner.rs,SparkPlan::new_with_additionalmoving itsVecinstead of copying it, and a reused line buffer inanalyze_trace.Left alone on purpose:
Rowsacross batches in range partitioning andRankLimitStream, which would keep memory that the pools do not track alive between batches.BufBatchWriter::write. Removing it changes which batches are dropped when a write fails.How are these changes tested?
None of the changes alter output, so they are covered by existing tests, for example
test_split_limit_zero, the bloom filter round-trip tests,CometColumnarShuffleSuite,CometNativeShuffleSuiteandCometGenerateExecSuite. The only test change is an import thatmulti_partition.rs's test module now needs.cargo clippy --all-targets --workspace -- -D warningspasses locally. The test suites were not run locally, and there are no benchmark numbers for these changes.