Skip to content

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

Description

@andygrove

What is the problem the feature request solves?

#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:

Column ArrowWriter row path Native row converter, after #6604
string 9.1-13 5.7
array<string> 32.9-35.7 24.8
map<string,string> 65-86 47.3
array<int> 11.0-11.8 9.1
struct<int,long,double,date> 22.3-22.7 11.4
int 4.0-5.2 2.3

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.

Activity

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

Metadata

Metadata

Assignees

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions