Skip to content

Row-input conversionTime metric includes the time the child spends producing rows #6723

Description

@andygrove

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.

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