Skip to content

perf: convert Spark's unsafe rows to Arrow straight from their memory - #6730

Merged
andygrove merged 6 commits into
apache:mainfrom
andygrove:perf/unsafe-row-to-arrow
Oct 8, 2026
Merged

andygrove merged 6 commits into
apache:mainfrom
andygrove:perf/unsafe-row-to-arrow

Conversation

@andygrove

@andygrove andygrove commented Oct 6, 2026 •

Copy link
Copy Markdown
Member

Which issue does this PR close?

Closes #6721.

Rationale for this change

ArrowWriter's row path serves CometSparkRowToColumnar, CometLocalTableScanExec and the Comet cache serializer's row input. Since #6607, CometSparkRowToColumnar also converts a shuffle's row input for native shuffle when spark.comet.convert.shuffleInput.enabled is set, and #5634 adds more callers. Nearly all of those rows are UnsafeRows, but the writer read them through Spark's generic getters. A string went through a UTF8String and Arrow's per-value setter, and every element of an array, map or struct took the generic write, which makes two or three virtual or interface calls per value. #6721 has the per-type comparison with the native row converter.

Per column type, in ns per row with one column of 8192 rows and every tenth row and element null, in the shapes CometArrowWriterBenchmark uses. A harness outside the repo ran main's writer and this one in alternating JVMs, three times each for two orders of the types. It warms every type up before measuring, so the JIT sees a mix of types as a real workload does. Each number is the median of the six runs' best.

Row input main This PR
int 6.1 5.5
long, double, date, boolean 6.2-6.5 5.5-6.4
decimal(18,2) 8.8 6.8
decimal(38,10) 17.5 10.2
string, 11 bytes 17.6 11.7
string, 48 bytes 17.9 13.7
binary 17.2 12.7
struct<int,long,double,date> 30.6 28.8
struct<int,string> 32.2 23.9
array 24.7 17.6
array 25.7 21.8
array<decimal(18,2)> 32.5 22.9
array 56.4 30.9
array<struct<int,string>> 86.1 65.8
array<array> 69.3 49.1
array<map<string,int>> 202.0 147.5
map<string,string> 114.4 55.1
map<int,long> 45.7 32.9
8 columns: int, long, double, string, date, decimal(18,2), boolean, timestamp 67.0 53.3

On array<map<string,int>>, a C2 inlining quirk costs this PR about a quarter of its time, see #6757.

CometArrowWriterBenchmark itself warms one type at a time, so the JIT sees fewer types at main's shared call sites, and the gains differ. They are smaller on strings, array<int> and the fixed-width rows, and larger on array<string> and structs. Best of two alternating runs each:

Spark rows to Arrow by type main This PR
int 5.8 5.0
decimal(18,2) 8.4 6.6
decimal(38,10) 15.6 9.8
string 13.1 10.9
struct<int,long,double,date> 30.3 27.0
array 18.7 14.9
array 59.0 24.5
map<string,string> 98.6 52.8
int, long, double (the "Spark rows to Arrow" case) 16.2 14.9

End to end, count aggregates over a table of 4M rows cached in Spark's default format, with a release libcomet at local[1] on the Spark 4.1 profile. The cache's vectorized reader is off, so CometSparkToColumnarExec converts every column from rows. Every column has 10% nulls, the arrays hold 0 to 3 strings and the maps 0 to 2 entries. Each number is the median of three alternating runs, each the best of 7. conversionTime also counts decoding the cached rows.

Columns read Wall time (ms) conversionTime (ms)
int, long, double, decimal(18,2), date 213 → 187 177 → 153
string 119 → 104 92 → 77
struct<int,string> 206 → 135 181 → 111
array, map<string,string> 402 → 299 368 → 267
all nine 853 → 657 807 → 614

All of this is on an M3 Max with JDK 17, on a machine shared with other jobs, so take the numbers as indicative.

What changes are included in this PR?

