Skip to content

Commit 452e517

Browse files
authored
domain: add idle GC for cross-keyspace runtimes (#69013)
ref #68883
1 parent d7cdd7a commit 452e517

4 files changed

Lines changed: 219 additions & 13 deletions

File tree

pkg/domain/crossks/BUILD.bazel

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -51,10 +51,11 @@ go_test(
5151
],
5252
embed = [":crossks"],
5353
flaky = True,
54-
shard_count = 4,
54+
shard_count = 7,
5555
deps = [
5656
"//pkg/config",
5757
"//pkg/config/kerneltype",
58+
"//pkg/ddl/schemaver",
5859
"//pkg/dxf/framework/storage",
5960
"//pkg/executor/importer",
6061
"//pkg/infoschema",

pkg/domain/crossks/cross_ks.go

Lines changed: 57 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -48,7 +48,11 @@ import (
4848
"go.uber.org/zap"
4949
)
5050

51-
const crossKSSessPoolSize = 5
51+
const (
52+
crossKSSessPoolSize = 5
53+
crossKSRuntimeIdleTimeout = 30 * time.Minute
54+
crossKSRuntimeSweepInterval = time.Minute
55+
)
5256

5357
type runtimeEntry struct {
5458
sessMgr *SessionManager
@@ -326,7 +330,7 @@ func (*Manager) createSessionManager(
326330

327331
// release releases the runtime handle for the specified keyspace and holderID.
328332
// the resources will be cleaned up if there is no active holder after enough
329-
// time. we will impl this part later.
333+
// time.
330334
func (m *Manager) release(targetKS string, holderID string) {
331335
m.mu.Lock()
332336
defer m.mu.Unlock()
@@ -345,6 +349,57 @@ func (m *Manager) release(targetKS string, holderID string) {
345349
zap.Int("activeHolderCount", len(entry.activeHolders)))
346350
}
347351

352+
// RunSystemKSGCLoop periodically evicts idle cross keyspace runtimes.
353+
// as the name noted, this loop only runs in SYSTEM keyspace, user keyspace
354+
// only access the SYSTEM ks, and will be kept alive.
355+
// Note: there are raw caller to GetOrCreate which conflicts with the GC loop.
356+
// we will make the API more clear in the future after we can refactor those
357+
// callers to use Acquire instead of GetOrCreate directly.
358+
func (m *Manager) RunSystemKSGCLoop(ctx context.Context) {
359+
interval := crossKSRuntimeSweepInterval
360+
ticker := time.NewTicker(interval)
361+
defer ticker.Stop()
362+
363+
for {
364+
select {
365+
case <-ctx.Done():
366+
return
367+
case <-ticker.C:
368+
m.sweepIdleRuntimes(crossKSRuntimeIdleTimeout)
369+
}
370+
}
371+
}
372+
373+
func (m *Manager) sweepIdleRuntimes(idleTimeout time.Duration) {
374+
type evictedRuntime struct {
375+
targetKS string
376+
entry *runtimeEntry
377+
}
378+
379+
evicted := make([]evictedRuntime, 0, 1)
380+
now := time.Now()
381+
m.mu.Lock()
382+
for targetKS, entry := range m.runtimes {
383+
if len(entry.activeHolders) > 0 || entry.lastReleaseAt.IsZero() {
384+
continue
385+
}
386+
if now.Sub(entry.lastReleaseAt) < idleTimeout {
387+
continue
388+
}
389+
delete(m.runtimes, targetKS)
390+
evicted = append(evicted, evictedRuntime{targetKS: targetKS, entry: entry})
391+
}
392+
m.mu.Unlock()
393+
394+
for _, item := range evicted {
395+
logutil.BgLogger().Info("evict idle cross keyspace runtime",
396+
zap.String("targetKS", item.targetKS),
397+
zap.Duration("idleTimeout", idleTimeout),
398+
zap.Time("lastReleaseAt", item.entry.lastReleaseAt))
399+
item.entry.sessMgr.close()
400+
}
401+
}
402+
348403
// Close closes all session managers and their associated resources.
349404
func (m *Manager) Close() {
350405
m.mu.Lock()

pkg/domain/crossks/cross_ks_internal_test.go

Lines changed: 152 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -15,52 +15,76 @@
1515
package crossks
1616

1717
import (
18+
"context"
1819
"fmt"
1920
"testing"
21+
"time"
2022

2123
"github.com/ngaut/pools"
2224
"github.com/pingcap/tidb/pkg/config/kerneltype"
25+
"github.com/pingcap/tidb/pkg/ddl/schemaver"
2326
"github.com/pingcap/tidb/pkg/infoschema/validatorapi"
2427
"github.com/pingcap/tidb/pkg/keyspace"
2528
"github.com/pingcap/tidb/pkg/kv"
2629
"github.com/pingcap/tidb/pkg/testkit/testfailpoint"
2730
"github.com/pingcap/tidb/pkg/util"
2831
"github.com/stretchr/testify/require"
32+
clientv3 "go.etcd.io/etcd/client/v3"
2933
)
3034

3135
type runtimeHandleTestStore struct {
3236
kv.Storage
33-
ks string
37+
ks string
38+
closeCount int
3439
}
3540

3641
func (s *runtimeHandleTestStore) GetKeyspace() string {
3742
return s.ks
3843
}
3944

45+
func (s *runtimeHandleTestStore) Close() error {
46+
s.closeCount++
47+
return nil
48+
}
49+
4050
type runtimeHandleTestSessPool struct {
4151
util.DestroyableSessionPool
4252
closeCount int
53+
onClose func()
4354
}
4455

4556
func (p *runtimeHandleTestSessPool) Close() {
4657
p.closeCount++
58+
if p.onClose != nil {
59+
p.onClose()
60+
}
4761
}
4862

4963
func newRuntimeHandleTestManager(targetKS string) (*Manager, *runtimeEntry, *runtimeHandleTestStore, *runtimeHandleTestSessPool) {
5064
mgr := NewManager(&runtimeHandleTestStore{ks: keyspace.System})
5165
targetStore := &runtimeHandleTestStore{ks: targetKS}
5266
sessPool := &runtimeHandleTestSessPool{}
5367
entry := &runtimeEntry{
54-
sessMgr: &SessionManager{
55-
store: targetStore,
56-
sessPool: sessPool,
57-
},
68+
sessMgr: newRuntimeHandleTestSessionManager(targetStore, sessPool),
5869
activeHolders: make(map[string]struct{}),
5970
}
6071
mgr.runtimes[targetKS] = entry
6172
return mgr, entry, targetStore, sessPool
6273
}
6374

75+
func newRuntimeHandleTestSessionManager(targetStore *runtimeHandleTestStore, sessPool *runtimeHandleTestSessPool) *SessionManager {
76+
ctx, cancel := context.WithCancel(context.Background())
77+
return &SessionManager{
78+
ctx: ctx,
79+
cancel: cancel,
80+
exitCh: make(chan struct{}),
81+
store: targetStore,
82+
etcdCli: clientv3.NewCtxClient(context.Background()),
83+
schemaVerSyncer: schemaver.NewMemSyncer(),
84+
sessPool: sessPool,
85+
}
86+
}
87+
6488
func unusedRuntimeHandleFactoryGetter(t *testing.T) func(string, validatorapi.Validator) pools.Factory {
6589
return func(string, validatorapi.Validator) pools.Factory {
6690
t.Fatal("test should use the pre-seeded runtime entry")
@@ -238,3 +262,126 @@ func TestAcquireRuntimeHandle(t *testing.T) {
238262
require.Zero(t, sessPool.closeCount)
239263
})
240264
}
265+
266+
func TestEvictRuntime(t *testing.T) {
267+
if kerneltype.IsClassic() {
268+
t.Skip("cross keyspace runtime acquire is supported only in nextgen kernel")
269+
}
270+
271+
t.Run("skips active holders", func(t *testing.T) {
272+
targetKS := "ks-evict-active"
273+
mgr, entry, _, sessPool := newRuntimeHandleTestManager(targetKS)
274+
factoryGetter := unusedRuntimeHandleFactoryGetter(t)
275+
276+
first, err := mgr.Acquire(targetKS, "holder-1", factoryGetter)
277+
require.NoError(t, err)
278+
second, err := mgr.Acquire(targetKS, "holder-2", factoryGetter)
279+
require.NoError(t, err)
280+
281+
first.Release()
282+
entry.lastReleaseAt = time.Now().Add(-crossKSRuntimeIdleTimeout - time.Second)
283+
mgr.sweepIdleRuntimes(crossKSRuntimeIdleTimeout)
284+
285+
_, ok := mgr.Get(targetKS)
286+
require.True(t, ok)
287+
require.Contains(t, entry.activeHolders, "holder-2")
288+
require.Zero(t, sessPool.closeCount)
289+
second.Release()
290+
})
291+
292+
t.Run("closes idle entry outside manager lock", func(t *testing.T) {
293+
targetKS := "ks-evict-idle"
294+
mgr, entry, targetStore, sessPool := newRuntimeHandleTestManager(targetKS)
295+
factoryGetter := unusedRuntimeHandleFactoryGetter(t)
296+
sessPool.onClose = func() {
297+
require.True(t, mgr.mu.TryLock())
298+
mgr.mu.Unlock()
299+
}
300+
301+
handle, err := mgr.Acquire(targetKS, "holder-1", factoryGetter)
302+
require.NoError(t, err)
303+
handle.Release()
304+
entry.lastReleaseAt = time.Now().Add(-crossKSRuntimeIdleTimeout - time.Second)
305+
306+
mgr.sweepIdleRuntimes(crossKSRuntimeIdleTimeout)
307+
308+
_, ok := mgr.Get(targetKS)
309+
require.False(t, ok)
310+
require.Equal(t, 1, sessPool.closeCount)
311+
require.Equal(t, 1, targetStore.closeCount)
312+
})
313+
314+
t.Run("reacquire creates new runtime", func(t *testing.T) {
315+
targetKS := "ks-evict-reacquire"
316+
mgr, entry, oldStore, oldSessPool := newRuntimeHandleTestManager(targetKS)
317+
factoryGetter := unusedRuntimeHandleFactoryGetter(t)
318+
createCount := 0
319+
testfailpoint.EnableCall(t, "github.com/pingcap/tidb/pkg/domain/crossks/mockCreateSessionManager",
320+
func(createSessionManager *func(string, func(string, validatorapi.Validator) pools.Factory) (*SessionManager, error)) {
321+
*createSessionManager = func(ks string, _ func(string, validatorapi.Validator) pools.Factory) (*SessionManager, error) {
322+
createCount++
323+
newStore := &runtimeHandleTestStore{ks: ks}
324+
newSessPool := &runtimeHandleTestSessPool{}
325+
return newRuntimeHandleTestSessionManager(newStore, newSessPool), nil
326+
}
327+
},
328+
)
329+
330+
handle, err := mgr.Acquire(targetKS, "holder-1", factoryGetter)
331+
require.NoError(t, err)
332+
handle.Release()
333+
entry.lastReleaseAt = time.Now().Add(-crossKSRuntimeIdleTimeout - time.Second)
334+
335+
mgr.sweepIdleRuntimes(crossKSRuntimeIdleTimeout)
336+
require.Equal(t, 1, oldSessPool.closeCount)
337+
require.Equal(t, 1, oldStore.closeCount)
338+
339+
reacquired, err := mgr.Acquire(targetKS, "holder-2", factoryGetter)
340+
require.NoError(t, err)
341+
require.NotSame(t, oldStore, reacquired.Store())
342+
require.Equal(t, 1, createCount)
343+
reacquired.Release()
344+
})
345+
}
346+
347+
func TestRuntimeHandleManagerCloseClosesAllEntriesRegardlessOfIdleTimeout(t *testing.T) {
348+
firstKS := "ks-close-runtime-1"
349+
secondKS := "ks-close-runtime-2"
350+
mgr, firstEntry, firstStore, firstSessPool := newRuntimeHandleTestManager(firstKS)
351+
secondStore := &runtimeHandleTestStore{ks: secondKS}
352+
secondSessPool := &runtimeHandleTestSessPool{}
353+
mgr.runtimes[secondKS] = &runtimeEntry{
354+
sessMgr: newRuntimeHandleTestSessionManager(secondStore, secondSessPool),
355+
activeHolders: make(map[string]struct{}),
356+
}
357+
firstEntry.lastReleaseAt = time.Now()
358+
359+
mgr.Close()
360+
361+
require.Empty(t, mgr.GetAllKeyspace())
362+
require.Equal(t, 1, firstSessPool.closeCount)
363+
require.Equal(t, 1, firstStore.closeCount)
364+
require.Equal(t, 1, secondSessPool.closeCount)
365+
require.Equal(t, 1, secondStore.closeCount)
366+
}
367+
368+
func TestGCLoopExitsWhenContextCancelled(t *testing.T) {
369+
mgr := NewManager(&runtimeHandleTestStore{ks: keyspace.System})
370+
ctx, cancel := context.WithCancel(context.Background())
371+
done := make(chan struct{})
372+
go func() {
373+
mgr.RunSystemKSGCLoop(ctx)
374+
close(done)
375+
}()
376+
377+
cancel()
378+
379+
require.Eventually(t, func() bool {
380+
select {
381+
case <-done:
382+
return true
383+
default:
384+
return false
385+
}
386+
}, 5*time.Second, 10*time.Millisecond)
387+
}

pkg/domain/domain.go

Lines changed: 8 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -852,18 +852,20 @@ func (do *Domain) Start(startMode ddl.StartMode) error {
852852
return nil
853853
}
854854

855-
// TODO: we should sync the system keyspace info schema, not just load it,
856-
// currently, we assume there is no upgrade, so we only load the info schema of
857-
// system keyspace once, and it will not change during the lifetime of the domain,
858-
// it's not right, but it's enough to push subtasks which depends on it forward,
859-
// we will fix it in the future.
860855
func (do *Domain) loadSysKSInfoSchema() error {
861856
logutil.BgLogger().Info("loading system keyspace info schema")
857+
// it will trigger the creation of system keyspace session manager,
858+
// which will load the info schema cache.
862859
_, err := do.GetKSStore(keyspace.System)
863860
return err
864861
}
865862

866863
// GetKSStore returns the kv.Storage for the given keyspace.
864+
// we should forbid direct access cross KS component through Domain. we should
865+
// use AcquireKSRuntime to manage their lifecycle.
866+
// but Session dependents on Domain, to create a session pool we need to access the
867+
// GetKSStore/GetKSInfoCache inside Session where we don't know the runtime holder.
868+
// and trying to refactor that part can cause import cycle easily.
867869
func (do *Domain) GetKSStore(targetKS string) (store kv.Storage, err error) {
868870
mgr, err := do.crossKSSessMgr.GetOrCreate(targetKS, do.crossKSSessFactoryGetter)
869871
if err != nil {
@@ -873,6 +875,7 @@ func (do *Domain) GetKSStore(targetKS string) (store kv.Storage, err error) {
873875
}
874876

875877
// GetKSInfoCache returns the system keyspace info cache.
878+
// see comments of GetKSStore too.
876879
func (do *Domain) GetKSInfoCache(targetKS string) (*infoschema.InfoCache, error) {
877880
mgr, err := do.crossKSSessMgr.GetOrCreate(targetKS, do.crossKSSessFactoryGetter)
878881
if err != nil {

0 commit comments

Comments
 (0)