Skip to content

perf: copy off-heap strings directly into Arrow - #6281

Open
pingzh wants to merge 4 commits into
apache:mainfrom
pingzh:pingzh-offheap-string-copy
Open

pingzh wants to merge 4 commits into
apache:mainfrom
pingzh:pingzh-offheap-string-copy

Conversation

@pingzh

@pingzh pingzh commented Sep 27, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

Closes #5910.

Rationale for this change

An off-heap Spark UTF8String currently passes through getByteBuffer(), which allocates a JVM byte array and copies the payload before Arrow copies it into its own buffer. Copying directly into reserved Arrow storage removes that intermediate allocation while retaining independent ownership of the output bytes.

What changes are included in this PR?

  • In ordinary StringWriter, reserve the value range with setValueLengthSafe, obtain the current buffer address after any reallocation, copy with writeToMemory, and mark the value defined. Reject negative string lengths and check destination bounds before the raw copy; Arrow enforces the 32-bit offset limit with the same exception for heap and off-heap inputs. Heap-backed inputs retain the existing setter path.
  • Add byte-level regression coverage for source reuse, null/empty/multibyte/invalid UTF-8 values, heap slices, separate data and metadata growth, retained C Data exports, allocation-failure cleanup, negative lengths, and consistent offset-overflow errors across all three backings, including failures after metadata growth. Register the suite in Linux and macOS CI.
  • Add CometStringWriterBenchmark, comparing the old setter against the actual writer with explicit off-heap inputs and heap-backed controls, preallocated and growing destinations, allocation counters, and output validation.

How are these changes tested?

  • make core and full Maven compilation succeeded. Spotless, Scalastyle, and dev/ci/check-suites.py passed.
  • All 37 tests in CometStringWriterSuite and CometArrowStreamSuite pass in the configurations below. Semantic Scalafix checks pass with both Scala versions; Spotless and Scalastyle pass.
Spark Scala JDK Tests
3.4.3 2.12 17 37/37
4.0.4 2.13 21 37/37

Focused Maven command (select the corresponding Spark profile):

./mvnw test -Pspark-4.1 -Dtest=none -Dsuites=org.apache.spark.sql.comet.execution.arrow.CometStringWriterSuite,org.apache.spark.sql.comet.execution.arrow.CometArrowStreamSuite

The full benchmark measurements were collected on commit 29da7e156 with Spark 4.1.3 / Temurin JDK 17.0.20.1, Linux x86-64, AMD EPYC 9V74. Each case uses a 1-second warmup and at least 1 second of measured writing. Source construction and initial destination allocation are outside the measured region; growth during writing is included.

Preallocated off-heap destinations:

Value size Baseline JVM bytes/row Direct copy JVM bytes/row Conversion throughput vs baseline
16 B 32 0 1.09x
1 KiB 1,096 0 1.82x
64 KiB 65,608 0 3.16x
1 MiB 1,048,648 0 1.58x
Varied (0 B–1 MiB) 209,163.5 0 3.47x

Growing off-heap destinations retain 0.8–81 JVM bytes/row of Arrow growth bookkeeping, while still removing the payload-sized array. Heap-array and heap-slice controls have identical allocation counts and throughput within approximately 2% of baseline across the measured cases. These are conversion microbenchmark results, not end-to-end query speedups, and the JVM allocation counter excludes Arrow's native storage.

A separate JDK 21 JFR allocation run (--smoke) attributes sampled baseline byte arrays to UTF8String.getBytes → getByteBuffer → ByteBufferStringWriter.setValue; the direct writer has no sampled payload-array allocation stack. The thread-allocation counters above provide the per-row measurement.

Reproduce with:

make benchmark-org.apache.spark.sql.comet.execution.arrow.CometStringWriterBenchmark

@github-actions github-actions Bot added enhancement New feature or request performance labels Sep 27, 2026

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