ArrowWriter.write hands an UnsafeRow to a new writeUnsafeRowField, and unsafe arrays go to a new writeArrayElements. Other rows take the generic path as before. Both read Spark's memory directly:

  • Fixed-width values are read from their slot with UnsafeRow's own getters and written at the vector's width, without the virtual call to setValueUnsafe. Arrays of these types are copied in one block, and their validity is the array's null words inverted, written at any bit offset. Booleans are converted from a byte to a bit per element.
  • Decimals up to 18 digits are written from the unscaled long, which UnsafeRow.getDecimal reads without checking it against the precision either. In an unsafe array the long is range-checked first, as UnsafeArrayData.getDecimal checks it, so a value past the precision still fails as before. Wider decimals are range-checked in place instead of being copied to a byte[] first.
  • Strings and binaries are copied using the offset and size in their slot, without a UTF8String or setSafe. An array's elements are all checked first, then the buffers are grown once, through the reserveValues that the columnar string paths now share, and the bytes copied. Values up to 64 bytes are copied a word at a time, which beats Unsafe.copyMemory's checks and call on short strings.
  • Arrays, maps and structs read each value through Spark's getters, which allocate a view of it as before, and write its elements, keys and values, or fields through these same paths. One view reused across values would save another 1-4 ns per row on nested types, but it keeps the last row it read reachable until the writer is dropped, so that is left to perf: reuse the nested unsafe views in ArrowWriter's row path, detaching them once per row #6765. A map's entries are marked valid in one call, and setValid now writes the bits around whole bytes with one masked write each, calling setOne (a JNI call to Unsafe.setMemory on JDK 17 and 21) only for runs of 64 bytes or more. A generic row holding an unsafe value, as a copied ColumnarBatchRow holds unsafe primitive arrays, takes these paths too.

ArrowFieldWriter.writeUnsafe is gone. Its fixed-width override skipped the capacity check and relied on ArrowWriter.create allocating top-level vectors at the batch size. Arrow 18 keeps a fixed-width vector's capacity in a field, so the new path checks it per value for nothing measurable. Struct fields can then use the same path as top-level fields, which addresses the issue's point about structs without allocating nested vectors at the batch size. FixedWidthArrowFieldWriter now takes its vector as a constructor parameter, so the per-value code reads it from a field instead of through each subclass's accessor.

The two failure modes that CometStringWriterSuite pins still hold on the new path, and more strictly than on Arrow's setters. An offset overflow throws Arrow's OversizedAllocationException and a negative length throws an IllegalArgumentException, before anything is written or allocated, for a single value and for every element of an array.

I also tried binding fixed-width fields statically from the row and struct loops. That gained another 12-16% on structs and wide rows in the harness, but it was a wash elsewhere, so I left it out.

How are these changes tested?

  • New CometArrowWriterSuite tests write unsafe rows, both on heap and copied to off-heap memory, and check that every vector matches what the generic row path writes for the same rows. They cover the suite's primitive types (18 on Spark 4.1), an array of each of them, the 11 nested types already there and 5 more (collections inside collections and inside a struct, a map whose decimal values are converted one at a time, and a struct of 70 fields), plus a row with a column of every type three times over. The arrays reach each of the twelve block copies of fixed-width elements, at widths of 1, 2, 4 and 8 bytes. Null fractions are 0%, 20% and 100%. A third writer alternates between unsafe and generic rows, so each path appends where the other left off. A fourth writes generic rows that hold unsafe values.
  • Collections of up to 150 elements cover null words past the first and every bit offset in the Arrow validity buffer. Strings and binaries of every length from 0 to 140 bytes cover both sides of the word-copy threshold.
  • The vector comparison now also checks that variable-width offsets never decrease, null rows included, since an unfilled hole would leave one smaller than the last.
  • New CometStringWriterSuite tests check the overflow and negative-length cases for an unsafe row and for an unsafe array, including that nothing is allocated.
  • A new CometArrowWriterSuite test reads an unsafe array of decimals at a narrower precision than it was written with, so that it holds 10^p or -10^p, the smallest values past the precision p. It checks that the write still fails as UnsafeArrayData.getDecimal does, for a decimal held as the unscaled long and one held as the unscaled bytes.
  • A new CometArrowWriterSuite test writes rows of nested arrays, maps and structs. After each row it walks the writers' fields and those of any unsafe view they hold, and checks that the row's backing array is not reachable. It fails when the nested writers reuse one view across values, and when a version that detaches its views after each value leaves out any one writer.
  • A new CometArrowWriterSuite test appends a dictionary-encoded string field as a column after its struct took the row path, which leaves the column offsets to fill, and past the writer's initial capacity. Leaving reserveValues out of the dictionary path, which no test caught before, now fails it.
  • Each of eleven deliberate bugs in the new code makes these tests fail: inverted validity, offset holes not filled, a missing sign extension, a wrong value width, an unshifted bit mask, map entries not marked valid, a short-copy tail, either check dropped, and the decimal array's range check skipped in either branch. An off-by-one at 10^p or a negation without the + 1 in the existing DecimalWriter.fitsPrecision, which none of these suites caught before, now fails the new decimal test.
  • CometArrowWriterSuite, CometStringWriterSuite and CometArrowStreamSuite pass on Spark 4.1 and on Spark 3.4 / Scala 2.12. With this branch merged into main, the SparkToColumnar and LocalTableScan tests in CometExecSuite, CometInMemoryCacheSuite, CometInMemoryCacheKryoSuite, CometInMemoryCachePruningSuite and CometShuffleInputConversionSuite pass on 4.1 too.
  • The strict Scala warnings build (./mvnw test-compile -Pspark-3.5 -Pstrict-warnings) and the semantic scalafix check on the 3.5 profile both pass.

ArrowWriter's row path read every row through Spark's generic getters. A string
went through a UTF8String and Arrow's per-value setter, each array, map and
struct allocated a view, and every element took a generic write with two or
three virtual or interface calls.

Unsafe rows now take a path of their own, and unsafe arrays are written in bulk:

- Fixed-width values are read from their slot at the vector's width, and arrays
  of them are copied in one block, with the validity inverted from the array's
  null words.
- Strings and binaries are copied from the offset and size in their slot, short
  ones a word at a time. A negative size or an offset overflow is rejected
  before anything is written or allocated, for a value and for every element of
  an array.
- Decimals up to 18 digits are written from the unscaled long, and wider ones
  are range-checked in place.
- Arrays, maps and structs point one reused view at each value.

Closes apache#6721.
@andygrove andygrove added enhancement New feature or request performance area:scan Parquet scan / data reading labels Oct 6, 2026
- Use Arrow's fillEmpties and setValueLengthSafe in place of hand-written hole
  filling and buffer growth, and share one variable-width overflow check.
- Read the fixed-width vector through its field in every per-value method,
  and have the base defaults delegate to write.
- Move the columnar code StringWriter and BinaryWriter shared into their base.
- Write unsafe arrays of decimals up to 18 digits from the unscaled long,
  range-checked as UnsafeArrayData.getDecimal checks it.
- Mark map entries valid with one masked write rather than a write per bit,
  and call setOne only for long runs.
- Share the tests' reject, root-comparison and fill helpers.
@andygrove

Copy link
Copy Markdown
Member Author

@comphead this could help with complex types

@mbutrovich mbutrovich left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Thanks @andygrove. The differential tests against the generic row path are a good fit for this change, and the unsafe and generic paths alternating in one writer covers the append boundaries well. I have one comment, about test coverage for a decimal check on the new path.

#6729 was closed without being merged. Should the paragraph about it come out of the description? It says the zero that #6729 writes under each null applies to the new path, and that would end up in the squash commit message.

