kazuyukitanimura opened a new pull request, #6703:
URL: https://github.com/apache/datafusion-comet/pull/6703
fix: keep a Spark scan unconverted when the plan reads input_file_name
## Which issue does this PR close?
Closes #6573.
## Rationale for this change
`input_file_name()`, `input_file_block_start()` and
`input_file_block_length()` read
`InputFileBlockHolder`, a thread-local that the reader producing the rows
(`FileScanRDD`, the V2
file readers, `HadoopRDD`) sets as it moves from file to file. #3312 made
`CometScanRule` keep
Spark's `FileSourceScanExec` when a plan uses them. With
`spark.comet.convert.parquet.enabled=true`,
though, `CometExecRule` then wrapped that scan in
`CometSparkToColumnarExec`, ran the operators above
it natively, and left the Spark `Project` that evaluates the expressions
above `CometColumnarToRow`.
The conversion fills a batch, and Comet operators above it pull more input
before they emit. By the
time Spark evaluates the expressions, the reader has moved on to a later
file or reached the end of
its input and unset the holder. Rows then report another file's values or
`""`, `-1`, `-1`. The query
succeeds, so the wrong answer is silent.
The issue notes that keeping such plans on Spark is acceptable, and that is
what this PR does.
## What changes are included in this PR?
- `CometScanRule.readsInputFileBlock(plan)`: the check `CometScanRule`
already made for the native
Parquet scan, moved into the companion object so both rules share it.
- `CometExecRule.shouldApplySparkToColumnar` now refuses to convert a leaf
when the plan reads these
expressions, and records the fallback reason "Spark to Arrow conversion is
not compatible with
input_file_name, input_file_block_start, or input_file_block_length". The
previous body is
unchanged, renamed to `canApplySparkToColumnar`. The plan walk is lazy, so
it only runs once a leaf
could otherwise be converted.
- The guard applies to every leaf, not just Parquet. That covers the CSV
and JSON conversions, V1
and V2 scans, and other `spark.comet.sparkToColumnar` leaves such as
`RDDScan`, whose `HadoopRDD`
sets the same holder.
- A note in the user guide's data sources page.
Not in this PR: the native Iceberg scan and the native CSV V2 scan
(`spark.comet.scan.csv.v2.enabled`) have no such guard and have the same
problem. That gap predates
this change and needs its own fix in `CometScanRule.transformV2Scan`, which
I will file separately.
## How are these changes tested?
Two new tests in `CometExecSuite`. Each puts a filter between the scan and
the projection, so a Comet
operator would sit above the conversion, then compares results with Spark
using
`checkSparkAnswerAndFallbackReason` and asserts the plan has no
`CometSparkToColumnarExec`:
- `input_file_name above a converted Spark file scan keeps the scan on
Spark`: Parquet, CSV and JSON,
each as a V1 and a V2 scan.
- `input_file_name above a converted RDD scan keeps the scan on Spark`: an
`RDDScan` over
`sparkContext.textFile`.
--
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]