Skip to content

perf: presize collections and reuse buffers in native loops - #6861

Open
comphead wants to merge 1 commit into
apache:mainfrom
comphead:rust-alloc-capacity-reuse
Open

comphead wants to merge 1 commit into
apache:mainfrom
comphead:rust-alloc-capacity-reuse

Conversation

@comphead

@comphead comphead commented Oct 10, 2026 •

Copy link
Copy Markdown
Contributor

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

  • A collection is created empty (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.
  • A loop over rows or batches allocates a fresh Vec, String or 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_batch built a ScalarValue per row, which allocates a String for Utf8 input. String arrays are now hashed in place.
  • split.rs: with limit = 0, each row's parts were collected into a Vec so 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_partition built a ShuffleBlockWriter, 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_cols and the column merge in build_batch allocated a one-element Vec per column before flattening. Both now append into one vector of known capacity.
  • multi_partition.rs (hash-all round robin), schema_align.rs and parquet_writer.rs borrow the batch's columns, or move them out with into_parts(), instead of copying them into a new Vec.
  • spark_bloom_filter.rs: serialization cloned the bit array and then copied it twice more. It now writes each word into one presized buffer, and SparkBitArray::data() returns a slice.
  • Smaller ones: presized output vectors in expand.rs and SparkPercentileGroupsAccumulator::state, an all-true MERGE mask built without a temporary Vec<bool> in merge_rows.rs, and no class_name clone per batch in JvmScalarUdfExpr.

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 in process_sorted_row_partition, explode's per-batch maps, a loop-invariant Schema::empty() hoisted in planner.rs, SparkPlan::new_with_additional moving its Vec instead of copying it, and a reused line buffer in analyze_trace.

Left alone on purpose:

  • The RSS writer's per-frame buffers, which are sized and released inside each frame's memory reservation.
  • Reusing arrow Rows across batches in range partitioning and RankLimitStream, which would keep memory that the pools do not track alive between batches.
  • The per-call vector of completed batches in 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, CometNativeShuffleSuite and CometGenerateExecSuite. The only test change is an import that multi_partition.rs's test module now needs.

cargo clippy --all-targets --workspace -- -D warnings passes locally. The test suites were not run locally, and there are no benchmark numbers for these changes.

@github-actions github-actions Bot added enhancement New feature or request performance area:writer Native Parquet writer area:shuffle Shuffle (JVM and native) area:aggregation Hash aggregates, aggregate expressions area:scan Parquet scan / data reading area:expressions Expression evaluation area:ffi Arrow FFI / JNI boundary area:Iceberg area:udf labels Oct 10, 2026
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.
@comphead
comphead force-pushed the rust-alloc-capacity-reuse branch from d60b157 to d156eff Compare October 10, 2026 19:54

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.1 at e9efd9f764ee0a59b7898ff028d6985d4a7a28e1 identified 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.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:aggregation Hash aggregates, aggregate expressions area:expressions Expression evaluation area:ffi Arrow FFI / JNI boundary area:Iceberg area:scan Parquet scan / data reading area:shuffle Shuffle (JVM and native) area:udf area:writer Native Parquet writer enhancement New feature or request performance

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants