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]

Reply via email to