Conversation
sunchao
left a comment
There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Off-heap
UTF8String.getByteBuffer()allocates an intermediate JVM byte array before Arrow copies the payload again. - Design approach:
StringWritercopies 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 withwriteToMemory, then callsetIndexDefined. 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.
sunchao
left a comment
There was a problem hiding this comment.
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
StringWriterreserves the value range in itsVarCharVectorand copies off-heap bytes directly into that Arrow-owned range. Heap-backed strings retain the existing setter. - Correctness / compatibility analysis: Reviewed head
f780275f855c6861687a4ba4f59107dc24b5e8bbagainst comparison basea12fd5902f1012a9e493739ecacf7dd8e22ca78a, 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, callwriteToMemory, then callsetIndexDefined. 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
OversizedAllocationExceptionfor 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.
|
This is a light fully automated review since there are so many PRs open.
|
ajsquared
left a comment
There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
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.
Which issue does this PR close?
Closes #5910.
Rationale for this change
An off-heap Spark
UTF8Stringcurrently passes throughgetByteBuffer(), 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?
StringWriter, reserve the value range withsetValueLengthSafe, obtain the current buffer address after any reallocation, copy withwriteToMemory, 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.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 coreand full Maven compilation succeeded. Spotless, Scalastyle, anddev/ci/check-suites.pypassed.CometStringWriterSuiteandCometArrowStreamSuitepass in the configurations below. Semantic Scalafix checks pass with both Scala versions; Spotless and Scalastyle pass.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.CometArrowStreamSuiteThe full benchmark measurements were collected on commit
29da7e156with 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:
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 toUTF8String.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: