You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
{{ message }}
Repository navigation
Write strings and nested values in bulk when converting Spark rows to Arrow #6721
#6566 made ArrowWriter copy Spark columns in bulk, but row input still goes through a field writer per column per row, and strings and nested values pay the most for that:
StringWriter gets a UTF8String for each value and appends it through Arrow's per-value path, which checks capacity and updates offsets one value at a time. perf: copy off-heap strings directly into Arrow #6281 removed the staging copy for off-heap strings, but not the per-value path.
ArrayWriter and MapWriter take each row's ArrayData or MapData and write it one element at a time through the element writer's generic write, which calls isNullAt and setSafe per element. MapWriter also marks each entry's struct slot defined.
StructWriter writes each field through the generic write as well. Only top-level fixed-width vectors are allocated at the batch size, so nested ones can't use the unchecked set.
I measured the row path on 10-04 against the native row converter that the JVM columnar shuffle uses, on identical UnsafeRows: one column of 8192 rows with every tenth row null, in the shapes CometArrowWriterBenchmark uses (arrays and maps hold 0-5 entries). In ns per row, on an M3 Max:
The native string number is with spark.comet.shuffle.jvm.preferDictionary.ratio at 0. With the default dictionary trial it is 17.2 (#6597).
Today the row path serves CometSparkRowToColumnar, CometLocalTableScanExec and the Comet cache serializer's row input (CometArrowConverters.rowToArrowBatchIter). Two approved PRs add more: #6607 converts the input of row-based shuffles when spark.comet.convert.shuffleInput.enabled is set, and #5634 turns Comet's cache on by default, so caching a row-based plan goes through it.
Describe the potential solution
Most row input is UnsafeRow: RDD scans, local tables and whole-stage-codegen stages all produce it. A path for UnsafeRow, falling back to the generic one for other InternalRows, could read the row's memory directly:
A string's offset and length are in its fixed-length slot. Copy the bytes straight into the data buffer, skipping the UTF8String and setSafe.
A fixed-width UnsafeArrayData stores its elements back to back at their natural width, so they can be copied in one block. Arrow validity is its null bitmap inverted: Spark sets a bit for null, Arrow for valid.
An array of strings, and a map's keys and values (two UnsafeArrayDatas), can be written the same way, element by element, without the generic writer.
Booleans (a byte each in UnsafeArrayData, a bit in Arrow) and decimals up to 18 digits (stored as longs, written as 128-bit values) still need converting per element.
Additional context
Part of #6565. CometStringWriterSuite (from #6281) pins how the string writer fails: an offset overflow throws Arrow's OversizedAllocationException before anything changes, and a negative length is rejected before anything is reserved or copied. A bulk path has to keep both.
What is the problem the feature request solves?
#6566 made
ArrowWritercopy Spark columns in bulk, but row input still goes through a field writer per column per row, and strings and nested values pay the most for that:StringWritergets aUTF8Stringfor each value and appends it through Arrow's per-value path, which checks capacity and updates offsets one value at a time. perf: copy off-heap strings directly into Arrow #6281 removed the staging copy for off-heap strings, but not the per-value path.ArrayWriterandMapWritertake each row'sArrayDataorMapDataand write it one element at a time through the element writer's genericwrite, which callsisNullAtandsetSafeper element.MapWriteralso marks each entry's struct slot defined.StructWriterwrites each field through the genericwriteas well. Only top-level fixed-width vectors are allocated at the batch size, so nested ones can't use the uncheckedset.I measured the row path on 10-04 against the native row converter that the JVM columnar shuffle uses, on identical
UnsafeRows: one column of 8192 rows with every tenth row null, in the shapesCometArrowWriterBenchmarkuses (arrays and maps hold 0-5 entries). In ns per row, on an M3 Max:ArrowWriterrow pathstringarray<string>map<string,string>array<int>struct<int,long,double,date>intThe native string number is with
spark.comet.shuffle.jvm.preferDictionary.ratioat 0. With the default dictionary trial it is 17.2 (#6597).Today the row path serves
CometSparkRowToColumnar,CometLocalTableScanExecand the Comet cache serializer's row input (CometArrowConverters.rowToArrowBatchIter). Two approved PRs add more: #6607 converts the input of row-based shuffles whenspark.comet.convert.shuffleInput.enabledis set, and #5634 turns Comet's cache on by default, so caching a row-based plan goes through it.Describe the potential solution
Most row input is
UnsafeRow: RDD scans, local tables and whole-stage-codegen stages all produce it. A path forUnsafeRow, falling back to the generic one for otherInternalRows, could read the row's memory directly:UTF8StringandsetSafe.UnsafeArrayDatastores its elements back to back at their natural width, so they can be copied in one block. Arrow validity is its null bitmap inverted: Spark sets a bit for null, Arrow for valid.UnsafeArrayDatas), can be written the same way, element by element, without the generic writer.Booleans (a byte each in
UnsafeArrayData, a bit in Arrow) and decimals up to 18 digits (stored as longs, written as 128-bit values) still need converting per element.Additional context
Part of #6565.
CometStringWriterSuite(from #6281) pins how the string writer fails: an offset overflow throws Arrow'sOversizedAllocationExceptionbefore anything changes, and a negative length is rejected before anything is reserved or copied. A bulk path has to keep both.