Comment on lines +1209 to +1215
case array: UnsafeArrayData if precision <= Decimal.MAX_LONG_DIGITS =>
val unscaled = array.getLong(ordinal)
val fits = DecimalWriter.fitsPrecision(unscaled >> 63, unscaled, precision)
if (fits) {
putLong(valueAddress, unscaled)
}
fits

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Could we add a test for this range check? The description says a value past the precision in an unsafe array still fails as before, but the only overflow test I found, a wide decimal past its precision still fails the row path, uses an unsafe row. If I replace the fitsPrecision call here with true, CometArrowWriterSuite, CometStringWriterSuite and CometArrowStreamSuite still pass, so nothing would catch the array path writing an out-of-range long without the error that UnsafeArrayData.getDecimal raises. Here is something along the lines of the existing test, covering both the long-backed branch and the wide putUnsafeBytesIfFits branch below it:

  test("a decimal past its precision in an unsafe array still fails the row path") {
    // Written at a wider precision and read at a narrower one, so the array holds more digits.
    Seq(
      (DecimalType(18, 2), DecimalType(5, 2), "1234567.89"),
      (DecimalType(38, 0), DecimalType(20, 0), "123456789012345678901")).foreach {
      case (written, read, value) =>
        val row = UnsafeProjection.create(new StructType().add("a", ArrayType(written)))(
          new GenericInternalRow(Array[Any](new GenericArrayData(Array[Any](
            Decimal(new JavaBigDecimal(value), written.precision, written.scale))))))
        val allocator = new RootAllocator(Long.MaxValue)
        val root = VectorSchemaRoot.create(
          Utils.toArrowSchema(new StructType().add("a", ArrayType(read)), "UTC"),
          allocator)
        try {
          val writer = ArrowWriter.create(root, 1)
          intercept[ArithmeticException](writer.write(row))
        } finally {
          root.close()
          allocator.close()
        }
    }
  }

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

Thanks, added in b441fe5 along the lines of yours. It uses 10^p and -10^p, the smallest values past the precision, rather than 1234567.89 and 123456789012345678901. With the check here replaced by true, or with the wide branch ignoring what putUnsafeBytesIfFits returns, only the new test fails. The same goes for an off-by-one at 10^p or a negation without the + 1 in fitsPrecision, which none of the three suites caught before.

I also took the #6729 paragraph out of the description.

@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 all four changed files and both commits in the full diff from 6940c2e584b96ae9e6666228ab14d5d8bcc693f6 to 0182cc9c03d599836157adeef5a851ce6260cbaa. The PR is not a draft. Read the existing review, conversation comment, inline comment, and unresolved thread.

Routed skills: review-comet-pr, review-comet-ffi-pr, and review-comet-memory-pr. Also checked the contributor documentation and timestamp invariants.

Summary

  • Prior state and problem: ArrowWriter accessed unsafe rows through generic getters, allocating intermediate string, decimal, and nested-value objects and dispatching individual array elements through scalar writers.
  • Design approach: Dispatch unsafe rows to specialized writers, copy compatible array payloads in bulk, and reuse nested Spark views. Generic inputs retain their existing conversion semantics.
  • Correctness: Checked widths, null-bit inversion, append offsets, capacity growth, decimal sign extension and range checks, and nested writer counts. Spark source confirms that short decimals require different validation in UnsafeRow and UnsafeArrayData. Focused tests and direct comparisons with Spark found no introduced P1/P2 issue.
  • Compatibility analysis: Compared relevant unsafe layouts and getters across Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3, and 4.2.0. Timestamps retain their raw microseconds and existing timezone metadata. Compared the touched writer against branch-1.1 at e9efd9f764ee0a59b7898ff028d6985d4a7a28e1. The intended change is conversion efficiency, with no unintended valid-input semantic change attributable to this PR identified.
  • Key design decisions: Raw fixed-width copying is gated on little-endian representation. Booleans and decimals retain conversion-specific handling. Variable-width array lengths are checked before buffer growth or payload writes, and nested fields now receive capacity checks.
  • Implementation sketch: writeUnsafeRowField handles row slots, writeArrayElements handles collection payloads, and shared helpers maintain Arrow validity and offsets. Arrays, maps, and structs reuse separate views within their writers.
  • Performance: Removing temporary objects, bulk-copying primitive arrays, and reserving variable-width array storage once directly address the stated overhead. The PR supplies microbenchmark and end-to-end measurements. I did not independently reproduce those timings and found no evidence-backed P1/P2 performance regression.
  • Design: The optimization remains inside the existing writer hierarchy. Arrow output remains copied into owned buffers, and stream ownership, release callbacks, and memory-pool accounting are unchanged.
  • Abstraction & complexity: VariableWidthArrowFieldWriter usefully consolidates string and binary handling. The additional low-level code corresponds to concrete Spark layouts, with shared bitmap and copy helpers limiting duplication. No complexity concern met the P1/P2 reporting bar.
  • Behavioral changes worth calling out: Generic rows containing unsafe nested values also use the new paths. Unsafe variable-width inputs receive earlier overflow and negative-length rejection. The change does not add fallback rules, configuration defaults, or timezone conversion.
  • Suggested improvements: No additional P1/P2 change is requested. The existing decimal-array test request and description correction are already covered by the discussion. Independent boundary checks confirmed the current decimal guards behave correctly, so these comments do not establish an unresolved runtime blocker.

