apache / apache/gluten

[VL] counterintuitive spills of join Operator

Open
#4,598 4 comments 0 reactions 0 assignees View on GitHub
bug triage
Dominant language
Scala
Stars
1.6k
Forks
657
Avg merge
2d 14h
Merged PRs (30d)
80

Description

### Backend

VL (Velox)

### Bug description

Hello, I am evaluating the performance of the gluten + velox backend on the TPC-DS 3000 benchmark and noticed that q95 performs weirdly.

While I set the spark.memory.offHeap.size to 16g, the entire execution time of q95 is **227.9 s**, and there is little data that has been spilled (part of the DAG shown below)

![image](https://github.com/oap-project/gluten/assets/43876761/edf9f0b6-3b4e-4b00-bd6b-c2029ce38fe2)

When I modify spark.memory.offHeap.size to 18g and keep the other configurations unchanged, the entire execution time of q95 degrades to **375.3 s,** and lots of spills have been triggered(several times higher than the 16g case.), which is counter-intuitive. In theory, allocate more memory to each task should trigger fewer spills.

when I disable the spill of join through `spark.gluten.sql.columnar.backend.velox.joinSpillEnabled false`, both of the execution times are the same. I am sure that the number of the executor is the same.

Any suggestions?

### Spark version

Spark-3.3.x

### Spark configurations

Here are the configs:

spark.master yarn
spark.deploy-mode client
spark.sql.cbo.enabled true

spark.executor.cores 14
spark.executor.memory 5g
spark.executor.memoryOverhead 944m
spark.memory.offHeap.size 18g
spark.memory.offHeap.enabled true
spark.gluten.enabled true
spark.plugins io.glutenproject.GlutenPlugin
spark.gluten.sql.columnar.backend.lib velox
spark.shuffle.manager org.apache.spark.shuffle.sort.ColumnarShuffleManager
spark.gluten.loadLibFromJar true
spark.sql.shuffle.partitions 208

spark.gluten.sql.columnar.backend.velox.maxSpillFileSize 1073741824
spark.default.parallelism 167

### System information

_No response_

### Relevant logs

_No response_

Contributor guide

Open the contributing guide

Research direction

Start with the q95 TPC-DS workload and compare the 16g and 18g configurations, focusing on the Velox join spill setting spark.gluten.sql.columnar.backend.velox.joinSpillEnabled. Gather the missing system information and relevant logs, then determine why the larger off-heap setting causes more spills and slower execution, with a reproducible explanation or fix as the outcome.

Written by the indexing model from the issue text.

Assessment

Tech stack
sql
Domain
backend, databases, performance
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Stale
Clarity
Needs clarification
Newbie friendliness
20/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.