parthchandra opened a new issue, #6008:
URL: https://github.com/apache/datafusion-comet/issues/6008

   ### What is the problem the feature request solves?
   
   Comet doesn't natively scan the Text file format (`[COMET: Unsupported file 
format Text]`). In isolation that's fine, but when the unsupported source sits 
on the **build side of a broadcast join**, it cascades: the build branch (text 
scan → transforms → `BroadcastExchange`) stays on Spark, so the exchange never 
becomes a `CometBroadcastExchange`, so the `BroadcastHashJoin` can't go native 
— even though the large probe side is a fully native scan+filter. A tiny lookup 
table read from text disqualifies acceleration of the entire join over the big 
stream.
   
   ### Describe the potential solution
   
   The core problem isn't that Text is unsupported — it's that an unsupported 
**leaf on a broadcast build side vetoes native execution of the whole join**, 
including the large probe stream. One possible fix is to give the build branch 
a columnar bridge so it can reach the `BroadcastExchange` as native. There are 
two options 
   
   #### Option A — auto-insert `SparkToColumnar` at unsupported leaves on 
broadcast build sides
   `CometSparkToColumnarExec` wraps a Spark leaf and is delivered as a 
`CometScanWrapper` (a `CometNativeExec`), so operators above it 
(`explode`/`Filter`/`Project`) can convert and the exchange becomes a 
`CometBroadcastExchange`. Today it's off by default and its 
`sparkToColumnar.supportedOperatorList` excludes `FileSourceScan`, so this path 
is never hit
   
   Change: when an unsupported leaf scan feeds a **broadcast** build side and 
nothing else bloc ks the branch, insert `CometSparkToColumnarExec` at that leaf 
automatically. This is safe precisely because the build side is broadcast — 
it's bounded and small, so the row→Arrow copy cost is negligible, and the 
payoff is a native join over the (large) probe input.
   
   Guards:
   - Only on broadcast build sides (or gate by a size/row threshold) to avoid 
converting large row inputs where the copy would dominate.
   - Keep the existing `isSchemaSupported` check.
   - Respect a config, e.g. 
`spark.comet.sparkToColumnar.broadcastBuildSide.enabled` (default true), plus 
the existing opt-in for the general case.
   
   As an interim, this is already achievable by config — document it:
   ```
   --conf spark.comet.sparkToColumnar.enabled=true
   --conf 
spark.comet.sparkToColumnar.supportedOperatorList=Range,InMemoryTableScan,RDDScan,OneRowRelation,FileSourceScan
   ```
   
   #### Option B — native Text scan 
   Add a minimal native Text reader (single `value: string` column, 
line-delimited), mirroring the CSV native scan path, gated by e.g. 
`spark.comet.scan.text.enabled`. This removes the row→columnar copy entirely 
and gives native I/O on the build side, not just a bridge. More work, but it's 
the general fix and closes the "unsupported file format" gap directly.
   
   #### Note / follow-up
   Whatever the bridge, plan-time conversion of the build chain still hinges on 
the transforms above the leaf converting (e.g. `explode` over a Scala UDF 
returning `array<string>` — supported at plan time via codegen dispatch). 
Runtime execution of a native `Explode` over a `JvmScalarUdf` child should be 
validated separately; it's orthogonal to enabling the columnar bridge but worth 
testing in the same change.
   
   
   ### Additional context
   
   #### Steps to reproduce
   Join a large Parquet/Iceberg fact against a small allowlist loaded from a 
text file:
   ```scala
   val big   = spark.read.format("iceberg").load("db.events")          // 
native Comet scan
   val allow = spark.read.text("s3a://example-bucket/config/allow.txt") // Text 
-> Spark only
                .select(explode(to_hosts($"value")).as("host"), 
$"value".as("pattern"))
   big.join(broadcast(allow), big("host") === allow("host"))
      .where(rlike(big("url"), allow("pattern")))
   ```
   Explain: only `CometIcebergNativeScan` + `CometFilter` on the probe side 
convert; `Generate`, `Filter`, `Project`, `BroadcastExchange`, and the 
`BroadcastHashJoin` all remain Spark. Summary: `accelerated 2 of 12 operators 
(16%)`.
   
   #### Expected
   One of:
   1. Native scan support for Text (even a minimal single-`value`-column 
reader), or
   2. Make `spark.comet.sparkToColumnar.enabled` cover file-source leaves like 
Text by default (or document adding `FileSourceScan` to 
`sparkToColumnar.supportedOperatorList`), so a small build side can be adapted 
to Arrow and let the join go native.
   
   #### Additional context
   - With `sparkToColumnar.enabled=true` and `FileSourceScan` added to the 
supported list, the build chain (`SparkToColumnar` leaf → explode → filter → 
project) does convert at plan time, which flips the exchange and join native — 
but this is non-obvious and off by default.
   - Impact is huge: a small text lookup on the build side blocks native 
execution of a join over a very large probe input.
   


-- 
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