Skip to content

perf: reuse the nested unsafe views in ArrowWriter's row path, detaching them once per row #6765

Description

@andygrove

What is the problem the feature request solves?

ArrowWriter's row path reads each array, map and struct of an UnsafeRow through Spark's getters, which allocate a view of the value for each value: an UnsafeArrayData, an UnsafeMapData with its two arrays, or an UnsafeRow. An earlier revision of #6730 pointed one reused view per writer at every value instead. In a harness that warms every type up together, that was 1 to 4 ns per row faster on nested types:

ns per row main reused views a view per value
array 59.1 26.8 30.1
map<string,string> 116.9 49.7 53.5
array<struct<int,string>> 86.8 57.8 61.1
array<array> 80.1 45.6 47.8
map<int,long> 49.5 30.2 31.9
struct<int,string> 32.0 22.2 23.0

#6730 dropped the reuse because a reused view keeps the last row it read reachable after the row's values are copied, along with the row's memory, until the writer is dropped. That is at most one row per nested writer for the life of a batch's writer, but with large rows it is heap that Spark does not account for.

Describe the potential solution

Reuse the views and detach them once per row, rather than after every value. Detaching after every value, by pointing each view at an empty layout, costs as much as it saves. It was 13% slower than a view per value on array<array>, because an array of nested values detaches once per element. Instead, ArrowWriter could keep the nested writers that hold a view and detach just those after each write, so flat rows pay nothing. An UnsafeArrayData can point at long[] {0} with a size of 8, an UnsafeMapData at long[] {8, 0, 0} with a size of 24, and an UnsafeRow at (null, 0, 0). A null base would crash the first two, because their pointTo reads a header.

CometArrowWriterSuite's the row path keeps no reference to an unsafe row once it is written checks that no row stays reachable from the writers, and it has to keep passing.

Additional context

Measured on an M3 Max with AppleJDK 17: one column of 8192 rows per case, every tenth row and element null, the median of three alternating JVMs. Part of #6565. See the review of #6730.

Activity

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

Metadata

Metadata

Assignees

No one assigned

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions