andygrove opened a new issue, #6708:
URL: https://github.com/apache/datafusion-comet/issues/6708
### Describe the bug
With `spark.comet.convert.typedDataset.enabled=true`, `input_file_name()`,
`input_file_block_start()` and `input_file_block_length()` return `""`, `-1`
and `-1` for every row when Spark evaluates them above a Comet operator that
reads the converted output of a typed Dataset operation. Spark returns each
row's file path and block offsets. The query succeeds, so the wrong values are
silent.
This is #6573 for the conversion that #6564 added.
`CometSparkToColumnarExec` fills an Arrow batch from `SerializeFromObjectExec`,
which pulls rows from the scan below the typed operation. By the time the Spark
`Project` evaluates the expressions, the reader has moved on or unset
`InputFileBlockHolder`. The `SerializeFromObjectExec` case in
`CometExecRule.transform` calls `convertTypedDatasetOutput` without going
through `shouldApplySparkToColumnar`, so the guard that #6703 adds there does
not cover it.
```
*(2) Project [input_file_name() AS input_file_name()#576,
input_file_block_start() AS input_file_block_start()#577L,
input_file_block_length() AS input_file_block_length()#578L, value#568L]
+- *(2) CometColumnarToRow
+- CometFilter [value#568L], (value#568L >= 0)
+- CometSparkRowToColumnar
+- *(1) SerializeFromObject [input[0, bigint, false] AS value#568L]
+- *(1) MapElements ..., obj#566: bigint
+- *(1) DeserializeToObject assertnotnull(id#558L), obj#564:
bigint
+- *(1) ColumnarToRow
+- FileScan parquet [id#558L] Batched: true, ...
```
### Steps to reproduce
```scala
spark.range(9000).repartition(3).write.parquet("/tmp/ifn")
spark.conf.set("spark.comet.convert.typedDataset.enabled", "true")
spark.read
.parquet("/tmp/ifn")
.as[Long]
.map(_ + 1)
.where("value >= 0")
.selectExpr("input_file_name()", "input_file_block_start()",
"input_file_block_length()", "value")
.show(3, false)
```
Comet returns `["", -1, -1, <value>]` for all 9000 rows.
### Expected behavior
Each row reports the file it was read from, as Spark does. Leaving the typed
operation's output unconverted when the plan uses these expressions would also
be acceptable.
### Additional context
- Reproduced on `main` at `b80bf4e08` with the default Spark 4.1 profile,
comparing against the same query with `spark.comet.enabled=false`. It still
reproduces with #6703 merged into `main`.
- The conversion is off by default and is not in a release yet.
- A fix could check `CometScanRule.readsInputFileBlock(plan)` from #6703 in
the `SerializeFromObjectExec` case, next to the partial-reader check, and add
the case to the typed Dataset paragraph in
`docs/source/contributor-guide/adding_a_new_operator.md`.
--
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]