Summary

  • Prior state and problem: Writing an off-heap Spark UTF8String through getByteBuffer() allocates an intermediate JVM byte array before Arrow copies the bytes into vector storage.
  • Design approach: StringWriter reserves the ordinary Arrow value range, obtains the data buffer address after reservation, copies bytes directly from the UTF8String source, then marks the value defined. Heap-backed values retain the existing setter.
  • Correctness / compatibility analysis: At commit 5e0bd66 against base a12fd59, the reservation follows Arrow's existing setSafe capacity, gap, and offset behavior. The inspected row and columnar callers keep the source alive during the synchronous copy; the destination is a distinct Arrow allocation. Raw UTF8String bytes are preserved without decoding. I found no introduced or materially worsened P1/P2 issue.
  • Key design decisions: Reading the destination address after reservation handles buffer growth. The checked bounds and offset arithmetic guard the raw copy. Keeping the heap path unchanged limits the new behavior to the off-heap case.
  • Implementation sketch: Reserve with setValueLengthSafe, compute the current destination address, copy with writeToMemory, and set validity only after the copy.
  • Behavioral changes worth calling out: Off-heap values avoid the intermediate JVM array. Expected value bytes, null handling, and Arrow ownership remain the same; performance gains have not been independently measured in this review.
  • Suggested improvements: None required for approval. The added source tests cover source reuse, null and empty values, arbitrary UTF-8 bytes, vector growth, retained exports, and reservation failure.

Validation: Immutable-source review of the exact PR head, relevant Arrow Java and Spark UTF8String contracts, callers, and test source. I did not run tests or builds or inspect CI logs or artifacts.

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

Summary

  • Prior state and problem: Off-heap UTF8String.getByteBuffer() allocates an intermediate JVM byte array before Arrow copies the payload again.
  • Design approach: StringWriter copies off-heap strings directly into reserved Arrow storage. Heap-backed inputs retain the existing setter.
  • Correctness / compatibility analysis: Compared Arrow 18.3.0 and Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. The new path preserves bytes, null gaps, offsets and independent output ownership. Inspected callers keep source storage alive through the synchronous copy.
  • Key design decisions: Reading the destination address after reservation handles reallocation. Checked offset arithmetic and destination bounds protect the copy. The implementation reuses existing Arrow APIs without adding an ownership abstraction.
  • Implementation sketch: Reserve with setValueLengthSafe, copy with writeToMemory, then call setIndexDefined. The PR adds regression tests, registers them in both workflows and supplies a comparative benchmark.
  • Behavioral changes worth calling out: The off-heap path removes payload-sized JVM allocations. Local benchmark smoke measurements confirmed this for 16-byte and 1-MiB inputs, with matching allocation counts for heap-backed controls. Arrow storage and memory accounting remain unchanged.
  • Suggested improvements: None warranted at P1/P2 priority. No introduced P1/P2 issues found within this review.

Reviewed both commits and all five files in the full PR diff at 5e0bd66fbbcc2648da49869af3bf19ec071a38e1, against base a12fd5902f1012a9e493739ecacf7dd8e22ca78a. The PR remains non-draft. Read the existing approval and confirmed there are no issue comments, inline comments or review threads. Routed skills: review-comet-pr, review-comet-ffi-pr and review-comet-memory-pr.

Exact-head CI: 23 checks passed and 14 were skipped, with no failures or unfinished checks. Linux execution logs confirm all eight new tests passed. CI ran merge commit c03662d8f43930db7a562805281ab5e408b8a4f8, whose five PR-changed files match the reviewed head. macOS, Spark SQL, Iceberg and benchmark checks were skipped.

Local validation: The exact writer source compiled and passed a bounded Spark 4.1.3 / JDK 21 probe covering 25,686 value, offset and validity checks, reset, columnar inputs, 20 retained FFI batches, allocation failure and offset guards. Benchmark compilation and all smoke cases passed. Suite registration and diff checks passed.

Validation limits: Local compilation used small schema/project-root adapters instead of a full Maven/native build. Other Spark versions received source review, not local runtime testing. Smoke measurements do not establish sustained throughput or end-to-end query gains.

@andygrove andygrove 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.

Thanks for this. The new path does the same Arrow bookkeeping as setSafe, and the tests and benchmark cover everything #5910 asked for. One question inline about the overflow pre-check.

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

