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.
What is the problem the feature request solves?
ArrowWriter's row path reads each array, map and struct of anUnsafeRowthrough Spark's getters, which allocate a view of the value for each value: anUnsafeArrayData, anUnsafeMapDatawith its two arrays, or anUnsafeRow. 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:main#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,
ArrowWritercould keep the nested writers that hold a view and detach just those after eachwrite, so flat rows pay nothing. AnUnsafeArrayDatacan point atlong[] {0}with a size of 8, anUnsafeMapDataatlong[] {8, 0, 0}with a size of 24, and anUnsafeRowat(null, 0, 0). A null base would crash the first two, because theirpointToreads a header.CometArrowWriterSuite'sthe row path keeps no reference to an unsafe row once it is writtenchecks 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.