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]