ReemaAlzaid opened a new issue, #13122:
URL: https://github.com/apache/gluten/issues/13122

   ### Description
   
   ### Description
   
   With cuDF enabled and `allowCpuFallback=false`, 48 of 103 TPC-DS SF1 queries 
pass. This issue tracks the 55 failures, grouped by root cause. Part of #9098; 
the function gaps also relate to #10752.
   
   The 55 failures group into 6 causes (by the first thing each query hits), 
and three of them already have open PRs.
   
   Two things to keep in mind:
   
   - `allowCpuFallback=false` doesn't mean the whole query runs on the GPU. 
Table scans still run on CPU (see **Scan stages** below).
   - This matters with default settings too. With `allowCpuFallback=true`, the 
stages behind causes 2–6 in the table don't fail, they just run on CPU, so 
fixing most of them moves more work onto the GPU for everyone. Cause 1 doesn't 
depend on that setting, so from reading the code (not tested with fallback on) 
it probably crashes in the default mode as well.
   
   ### How to reproduce
   
   Gluten `main` @ 9e944eca7, Spark 3.5.5, JDK 17, Velox `IBM/velox` 
`dft-2026_09_22` (cuDF 26.10), 2× L40S, CUDA 12.9, local mode. Built with 
`dev/buildbundle-veloxbe.sh --enable_gpu=ON --spark_version=3.5`. 
`spark.gluten.sql.columnar.backend.velox.cudf.enableTableScan` is left at its 
default (`false`).
   
   ```bash
   sbin/gluten-it.sh queries-compare \
     --local --preset=velox --benchmark-type=ds --error-on-memleak \
     --off-heap-size=10g -s=1 --threads=4 --iterations=1 
--decimal-as-double=true \
     --extra-conf=spark.gluten.sql.columnar.cudf=true \
     
--extra-conf=spark.gluten.sql.columnar.backend.velox.cudf.allowCpuFallback=false
   ```
   
   Add `--queries=q4,q11,q51` to rerun a single group. Decimals are stored as 
doubles here (`--decimal-as-double=true`), so most decimal paths aren't tested.
   
   Most failures show up as `Replacement with cuDF operator failed` 
(`ToCudf.cpp:242`). The error doesn't say which operator failed; the 
`Replacement Failed Operator:` warning just before it does.
   
   ### Failures by cause
   
   | # | Cause | Queries | Count | Fix in | Existing work |
   |---|---|---|---|---|---|
   | 1 | `CudfTopN` gets a CPU batch (bug) | q4, q11, q51 | 3 | Gluten | #13067 
|
   | 2 | `row_constructor_with_null` not supported | q1, q6, q7, q9, q13, q26, 
q28, q30, q32, q35, q44, q65, q81, q85, q92 | 15 | Gluten | #13104 |
   | 3 | `spark_legacy_cast` / `decimal_lessthan` not supported | q17, q21, 
q23a, q23b, q39a, q39b, q49, q54, q61, q66, q75, q83, q90, q93 | 14 | Velox 
cuDF or Gluten | — |
   | 4 | `isnull` not supported | q8, q14a, q14b, q38, q72, q78, q87, q97 | 8 | 
Velox cuDF or Gluten | — |
   | 5 | No cuDF `Expand` (ROLLUP / GROUPING SETS) | q5, q18, q22, q27, q36, 
q67, q77, q80, q86 | 9 | Velox cuDF | facebookincubator/velox#17146 |
   | 6 | Window / TopNRowNumber gaps | q47, q57, q53, q63, q89, q70 | 6 | Velox 
cuDF | facebookincubator/velox#18013 (avg only) |
   
   A query stops at its first failure, so some of these will hit another 
problem once their cause is fixed. From reading the SQL:
   
   - q14a also has an `avg`, a cast and a ROLLUP.
   - q18 and q67 have casts behind their `Expand`.
   - q70 has a ROLLUP after its `rank()`.
   - q44 and q67 filter on `rank()`, which turns into the same unsupported rank 
TopNRowNumber as q70 (cause 6).
   
   <details>
   <summary><b>1. CudfTopN gets a CPU batch</b></summary>
   
   This one is a real bug, not a missing feature.
   
   For `ORDER BY ... LIMIT`, when the input has more than one partition, 
`TakeOrderedAndProjectExecTransformer` creates its own shuffle internally. 
`AdjustStageExecutionMode` never sees that shuffle, so it gets read in CPU 
mode. The final stage is still tagged for GPU, though. The CPU batch goes 
through `CudfValueStream` unchanged, and `CudfTopN` crashes on it 
(`CudfTopN.cpp:148`, `cudfInput != nullptr`). TPC-H q3, q10 and q18 fail the 
same way.
   
   #13067 should fix it by uploading CPU batches to the GPU in 
`CudfVectorStream` (the reader behind `CudfValueStream`). So far I've only 
tested it on TPC-H; I still need to rerun q4, q11 and q51 with these settings. 
A follow-up could run that internal shuffle in GPU mode and skip the round trip.
   </details>
   
   <details>
   <summary><b>2. row_constructor_with_null</b></summary>
   
   For `avg` (and stddev, variance, corr, covar), Gluten uses 
`row_constructor_with_null` to pack the partial results before the final 
aggregation. cuDF only knows `row_constructor`, so the projection gets 
rejected, and the aggregation after it too.
   
   #13104 registers the function for cuDF from the Gluten side. I still need to 
test the global avg case (no GROUP BY: q9, q13, q28).
   </details>
   
   <details>
   <summary><b>3. spark_legacy_cast and decimal comparisons</b></summary>
   
   Since my PR #12051, Gluten sends regular casts to Velox as a function called 
`spark_legacy_cast` instead of a cast expression. cuDF only recognizes casts by 
expression type or by the names `cast` / `try_cast`, so any stage with a cast 
gets rejected: BIGINT→DOUBLE, DOUBLE→INT, decimal divides, and so on.
   
   With fallback on, these stages have been silently running on CPU since 