Exact-head CI: 26 checks succeeded and 15 were skipped, with no failed checks. Successful checks include Linux Spark 4.1 suites, TPC-H/TPC-DS verification, native tests, and multi-profile compilation/lint checks. Spark SQL suites, Iceberg suites, macOS, and benchmark jobs were skipped.

Validation: Fresh isolated compilation of the pinned Arrow package and its three unmodified suites passed all 74 tests in CometArrowWriterSuite, CometStringWriterSuite, and CometArrowStreamSuite. Direct comparisons with Spark’s ArrowWriter passed 73,024 values on each of Spark 3.5.9 and 4.1.3, including mixed generic/on-heap/off-heap inputs, nested values, nulls, signed zero, NaN, and negative timestamps. Another 124 decimal-array boundary cases per version matched Spark’s values or exception classes.

Limits: Local validation used cached supporting dependencies rather than a full Maven/native rebuild. Full Spark SQL/Iceberg integration suites, other runtime profiles, and performance benchmarks were not rerun. Project files remain unchanged.

No introduced P1/P2 issues found within this review.

@viirya viirya 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 @andygrove, this is a nice speedup, and the shape of the change makes it easy to trust. Comparing four writers against each other (generic, unsafe on and off heap, the two alternating, and generic rows holding unsafe values) checks the new paths against the existing one at every append boundary. I also like that the variable-width paths reject a bad value before anything is grown or written.

I went through the new paths against Spark's unsafe layouts and Arrow 18.3's variable-width internals and didn't find anything that changes results for valid input compared with branch-1.1. Passing the vector to FixedWidthArrowFieldWriter's constructor also avoids reading the subclass's valueVector before it is initialized, which matters for unsafeValueWidth.

+1 to @mbutrovich's points about a test for the decimal range check in unsafe arrays and the #6729 paragraph in the description. I left two more inline: test coverage for the block copy of primitive arrays, and the code that the three nested writers repeat.


// Nested shapes whose unsafe forms take paths of their own: elements converted one at a time,
// collections inside collections, and a struct wider than one word of null bits.
private val moreNestedTypes: Seq[DataType] = Seq(

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.

unsafeValueWidth in ArrowWriters.scala lists twelve vector types whose values are copied out of an unsafe array in one block. That relies on Spark's element size matching the Arrow type width for each of them. In these tests only int and long elements reach that copy (ArrayType(IntegerType), ArrayType(ArrayType(LongType)), and the int and long map keys). Byte, short, float, double, date, timestamp, timestamp_ntz, the two interval types and time never go through it as array elements.

Could we add primitiveTypes.map(ArrayType(_)) here? fill already handles every one of those element types, so this should be a one-line change that covers each entry in the list, including the 1- and 2-byte widths.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

Done in 61ec88b. moreNestedTypes now starts with primitiveTypes.map(ArrayType(_)), which also replaces the boolean, decimal(38,10) and binary arrays that were listed by hand. I logged which vectors reach the block copy while the suite runs. Before, only IntVector and BigIntVector did. Now all twelve in unsafeValueWidth do, including TinyIntVector and SmallIntVector, and TimeNanoVector on 4.1.

}
}

