Skip to content

Stop ending Spark-to-Arrow batches at input file boundaries #6752

Description

@andygrove

What is the problem the feature request solves?

SparkColumnarArrowReader, which CometSparkToColumnarExec uses for columnar input, ends an Arrow batch early whenever the next Spark batch comes from a different file block. #6566 added this so that Spark expressions above the conversion would read the right file from InputFileBlockHolder. When the reader pulls a batch from the next file, it keeps that batch buffered, returns the partial Arrow batch, and sets the holder back to the returned batch's file.

#6703 makes CometExecRule leave a leaf, or the output of a typed Dataset operation, unconverted whenever the plan uses input_file_name(), input_file_block_start() or input_file_block_length(). The other conversions, the shuffle input from #6607 and the revert of transition-heavy stages, sit at the root of the map stage, so Spark has already evaluated these expressions below them. Once #6703 lands, no expression above a conversion reads the holder, and the split protects nothing.

It still costs batch size. Each file block boundary ends the Arrow batch, so a task that reads many small files, which Spark packs into one partition, produces one Arrow batch per file. Files of 1000 rows give native operators 1000-row batches instead of spark.comet.batchSize (8192). A large file costs one partial batch at its end.

The end of the input also differs from Spark. When the source runs out, the reader sets the holder back to the last batch's file, while FileScanRDD unsets it at the end of its input. Nothing can observe that after #6703, but it would show up if a conversion ever ran under these expressions again.

Describe the potential solution

After #6703 merges, remove the file tracking from SparkColumnarArrowReader (the currentInputFile* fields, restoreInputFile and the sameInputFile check) so batches fill across files. The class doc can then point at the guard in CometExecRule.shouldApplySparkToColumnar as the reason the reader doesn't track files. The CometArrowStreamSuite test "Spark columnar reader does not fill batches across input files" would become a test that batches do fill across files.

It's worth measuring a Parquet scan over many small files with spark.comet.convert.parquet.enabled=true before and after, for both the batch count and the query time.

Additional context

Part of #6565. Depends on #6703.

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