Skip to content

Commit 33038e1

Browse files
authored
sql/colexec: fix spill threshold (#23862)
- Auto spill threshold now uses (MemoryTotal - fileCacheMemory) / GoMaxProcs / 8 for both hashbuild and group operators, avoiding over-commitment when filecache is large - Add CanSpill bool to HashBuild; set true only for hashjoin so dedupjoin never triggers spill Approved by: @ouyuanning
1 parent 392b482 commit 33038e1

6 files changed

Lines changed: 19 additions & 4 deletions

File tree

‎pkg/sql/colexec/group/types2.go‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@ import (
2626
"github.com/matrixorigin/matrixone/pkg/container/batch"
2727
"github.com/matrixorigin/matrixone/pkg/container/types"
2828
"github.com/matrixorigin/matrixone/pkg/container/vector"
29+
"github.com/matrixorigin/matrixone/pkg/fileservice"
2930
"github.com/matrixorigin/matrixone/pkg/pb/plan"
3031
"github.com/matrixorigin/matrixone/pkg/sql/colexec"
3132
"github.com/matrixorigin/matrixone/pkg/sql/colexec/aggexec"
@@ -130,7 +131,8 @@ func (ctr *container) setSpillMem(m int64, aggs []aggexec.AggFuncExecExpression)
130131

131132
if m == 0 {
132133
// 0 means auto config. Here the formula is made up on the fly.
133-
mem := int64(system.MemoryTotal()) / int64(system.GoMaxProcs()) / 8
134+
fileCacheMem := fileservice.GlobalMemoryCacheSizeHint.Load()
135+
mem := (int64(system.MemoryTotal()) - fileCacheMem) / int64(system.GoMaxProcs()) / 8
134136
// min 128MB
135137
if mem < common.MiB*128 {
136138
mem = common.MiB * 128

‎pkg/sql/colexec/hashbuild/build.go‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -173,7 +173,7 @@ func (ctr *container) build(hashBuild *HashBuild, proc *process.Process, analyze
173173
}
174174

175175
// Check if we should enter spill mode based on batch memory size
176-
if hashBuild.NeedHashMap && hashBuild.shouldSpillBatches() {
176+
if hashBuild.shouldSpillBatches() {
177177
spillMode = true
178178
// Create spill files once
179179
spilledBuckets, spillFiles, err = createSpillFiles(proc)

‎pkg/sql/colexec/hashbuild/spill.go‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -206,8 +206,11 @@ func (ctr *container) rowCnt() int64 {
206206
}
207207

208208
func (hashBuild *HashBuild) shouldSpillBatches() bool {
209+
if !hashBuild.CanSpill || !hashBuild.IsShuffle || !hashBuild.NeedHashMap {
210+
return false
211+
}
209212
ctr := &hashBuild.ctr
210-
if !hashBuild.IsShuffle || ctr.spillThreshold <= 0 {
213+
if ctr.spillThreshold <= 0 {
211214
return false
212215
}
213216
if ctr.spillThreshold <= 100000 {

‎pkg/sql/colexec/hashbuild/spill_test.go‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -150,6 +150,8 @@ func TestShouldSpillBatches(t *testing.T) {
150150
hb := &HashBuild{
151151
IsShuffle: true,
152152
SpillThreshold: 1, // 1 byte
153+
CanSpill: true,
154+
NeedHashMap: true,
153155
}
154156
hb.ctr.setSpillThreshold(1)
155157
bat := batch.NewWithSize(1)
@@ -528,6 +530,8 @@ func TestShouldSpillBatchesRowThreshold(t *testing.T) {
528530
hb := &HashBuild{
529531
IsShuffle: true,
530532
SpillThreshold: 10, // Small row threshold
533+
CanSpill: true,
534+
NeedHashMap: true,
531535
}
532536
hb.ctr.setSpillThreshold(10)
533537

‎pkg/sql/colexec/hashbuild/types.go‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@ import (
2121
"github.com/matrixorigin/matrixone/pkg/common/reuse"
2222
"github.com/matrixorigin/matrixone/pkg/common/system"
2323
"github.com/matrixorigin/matrixone/pkg/container/batch"
24+
"github.com/matrixorigin/matrixone/pkg/fileservice"
2425
"github.com/matrixorigin/matrixone/pkg/pb/plan"
2526
"github.com/matrixorigin/matrixone/pkg/vm"
2627
"github.com/matrixorigin/matrixone/pkg/vm/message"
@@ -56,6 +57,7 @@ type HashBuild struct {
5657
NeedBatches bool
5758
NeedAllocateSels bool
5859
IsShuffle bool
60+
CanSpill bool
5961
Conditions []*plan.Expr
6062
JoinMapTag int32
6163
JoinMapRefCnt int32
@@ -145,7 +147,8 @@ func (hashBuild *HashBuild) ExecProjection(proc *process.Process, input *batch.B
145147
func (ctr *container) setSpillThreshold(threshold int64) {
146148
if threshold == 0 {
147149
// 0 means auto config
148-
mem := int64(system.MemoryTotal()) / int64(system.GoMaxProcs()) / 8
150+
fileCacheMem := fileservice.GlobalMemoryCacheSizeHint.Load()
151+
mem := (int64(system.MemoryTotal()) - fileCacheMem) / int64(system.GoMaxProcs()) / 8
149152
// min 128MB
150153
if mem < common.MiB*128 {
151154
mem = common.MiB * 128

‎pkg/sql/compile/operator.go‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -134,6 +134,7 @@ func dupOperator(sourceOp vm.Operator, index int, maxParallel int) vm.Operator {
134134
op.NeedBatches = t.NeedBatches
135135
op.NeedAllocateSels = t.NeedAllocateSels
136136
op.IsShuffle = t.IsShuffle
137+
op.CanSpill = t.CanSpill
137138
op.Conditions = t.Conditions
138139
op.JoinMapTag = t.JoinMapTag
139140
op.JoinMapRefCnt = t.JoinMapRefCnt
@@ -1565,6 +1566,7 @@ func constructBroadcastHashBuild(op vm.Operator, proc *process.Process, mcpu int
15651566

15661567
ret.HashOnPK = arg.HashOnPK
15671568
ret.NeedAllocateSels = !arg.HashOnPK
1569+
ret.CanSpill = true
15681570
if len(arg.RuntimeFilterSpecs) > 0 {
15691571
ret.RuntimeFilterSpec = arg.RuntimeFilterSpecs[0]
15701572
}
@@ -1657,6 +1659,7 @@ func constructShuffleHashBuild(node *plan.Node, op vm.Operator, proc *process.Pr
16571659

16581660
ret.HashOnPK = arg.HashOnPK
16591661
ret.NeedAllocateSels = !arg.HashOnPK
1662+
ret.CanSpill = true
16601663
if len(arg.RuntimeFilterSpecs) > 0 {
16611664
ret.RuntimeFilterSpec = plan2.DeepCopyRuntimeFilterSpec(arg.RuntimeFilterSpecs[0])
16621665
}

0 commit comments

Comments
 (0)