Skip to content

Commit 767c75a

Browse files
authored
*: support memTracker.detach for HashJoin, Apply and IndexLookUp in Close func (#54095) (#54261)
close #54005
1 parent b90e57a commit 767c75a

5 files changed

Lines changed: 52 additions & 10 deletions

File tree

pkg/executor/aggregate/agg_hash_executor.go

Lines changed: 10 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -273,7 +273,11 @@ func (e *HashAggExec) initForUnparallelExec() {
273273

274274
e.tmpChkForSpill = exec.TryNewCacheChunk(e.Children(0))
275275
if vars := e.Ctx().GetSessionVars(); vars.TrackAggregateMemoryUsage && variable.EnableTmpStorageOnOOM.Load() {
276-
e.diskTracker = disk.NewTracker(e.ID(), -1)
276+
if e.diskTracker != nil {
277+
e.diskTracker.Reset()
278+
} else {
279+
e.diskTracker = disk.NewTracker(e.ID(), -1)
280+
}
277281
e.diskTracker.AttachTo(vars.StmtCtx.DiskTracker)
278282
e.dataInDisk.GetDiskTracker().AttachTo(e.diskTracker)
279283
vars.MemTracker.FallbackOldAndSetNewActionForSoftLimit(e.ActionSpill())
@@ -396,7 +400,11 @@ func (e *HashAggExec) initForParallelExec(ctx sessionctx.Context) error {
396400
}, spillChunkFieldTypes)
397401

398402
if isTrackerEnabled && isParallelHashAggSpillEnabled {
399-
e.diskTracker = disk.NewTracker(e.ID(), -1)
403+
if e.diskTracker != nil {
404+
e.diskTracker.Reset()
405+
} else {
406+
e.diskTracker = disk.NewTracker(e.ID(), -1)
407+
}
400408
e.diskTracker.AttachTo(sessionVars.StmtCtx.DiskTracker)
401409
e.spillHelper.diskTracker = e.diskTracker
402410
sessionVars.MemTracker.FallbackOldAndSetNewActionForSoftLimit(e.ActionSpill())

pkg/executor/distsql.go

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -574,7 +574,11 @@ func (e *IndexLookUpExecutor) open(_ context.Context) error {
574574
// constructed by a "IndexLookUpJoin" and "Open" will not be called in that
575575
// situation.
576576
e.initRuntimeStats()
577-
e.memTracker = memory.NewTracker(e.ID(), -1)
577+
if e.memTracker != nil {
578+
e.memTracker.Reset()
579+
} else {
580+
e.memTracker = memory.NewTracker(e.ID(), -1)
581+
}
578582
e.memTracker.AttachTo(e.Ctx().GetSessionVars().StmtCtx.MemTracker)
579583

580584
e.finished = make(chan struct{})
@@ -858,7 +862,6 @@ func (e *IndexLookUpExecutor) Close() error {
858862
e.tblWorkerWg.Wait()
859863
e.finished = nil
860864
e.workerStarted = false
861-
e.memTracker = nil
862865
e.resultCurr = nil
863866
return nil
864867
}

pkg/executor/join.go

Lines changed: 8 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -176,6 +176,7 @@ func (e *HashJoinExec) Close() error {
176176
}
177177
e.probeSideTupleFetcher.probeChkResourceCh = nil
178178
terror.Call(e.rowContainer.Close)
179+
e.hashJoinCtx.sessCtx.GetSessionVars().MemTracker.UnbindActionFromHardLimit(e.rowContainer.ActionSpill())
179180
e.waiterWg.Wait()
180181
}
181182
e.outerMatchedStatus = e.outerMatchedStatus[:0]
@@ -214,8 +215,12 @@ func (e *HashJoinExec) Open(ctx context.Context) error {
214215
}
215216
e.hashJoinCtx.memTracker.AttachTo(e.Ctx().GetSessionVars().StmtCtx.MemTracker)
216217

217-
e.diskTracker = disk.NewTracker(e.ID(), -1)
218-
e.diskTracker.AttachTo(e.Ctx().GetSessionVars().StmtCtx.DiskTracker)
218+
if e.hashJoinCtx.diskTracker != nil {
219+
e.hashJoinCtx.diskTracker.Reset()
220+
} else {
221+
e.hashJoinCtx.diskTracker = disk.NewTracker(e.ID(), -1)
222+
}
223+
e.hashJoinCtx.diskTracker.AttachTo(e.Ctx().GetSessionVars().StmtCtx.DiskTracker)
219224

220225
e.workerWg = util.WaitGroupWrapper{}
221226
e.waiterWg = util.WaitGroupWrapper{}
@@ -1468,7 +1473,7 @@ func (e *NestedLoopApplyExec) fetchAllInners(ctx context.Context) error {
14681473

14691474
if e.canUseCache {
14701475
// create a new one since it may be in the cache
1471-
e.innerList = chunk.NewList(exec.RetTypes(e.innerExec), e.InitCap(), e.MaxChunkSize())
1476+
e.innerList = chunk.NewListWithMemTracker(exec.RetTypes(e.innerExec), e.InitCap(), e.MaxChunkSize(), e.innerList.GetMemTracker())
14721477
} else {
14731478
e.innerList.Reset()
14741479
}

pkg/util/chunk/list.go

Lines changed: 8 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -40,18 +40,23 @@ type RowPtr struct {
4040
RowIdx uint32
4141
}
4242

43-
// NewList creates a new List with field types, init chunk size and max chunk size.
44-
func NewList(fieldTypes []*types.FieldType, initChunkSize, maxChunkSize int) *List {
43+
// NewListWithMemTracker creates a new List with field types, init chunk size, max chunk size and memory tracker.
44+
func NewListWithMemTracker(fieldTypes []*types.FieldType, initChunkSize, maxChunkSize int, tracker *memory.Tracker) *List {
4545
l := &List{
4646
fieldTypes: fieldTypes,
4747
initChunkSize: initChunkSize,
4848
maxChunkSize: maxChunkSize,
49-
memTracker: memory.NewTracker(memory.LabelForChunkList, -1),
49+
memTracker: tracker,
5050
consumedIdx: -1,
5151
}
5252
return l
5353
}
5454

55+
// NewList creates a new List with field types, init chunk size and max chunk size.
56+
func NewList(fieldTypes []*types.FieldType, initChunkSize, maxChunkSize int) *List {
57+
return NewListWithMemTracker(fieldTypes, initChunkSize, maxChunkSize, memory.NewTracker(memory.LabelForChunkList, -1))
58+
}
59+
5560
// GetMemTracker returns the memory tracker of this List.
5661
func (l *List) GetMemTracker() *memory.Tracker {
5762
return l.memTracker

pkg/util/memory/tracker.go

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -245,6 +245,27 @@ func (t *Tracker) UnbindActions() {
245245
t.actionMuForHardLimit.actionOnExceed = &LogOnExceed{}
246246
}
247247

248+
// UnbindActionFromHardLimit unbinds action from hardLimit.
249+
func (t *Tracker) UnbindActionFromHardLimit(actionToUnbind ActionOnExceed) {
250+
t.actionMuForHardLimit.Lock()
251+
defer t.actionMuForHardLimit.Unlock()
252+
253+
var prev ActionOnExceed
254+
for current := t.actionMuForHardLimit.actionOnExceed; current != nil; current = current.GetFallback() {
255+
if current == actionToUnbind {
256+
if prev == nil {
257+
// actionToUnbind is the first element
258+
t.actionMuForHardLimit.actionOnExceed = current.GetFallback()
259+
} else {
260+
// actionToUnbind is not the first element
261+
prev.SetFallback(current.GetFallback())
262+
}
263+
break
264+
}
265+
prev = current
266+
}
267+
}
268+
248269
// reArrangeFallback merge two action chains and rearrange them by priority in descending order.
249270
func reArrangeFallback(a ActionOnExceed, b ActionOnExceed) ActionOnExceed {
250271
if a == nil {

0 commit comments

Comments
 (0)