override private[arrow] def writeUnsafeRowField(row: UnsafeRow, ordinal: Int): Unit = {

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.

ArrayWriter, StructWriter and MapWriter each have the same writeUnsafeRowField and writeArrayElements. The only difference is which reused view writeUnsafe points at the value. Would it work to move those two methods into a small shared base class with an abstract writeUnsafe(base, baseOffset, offsetAndSize) that each writer implements? Then the null and count handling would live in one place rather than three. The stale-view mutations in the description show that a change here currently has to be made in each writer separately.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

I tried this, and it costs more than I expected. As a private method in each writer, writeUnsafe is bound statically and inlined. In a shared writeArrayElements and writeUnsafeRowField it becomes one virtual call site for all three writers, which goes megamorphic once a JVM has converted arrays, maps and structs. I warmed every case together and ran four case orders, with three alternating JVMs each. The shared class was 2-10% slower on most nested shapes, most of all on array, array<struct<int,string>>, array<array> and struct<int,string>, and level on flat ones. So I'd rather keep the three copies. The differential tests cover each of them, as the three stale-view mutations show, so a change that misses one of them fails.

The shared class did win on array<map<string,int>>, at 103 ns per row against 120-178, but that turned out to be a separate C2 quirk. When C2 inlines MapWriter.writeEntries into the loop over the array's maps, that loop is slow under a mixed profile, and the shared class only avoided it because its megamorphic call stops the inlining. -XX:CompileCommand=dontinline on writeEntries also gives 103, and so does warming only that shape. main takes 211 on it, so this isn't a regression, and I filed #6757 to dig into it.

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 measuring it. A shared call site going megamorphic is a good reason to keep the three copies, and the stale-view mutations show each copy is covered. I'm fine leaving this as is.

…l fails

Reads an unsafe array of decimals at a narrower precision than it was
written with, holding 10^p or -10^p, for both the unscaled long and the
unscaled bytes.
Each of the twelve vector types whose unsafe array elements are copied in
one block now reaches that copy, at every element width. Only int and long
arrays did before.

@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 commit 61ec88b771a77562b3ebe058183b2c84596852f2. New P1/P2 findings are detailed inline.


/** Sets value `count` to the unsafe array that `offsetAndSize` locates from `baseOffset`. */
private def writeUnsafe(base: AnyRef, baseOffset: Long, offsetAndSize: Long): Unit = {
elements.pointTo(base, baseOffset + (offsetAndSize >> 32), offsetAndSize.toInt)

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.

[P2] [P2] Release backing storage held by reusable nested views

elements.pointTo leaves the reusable view holding the entire input row’s backing array after its contents have been copied into Arrow. StructWriter.fields and MapWriter.entries do the same. When nullable nested fields have their last non-null values in different rows, each view retains a different complete row, including unrelated large payloads. Subsequent nulls do not release those references.

Through the production RowArrowReader, 24 streamed unsafe rows with a 2 MiB binary payload and distinct populated nested fields retain 50,337,600 bytes of consumed input storage at this head. The base retains none. Arrow already owns the output copies, so these input arrays should be collectible after conversion. This adds approximately 48 MiB of unnecessary JVM heap residency during batch construction. The arrays become collectible when the batch writer dies, so this increases peak residency rather than permanently leaking memory.

Could we detach each reusable view after its synchronous copy, including the map’s key/value views, and add an input-lifetime regression test?

Evidence: Run python /tmp/pr6730-61ec-review-wjvupbl7/run-retention.py. It separately compiles the base and exact-head ArrowWriters.scala with the production RowArrowReader, then runs ReaderRetentionProbe.scala using Spark 4.1.3, JDK 17, -Xmx160m and Serial GC. The iterator generates 24 UnsafeProjection(...).copy() rows containing a 2 MiB binary payload plus eight array, eight map and eight struct fields, with only one nested field populated per row. The harness keeps only weak references to consumed backing arrays. Forced GC inside the final hasNext, while batch construction is active, reports base (0,0) surviving arrays/bytes versus head (24,50337600). After loadNextBatch returns, both report (0,0). Outputs are in base-reader.log and head-reader.log. Spark’s pointTo implementations confirm the strong backing-object references.

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.

+1, and I should have caught this in my first pass. Old code allocated a fresh view per value, so nothing outlived the copy. Here each of elements, fields and entries keeps the last backing object it was pointed at, and entries does the same through its key and value arrays. It is at most one object per nested writer, and only until the batch's writer is dropped. But when those objects are per-row copies, or on-heap pages that the memory manager has already freed, that is heap Spark doesn't account for.

One thing to watch in the fix: pointTo(null, 0L, 0) is fine for UnsafeRow, but UnsafeArrayData.pointTo and UnsafeMapData.pointTo read their header with Platform.getLong(baseObject, baseOffset). With a null base and a zero offset, that reads address 0 and would crash the JVM. Pointing them at a shared empty layout avoids that. That would be a long[] {0} for an array of no elements, and long[] {8, 0, 0} with a size of 24 for a map, which is a key array of 8 bytes holding 0 elements and then an empty value array. Detaching once per row, after ArrowWriter.write finishes the row, may be cheaper than after every nested value, since an array<struct> would otherwise detach once per element. Whichever measures better is fine with me.

For the test, could we check that each view no longer references the input after a write rather than relying on GC with weak references? Something like a package-private accessor for the views' base objects would keep it deterministic.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

Thanks, confirmed. Each of the three views held the last row it read until the writer was dropped. Rather than detach them, 0623dba takes the reuse out. Each array, map and struct is read through Spark's getters again, which allocate a view per value as the generic path does, so nothing a writer holds outlives the copy.

In the mixed harness the reuse was only worth 1-4 ns per row on nested types, about a tenth of what this PR gains on them. array goes from 59 ns per row on main to 27 with reuse and 30 without, and map<string,string> from 117 to 50 and 53. Detaching after each value, as suggested, came out no better than allocating. It was 13% slower on array<array>, since an array of nested values detaches once per element. So reuse would need a detach once per row to pay off, and I filed #6765 for that.

The new test, the row path keeps no reference to an unsafe row once it is written, walks the writers' fields after each row, along with those of any unsafe view they hold, and checks that the row's backing array isn't reachable. So it doesn't depend on GC, as @viirya suggested. It fails on the previous head, and on a detaching version that leaves out any one of the three writers' detach.

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.

The retention is fixed in 0623dba by allocating a view per value again, and the new reachability test fails on the previous head. Resolved from my side.

@viirya viirya 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 the quick turnaround @andygrove. The new decimal test pins the boundary better than the values I would have picked. The primitive arrays now reach all twelve block copies, and the measurements on the shared base settle that question for me, so I'm fine keeping the three copies.

The one thing I'd like to see before this merges is @sunchao's point about the reused views keeping input rows alive. That's a memory change on this path, and I missed it in my first pass. I added some detail in that thread on how to detach the views safely.

@comphead comphead left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Thanks @andygrove. From reading the code, I agree the unsafe paths write the same vectors as the generic one for valid input. I checked them against Spark's unsafe getters and array element widths on 3.4.3, 3.5.8, 4.0.1 and 4.1.1, and against Arrow 18.3's variable-width internals. Two more points inline, one on sharing a buffer-growth step and one on repeated test cases. +1 to resolving the view retention that @sunchao and @viirya raised before this merges.

Comment on lines +1477 to +1483
while (valueVector.getValueCapacity < count + numElements) {
valueVector.reallocValidityAndOffsetBuffers()
}
reserveData(valueVector, end)
if (valueVector.getLastSet < count - 1) {
valueVector.fillEmpties(count)
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

writeVariableWidth and writeDictionaryVariableWidth open with the same two steps as these lines: fillEmpties when getLastSet < outStart - 1, then reallocValidityAndOffsetBuffers until getValueCapacity reaches outStart + numRows. Could the three share one helper in ArrowFieldWriter, next to reserveData? Then a change to how a variable-width append grows its buffers has one place to go. The order relative to reserveData does not matter. In Arrow 18.3 fillEmpties goes through handleSafe(index, 0), which here can grow only the validity and offset buffers.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

Done in 90c3bba. reserveValues sits next to reserveData, and all three paths call it, with the unsafe array path calling it after its checks. While I checked each call site with a mutation, dropping it from the dictionary path failed nothing, because no test appended dictionary strings after nulls the row path had left, or past the initial capacity. There's a test for that now, which appends a dictionary-encoded field as a column after its struct took the row path.

// Nested shapes whose unsafe forms take paths of their own: an array of each primitive type,
// whose elements are either copied in one block or converted one at a time, collections inside
// collections, and a struct wider than one word of null bits.
private val moreNestedTypes: Seq[DataType] = primitiveTypes.map(ArrayType(_)) ++ Seq(

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

primitiveTypes.map(ArrayType(_)) includes ArrayType(IntegerType) and ArrayType(StringType), which nestedTypes already lists. Both tests below concatenate the two lists, so the single-column array<int> and array<string> cases run twice in each of the six parameterized tests and in the 64-element test. The second run repeats the first exactly, because assertFilledRowsMatch seeds its Random with the schema's hash. Could we add a .distinct to the two types lists, or filterNot(nestedTypes.contains) here?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

Thanks, fixed in 90c3bba with filterNot(nestedTypes.contains), so array and array each run once.

The array, map and struct writers pointed one reused unsafe view at every
value they read. Each view then kept the last row it read reachable after
that row's values were copied to Arrow, along with the row's memory, until
the writer was dropped. They now read each value through Spark's getters,
which allocate a view per value as the generic path does. That keeps about
nine tenths of the gain on nested types.

A test walks the writers' fields after each row and checks that the row's
memory is not reachable from them.
writeVariableWidth, writeDictionaryVariableWidth and the unsafe array
path each filled the offsets of values skipped since the last one set and
grew the validity and offset buffers. They now call reserveValues, next to
reserveData.

A new test appends a dictionary-encoded string field as a column after
its struct took the row path, which the dictionary path's growth step had
no test for. The array of each primitive type no longer repeats the two
that nestedTypes lists.
@andygrove

Copy link
Copy Markdown
Member Author

@viirya 0623dba fixes the view retention by giving each nested value its own view again, and the new test checks that no writer still holds a written row. The details are in the thread on ArrowWriters.scala. Could you take another look?

@andygrove
andygrove requested a review from viirya October 8, 2026 03:28

@viirya viirya 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 @andygrove. Giving each nested value its own view again removes the retention outright, and it sidesteps the empty-header detail I mentioned, so I'm happy with that over detaching. The 1-4 ns you measured is small next to what the PR gains on nested types, and #6765 keeps the per-row detach on the table. The reachability test is a nice way to pin this without depending on GC, and it covers all three writers at the top level, nested in one another and as array elements.

I also checked the reserveValues change on the unsafe array path. It now runs after the element checks, so a rejected array still changes and allocates nothing. Filling the holes before reserveData is safe because fillEmpties only grows the validity and offset buffers there. Good catch on the dictionary path that no test covered.

LGTM.

@andygrove
andygrove added this pull request to the merge queue Oct 8, 2026
@andygrove

Copy link
Copy Markdown
Member Author

Thanks for the reviews @viirya @comphead @sunchao @mbutrovich!

Merged via the queue into apache:main with commit 5db0efe Oct 8, 2026
41 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:scan Parquet scan / data reading enhancement New feature or request performance

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Write strings and nested values in bulk when converting Spark rows to Arrow

5 participants