#12051. Nothing caught it because CI has no GPU runner, so nothing checks 
Gluten's expressions against cuDF.
   
   q75 also hits `decimal_lessthan`: Gluten renames decimal comparisons (since 
#3593), and cuDF only knows `lessthan`.
   
   The fix: teach cuDF that `spark_legacy_cast` is a normal cast, keeping the 
`cudf::is_supported_cast` check, and add the `decimal_*` names as aliases. This 
could also be done from Gluten with `registerCudfFunction`, like #13104 does.
   
   Just adding the name to the existing `cast` / `try_cast` alias list isn't 
enough. That entry skips the type check, so unsupported casts would pass 
validation and fail later. `spark_ansi_cast` should stay unsupported for now, 
since cuDF can't raise ANSI errors. (Some of Gluten's internal casts become 
`spark_ansi_cast` even with ANSI off, but none showed up in this run.)
   </details>
   
   <details>
   <summary><b>4. isnull</b></summary>
   
   Gluten uses Spark's name `isnull`, but cuDF only has Presto's `is_null` (and 
Spark's `isnotnull`). It mostly comes from INTERSECT/EXCEPT, where null-safe 
join keys become `hash_with_seed(42, coalesce(c, ''), isnull(c))`, and from 
`if(isnull(c), 1, 0)`.
   
   The fix is an `isnull` alias in the cuDF function registry plus one entry in 
`sparkUnaryOps` in `AstExpressionUtils.h`. Until that lands, a Gluten-side 
alias next to the #13104 registration would work as a stopgap. Probably the 
easiest item on this list.
   </details>
   
   <details>
   <summary><b>5. Expand (ROLLUP / GROUPING SETS)</b></summary>
   
   ROLLUP and GROUPING SETS turn into an `Expand`, and cuDF doesn't have one. 
`CudfGroupId` exists, but it handles a different plan node (`GroupIdNode`).
   
   facebookincubator/velox#17146 adds `CudfExpand` for the same node Gluten 
produces, so nothing needs to change on the Gluten side. The PR needs a rebase, 
though (see the task list).
   </details>
   
   <details>
   <summary><b>6. Window / TopNRowNumber</b></summary>
   
   Three separate gaps:
   
   - **Windows with more than one ORDER BY column** (q47, q57). The cuDF change 
this needs (NVIDIA/cudf#22627) is already in the cuDF version we build against 
(26.10); it just isn't wired up in Velox yet.
   - **`avg` over a whole partition** (q53, q63, q89, plus q47 and q57). 
Tracked in facebookincubator/velox#18013. facebookincubator/velox#18784 only 
covers DECIMAL avg where the result type matches the input type, so it doesn't 
help this run (DOUBLE) or Spark's DECIMAL avg, which widens the scale.
   - **`rank()` in TopNRowNumber** (q70). cuDF only supports `row_number` there 
(facebookincubator/velox#17805). Setting 
`spark.sql.optimizer.windowGroupLimitThreshold=-1` works around it.
   </details>
   
   ### Scan stages
   
   Scans stay on CPU by design for now (see the discussion in #9098 and 
`VeloxGPU.md`), and `allowCpuFallback=false` doesn't change that:
   
   - With `enableTableScan=false` (the default), any stage with a scan is 
marked CPU without being checked. In the 48 passing queries, that's 368 stages.
   - Velox's fallback check skips CPU table scans, so they never cause an 
error, even with fallback off.
   
   It gets confusing because, with fallback off, those "CPU" stages still run 
their other operators on cuDF after the CPU scan. Some of the failures in the 
table actually happen inside these stages (q93's cast failure, for example). 
I'll open a separate issue for this.
   
   Open question for the maintainers: should `allowCpuFallback=false` also 
cover scans? That would need `enableTableScan=true`, proper validation of scan 
stages (#12383, closed by the stale bot; I can reopen it), and a decision on 
whether a CPU `TableScan` should fail the query. Turning it on today would 
probably surface more failures first, like CPU batches reaching 
`CudfValueStream` (#13067) and gaps in the cuDF Parquet reader. I'm happy to 
work on it if we agree on the direction.
   
   ### Possible wrong results
   
   - With fallback off, multi-column `hash_with_seed` runs on the GPU, and cuDF 
combines columns differently from Spark (NVIDIA/cudf#21720). It also hashes 
booleans as 1 byte instead of 4. That's fine as long as both sides of a join 
are hashed on the GPU, but results will be wrong if one side is hashed on CPU. 
This already applies to passing queries like q3, and the `isnull` fix adds 
boolean keys, so it matters more once that lands.
   - Once casts are supported (cause 3), watch the cast semantics: 
DOUBLE→DECIMAL truncates on cuDF, but Spark rounds (q61). DOUBLE→INT is a plain 
`static_cast`, so NaN and out-of-range values won't match Spark (q54).
   
   ### Tasks
   
   - [ ] Review #13067 (q4, q11, q51)
   - [ ] Review #13104 (15 queries)
   - [ ] Support `spark_legacy_cast` and `decimal_*` comparisons in cuDF (14 
queries)
   - [ ] Add an `isnull` alias in cuDF (8 queries). Good first issue.
   - [ ] Rebase facebookincubator/velox#17146 for `Expand` (9 queries, plus 
q14a and q70). @jinchengchenghh, are you planning to pick this up again, or can 
I help?
   - [ ] Window with multi-column ORDER BY (needs a Velox issue)
   - [ ] Full-partition `avg` for DOUBLE (facebookincubator/velox#18013)
   - [ ] `rank()` / `dense_rank()` in cuDF TopNRowNumber (needs a Velox issue)
   - [ ] `CudfPlanValidator` only checks the last pipeline 
(`CudfPlanValidator.cc:72`), so in stages with a join, everything above the 
join goes unchecked. Small Gluten fix.
   - [ ] Make the `Replacement with cuDF operator failed` error say which 
operator failed. Small Velox fix.
   - [ ] Decide whether `allowCpuFallback=false` should cover scans (separate 
issue)
   
   
   ### Gluten version
   
   main branch


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