Summary

  • Prior state and problem: Off-heap Spark UTF8String.getByteBuffer() allocates a temporary JVM byte array before Arrow copies the same bytes into its own storage. Large strings therefore incur an avoidable payload allocation and copy at conversion boundaries.
  • Design approach: Ordinary StringWriter reserves the value range in its VarCharVector and copies off-heap bytes directly into that Arrow-owned range. Heap-backed strings retain the existing setter.
  • Correctness / compatibility analysis: Reviewed head f780275f855c6861687a4ba4f59107dc24b5e8bb against comparison base a12fd5902f1012a9e493739ecacf7dd8e22ca78a, including Arrow Java 18.3.0 and Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3, and 4.2.0 source contracts. Reservation preserves Arrow's null-gap and offset bookkeeping. The inspected callers keep source storage available during the synchronous copy, and output retains independent ownership. Raw bytes and downstream invalid-UTF-8 handling are preserved. No introduced P1/P2 issue was found.
  • Key design decisions: Retrieve the destination address after reservation can replace the buffer, check the destination range, copy before setting validity, and delegate offset-size limits to Arrow. The current head fixes the prior review's overflow-error inconsistency while keeping the negative-length guard.
  • Implementation sketch: Use setValueLengthSafe, obtain the current buffer and start offset, call writeToMemory, then call setIndexDefined. The PR adds a ten-test regression suite, registers it in both workflows, and adds a benchmark comparing the actual writer with the base setter across off-heap and heap controls.
  • Behavioral changes worth calling out: Off-heap ordinary strings avoid the intermediate payload-sized JVM array. Arrow still owns a copy, and its allocation accounting and batch lifetimes are unchanged. Overflow uses Arrow's OversizedAllocationException for all three tested source backings. The benchmark measures conversion and JVM allocation, not whole-query speedups.
  • Suggested improvements: None warranted at P1/P2 priority within this source review.

Validation: Source inspection only. Check metadata associated with the reviewed head reports 24 successful and 14 skipped checks. No tests, builds, runtime probes, benchmarks, package fetching, CI logs, or CI artifacts were run or inspected. Test results and performance figures in the PR description and older reviews were not independently reproduced in this task.

@andygrove

Copy link
Copy Markdown
Member

This is a light fully automated review since there are so many PRs open.

benchmarks/micro/run.py discovers every benchmark object in spark/src/test/scala/org/apache/spark/sql/benchmark/ (discover_suites, run.py:469) and launches each one as -Dexec.mainClass=org.apache.spark.sql.benchmark.<name> (suite_command, run.py:545). CometStringWriterBenchmark sits in that directory but declares package org.apache.spark.sql.comet.execution.arrow (spark/src/test/scala/org/apache/spark/sql/benchmark/CometStringWriterBenchmark.scala:20). So a default run.py run will pick it up and fail it with a ClassNotFoundException for org.apache.spark.sql.benchmark.CometStringWriterBenchmark. Naming the real class in a --suites file doesn't get around it either, because load_suites keeps only the simple name and suite_command puts the fixed package back on. CometArrowWriterBenchmark in the same directory already has this problem. Could this file move to spark/src/test/scala/org/apache/spark/sql/comet/execution/arrow/ next to CometStringWriterSuite? The make benchmark-... command in the docstring would still work, and the runner would stop picking it up. If you'd rather have it in the runner's results, would it make sense to have discover_suites read each file's package line and pass the full class name to suite_command? That would fix CometArrowWriterBenchmark at the same time.

@ajsquared ajsquared left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Reviewed 3e8c980: no actionable findings. The off-heap copy preserves Arrow reservation, offset and validity bookkeeping and independent buffer ownership. The earlier overflow-error and benchmark-discovery concerns are addressed.

Remote source review only; tests were not run, and CI and production were not inspected.

@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 3e8c980c7760d85339b25525aa4cbe853adeb64a against a12fd5902f1012a9e493739ecacf7dd8e22ca78a, including the complete PR diff. The sole change since f780275 moves the unchanged benchmark into its declared package directory. This addresses the benchmark discovery comment: the generic runner no longer discovers it under the wrong class name, and the documented Makefile command remains valid. The earlier overflow-error fix is retained. No introduced P1/P2 issue found.

Validation was source-only. No tests, builds, benchmarks, runtime probes, package fetching, CI logs, or artifacts were run or inspected. CI check metadata remained in progress at the final refresh.

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

Labels

enhancement New feature or request performance

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Avoid temporary heap allocation when copying off-heap UTF8String into Arrow

4 participants