andygrove opened a new issue, #5625:
URL: https://github.com/apache/datafusion-comet/issues/5625

   ## What is the problem the feature request solves?
   
   A write from a Comet plan pays for a columnar-to-row transition that 
materialises an `UnsafeRow`
   per row, and none of the write path needs one. `OutputWriter.write`, 
`FileFormatDataWriter.write`
   and `WriteTaskStatsTracker.newRow` are all typed on `InternalRow`, 
`ParquetWriteSupport extends
   WriteSupport<InternalRow>` and reads fields through `SpecializedGetters`, and
   `BasicWriteTaskStatsTracker.newRow` ignores the row entirely. So the 
`UnsafeProjection` in
   `CometColumnarToRowExec` builds a row that the writer immediately decodes 
again.
   
   `ColumnarBatch.rowIterator()` already gives a reused `ColumnarBatchRow` that 
is a zero-copy view
   over the Arrow buffers - `CometColumnarToRowExec.doExecute` produces exactly 
that and then throws
   it away by projecting. Handing the writer the view instead skips the copy.
   
   For flat schemas this is not worth doing: the projection there is a 
generated fixed-width copy and
   it measures inside the run-to-run noise of a Parquet write. It becomes 
worthwhile once a struct,
   array or map is present, because the projection then has to build nested 
`UnsafeRow` /
   `UnsafeArrayData` with offset-and-length bookkeeping.
   
   Measured on an M3 Max, 1M rows, Spark 4.1, release build, best-of, comparing 
against today's
   `CometColumnarToRowExec`:
   
   | schema | uncompressed | snappy |
   | --- | --- | --- |
   | flat (fixed width / strings / 50 columns) | 0-3% (noise) | 0-3% (noise) |
   | one struct + array + map | 10% | 15% |
   | struct + array-of-structs + map-of-array-of-structs | 11% | 13% |
   | single struct, 1 level | 10% | 8% |
   | single struct, 8 levels | 7% | 8% |
   
   The gain comes from the presence of complex types rather than from depth - 
one level of nesting
   already captures it.
   
   This does not touch the real cost of a write, which is parquet-mr encoding. 
Only the native writer
   changes that. It is a cheap improvement to the interim path while native 
writes remain
   experimental, and unlike the native writer it is Spark-compatible by 
construction because Spark's
   own writer still does the encoding.
   
   ## Describe the potential solution
   
   A `ColumnarToRowTransition` that returns `batch.rowIterator()` unprojected, 
deliberately not
   `CodegenSupport` so that whole-stage codegen does not regenerate the copy, 
planted by
   `EliminateRedundantTransitions` under `WriteFilesExec` / 
`DataWritingCommandExec`.
   
   The row view is a reused mutable row, so it is only correct for a consumer 
that finishes with a row
   before pulling the next one. That restricts it to:
   
   - unpartitioned and unbucketed writes, which is what makes 
`FileFormatWriter` pick
     `SingleDirectoryDataWriter`. The partitioned and bucketed writers do not 
qualify - the required
     ordering puts a `SortExec` in between and `UnsafeExternalSorter` needs 
`UnsafeRow`, and
     `DynamicPartitionDataConcurrentWriter` spills through 
`UnsafeKVExternalSorter.insertKV` which is
     typed on `UnsafeRow`
   - Spark's own `FileFormat` implementations, whose `OutputWriter`s encode 
each row on the spot,
     rather than a third-party format that may buffer the `InternalRow` it is 
handed
   - schemas containing a complex type, per the measurements above
   
   Behind an off-by-default config while it is experimental.
   
   ## Additional context
   
   None.
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to