Repository navigation
perf: convert Spark's unsafe rows to Arrow straight from their memory - #6730
Conversation
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.
- 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.
|
@comphead this could help with complex types |
mbutrovich
left a comment
There was a problem hiding this comment.
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.
| 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 |
There was a problem hiding this comment.
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()
}
}
}There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
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:
ArrowWriteraccessed 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
UnsafeRowandUnsafeArrayData. 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, and4.2.0. Timestamps retain their raw microseconds and existing timezone metadata. Compared the touched writer againstbranch-1.1ate9efd9f764ee0a59b7898ff028d6985d4a7a28e1. 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:
writeUnsafeRowFieldhandles row slots,writeArrayElementshandles 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:
VariableWidthArrowFieldWriterusefully 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
left a comment
There was a problem hiding this comment.
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( |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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 = { |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
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) |
There was a problem hiding this comment.
[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.
There was a problem hiding this comment.
+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.
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
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.
| while (valueVector.getValueCapacity < count + numElements) { | ||
| valueVector.reallocValidityAndOffsetBuffers() | ||
| } | ||
| reserveData(valueVector, end) | ||
| if (valueVector.getLastSet < count - 1) { | ||
| valueVector.fillEmpties(count) | ||
| } |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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( |
There was a problem hiding this comment.
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?
There was a problem hiding this comment.
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.
viirya
left a comment
There was a problem hiding this comment.
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.
|
Thanks for the reviews @viirya @comphead @sunchao @mbutrovich! |
Which issue does this PR close?
Closes #6721.
Rationale for this change
ArrowWriter's row path servesCometSparkRowToColumnar,CometLocalTableScanExecand the Comet cache serializer's row input. Since #6607,CometSparkRowToColumnaralso converts a shuffle's row input for native shuffle whenspark.comet.convert.shuffleInput.enabledis set, and #5634 adds more callers. Nearly all of those rows areUnsafeRows, but the writer read them through Spark's generic getters. A string went through aUTF8Stringand Arrow's per-value setter, and every element of an array, map or struct took the genericwrite, 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
CometArrowWriterBenchmarkuses. A harness outside the repo ranmain'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.mainOn
array<map<string,int>>, a C2 inlining quirk costs this PR about a quarter of its time, see #6757.CometArrowWriterBenchmarkitself warms one type at a time, so the JIT sees fewer types atmain's shared call sites, and the gains differ. They are smaller on strings,array<int>and the fixed-width rows, and larger onarray<string>and structs. Best of two alternating runs each:mainEnd to end,
countaggregates over a table of 4M rows cached in Spark's default format, with a releaselibcometatlocal[1]on the Spark 4.1 profile. The cache's vectorized reader is off, soCometSparkToColumnarExecconverts 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.conversionTimealso counts decoding the cached rows.conversionTime(ms)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.writehands anUnsafeRowto a newwriteUnsafeRowField, and unsafe arrays go to a newwriteArrayElements. Other rows take the generic path as before. Both read Spark's memory directly:UnsafeRow's own getters and written at the vector's width, without the virtual call tosetValueUnsafe. 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.UnsafeRow.getDecimalreads without checking it against the precision either. In an unsafe array the long is range-checked first, asUnsafeArrayData.getDecimalchecks it, so a value past the precision still fails as before. Wider decimals are range-checked in place instead of being copied to abyte[]first.UTF8StringorsetSafe. An array's elements are all checked first, then the buffers are grown once, through thereserveValuesthat the columnar string paths now share, and the bytes copied. Values up to 64 bytes are copied a word at a time, which beatsUnsafe.copyMemory's checks and call on short strings.setValidnow writes the bits around whole bytes with one masked write each, callingsetOne(a JNI call toUnsafe.setMemoryon JDK 17 and 21) only for runs of 64 bytes or more. A generic row holding an unsafe value, as a copiedColumnarBatchRowholds unsafe primitive arrays, takes these paths too.ArrowFieldWriter.writeUnsafeis gone. Its fixed-width override skipped the capacity check and relied onArrowWriter.createallocating 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.FixedWidthArrowFieldWriternow 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
CometStringWriterSuitepins still hold on the new path, and more strictly than on Arrow's setters. An offset overflow throws Arrow'sOversizedAllocationExceptionand a negative length throws anIllegalArgumentException, 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?
CometArrowWriterSuitetests 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.CometStringWriterSuitetests check the overflow and negative-length cases for an unsafe row and for an unsafe array, including that nothing is allocated.CometArrowWriterSuitetest 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 asUnsafeArrayData.getDecimaldoes, for a decimal held as the unscaled long and one held as the unscaled bytes.CometArrowWriterSuitetest 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.CometArrowWriterSuitetest 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. LeavingreserveValuesout of the dictionary path, which no test caught before, now fails it.+ 1in the existingDecimalWriter.fitsPrecision, which none of these suites caught before, now fails the new decimal test.CometArrowWriterSuite,CometStringWriterSuiteandCometArrowStreamSuitepass on Spark 4.1 and on Spark 3.4 / Scala 2.12. With this branch merged intomain, theSparkToColumnarandLocalTableScantests inCometExecSuite,CometInMemoryCacheSuite,CometInMemoryCacheKryoSuite,CometInMemoryCachePruningSuiteandCometShuffleInputConversionSuitepass on 4.1 too../mvnw test-compile -Pspark-3.5 -Pstrict-warnings) and the semantic scalafix check on the 3.5 profile both pass.