What is the problem the feature request solves?
CometSparkToColumnarExec reports conversionTime, "time converting Spark batches to Arrow batches". For row input (CometSparkRowToColumnar in the plan), RowArrowReader.loadNextBatch starts the clock before the loop that calls both rowIter.next() and writer.write(row), so the metric also counts whatever the child does to produce each row. That can be most of a stage: a JDBC read under spark.comet.convert.rowDataSource.enabled, the user's functions under a typed Dataset conversion (#6564), decoding Spark's cache into rows, or, with #6607, the Spark operators below a converted shuffle.
In #6566's end-to-end numbers, the cache query with row output reported 485 ms of conversionTime in an 814 ms run, and that figure includes decoding the cached batches into rows. It can't say how much the conversion itself costs, which also makes the row-path work in #6721 hard to judge from the metric.
Columnar input doesn't have this problem: SparkColumnarArrowReader stops the clock while it pulls the next Spark batch.
Describe the potential solution
Timing each row would cost more than the conversion it measures: two clock reads per row take longer than writing a narrow row. So the options are:
- Describe the metric differently for row input, for example "time producing and converting rows".
- Drop it for row input. Spark's
RowToColumnarExec reports only input rows and output batches.
Columnar input keeps the metric as it is.
Additional context
Part of #6565. CometLocalTableScanExec also uses RowArrowReader but doesn't report a conversion time.
What is the problem the feature request solves?
CometSparkToColumnarExecreportsconversionTime, "time converting Spark batches to Arrow batches". For row input (CometSparkRowToColumnarin the plan),RowArrowReader.loadNextBatchstarts the clock before the loop that calls bothrowIter.next()andwriter.write(row), so the metric also counts whatever the child does to produce each row. That can be most of a stage: a JDBC read underspark.comet.convert.rowDataSource.enabled, the user's functions under a typed Dataset conversion (#6564), decoding Spark's cache into rows, or, with #6607, the Spark operators below a converted shuffle.In #6566's end-to-end numbers, the cache query with row output reported 485 ms of
conversionTimein an 814 ms run, and that figure includes decoding the cached batches into rows. It can't say how much the conversion itself costs, which also makes the row-path work in #6721 hard to judge from the metric.Columnar input doesn't have this problem:
SparkColumnarArrowReaderstops the clock while it pulls the next Spark batch.Describe the potential solution
Timing each row would cost more than the conversion it measures: two clock reads per row take longer than writing a narrow row. So the options are:
RowToColumnarExecreports only input rows and output batches.Columnar input keeps the metric as it is.
Additional context
Part of #6565.
CometLocalTableScanExecalso usesRowArrowReaderbut doesn't report a conversion time.