From cd15aac134d86835ef364dd7300ebc591f2c2de4 Mon Sep 17 00:00:00 2001 From: Ping Yu Date: Fri, 12 Jun 2026 17:40:19 +0800 Subject: [PATCH 1/3] txnkv: Support resolve txn file locks Signed-off-by: Ping Yu --- txnkv/txnlock/lock_resolver.go | 19 +++++++++++++++---- 1 file changed, 15 insertions(+), 4 deletions(-) diff --git a/txnkv/txnlock/lock_resolver.go b/txnkv/txnlock/lock_resolver.go index 15c3d74e46..d60a3b5223 100644 --- a/txnkv/txnlock/lock_resolver.go +++ b/txnkv/txnlock/lock_resolver.go @@ -208,6 +208,7 @@ type Lock struct { UseAsyncCommit bool LockForUpdateTS uint64 MinCommitTS uint64 + IsTxnFile bool } func (l *Lock) IsPessimistic() bool { @@ -226,8 +227,8 @@ func (l *Lock) String() string { buf.WriteString(redact.Key(l.Key)) buf.WriteString(", primary: ") buf.WriteString(redact.Key(l.Primary)) - return fmt.Sprintf("%s, txnStartTS: %d, lockForUpdateTS:%d, minCommitTs:%d, ttl: %d, type: %s, UseAsyncCommit: %t, txnSize: %d", - buf.String(), l.TxnID, l.LockForUpdateTS, l.MinCommitTS, l.TTL, l.LockType, l.UseAsyncCommit, l.TxnSize) + return fmt.Sprintf("%s, txnStartTS: %d, lockForUpdateTS:%d, minCommitTs:%d, ttl: %d, type: %s, UseAsyncCommit: %t, txnSize: %d, isTxnFile: %t", + buf.String(), l.TxnID, l.LockForUpdateTS, l.MinCommitTS, l.TTL, l.LockType, l.UseAsyncCommit, l.TxnSize, l.IsTxnFile) } // NewLock creates a new *Lock. @@ -242,6 +243,7 @@ func NewLock(l *kvrpcpb.LockInfo) *Lock { UseAsyncCommit: l.UseAsyncCommit, LockForUpdateTS: l.LockForUpdateTs, MinCommitTS: l.MinCommitTs, + IsTxnFile: l.IsTxnFile, } } @@ -296,6 +298,7 @@ func (lr *LockResolver) BatchResolveLocks(bo *retry.Backoffer, locks []*Lock, lo expiredLocks := locks txnInfos := make(map[uint64]uint64) + txnFileIDs := make(map[uint64]bool) startTime := time.Now() for _, l := range expiredLocks { logutil.Logger(bo.GetCtx()).Debug("BatchResolveLocks handling lock", zap.Stringer("lock", l)) @@ -357,6 +360,9 @@ func (lr *LockResolver) BatchResolveLocks(bo *retry.Backoffer, locks []*Lock, lo } txnInfos[l.TxnID] = status.commitTS + if l.IsTxnFile { + txnFileIDs[l.TxnID] = true + } } logutil.BgLogger().Info("BatchResolveLocks: lookup txn status", zap.Duration("cost time", time.Since(startTime)), @@ -365,8 +371,9 @@ func (lr *LockResolver) BatchResolveLocks(bo *retry.Backoffer, locks []*Lock, lo listTxnInfos := make([]*kvrpcpb.TxnInfo, 0, len(txnInfos)) for txnID, status := range txnInfos { listTxnInfos = append(listTxnInfos, &kvrpcpb.TxnInfo{ - Txn: txnID, - Status: status, + Txn: txnID, + Status: status, + IsTxnFile: txnFileIDs[txnID], }) } @@ -1008,6 +1015,7 @@ func (lr *LockResolver) getTxnStatus(bo *retry.Backoffer, txnID uint64, primary var status TxnStatus resolvingPessimisticLock := lockInfo != nil && lockInfo.IsPessimistic() + isTxnFile := lockInfo != nil && lockInfo.IsTxnFile req := tikvrpc.NewRequest(tikvrpc.CmdCheckTxnStatus, &kvrpcpb.CheckTxnStatusRequest{ PrimaryKey: primary, LockTs: txnID, @@ -1017,6 +1025,7 @@ func (lr *LockResolver) getTxnStatus(bo *retry.Backoffer, txnID uint64, primary ForceSyncCommit: forceSyncCommit, ResolvingPessimisticLock: resolvingPessimisticLock, VerifyIsPrimary: true, + IsTxnFile: isTxnFile, }, kvrpcpb.Context{ RequestSource: util.RequestSourceFromCtx(bo.GetCtx()), ResourceControlContext: &kvrpcpb.ResourceControlContext{ @@ -1460,6 +1469,7 @@ func (lr *LockResolver) batchLiteResolveLocks(bo *retry.Backoffer, l *Lock, keys func (lr *LockResolver) resolveRegionLocks(bo *retry.Backoffer, l *Lock, region locate.RegionVerID, keys [][]byte, status TxnStatus) error { lreq := &kvrpcpb.ResolveLockRequest{ StartVersion: l.TxnID, + IsTxnFile: l.IsTxnFile, } if status.IsCommitted() { lreq.CommitVersion = status.CommitTS() @@ -1560,6 +1570,7 @@ func (lr *LockResolver) resolveLock(bo *retry.Backoffer, l *Lock, status TxnStat } lreq := &kvrpcpb.ResolveLockRequest{ StartVersion: l.TxnID, + IsTxnFile: l.IsTxnFile, } if status.IsCommitted() { lreq.CommitVersion = status.CommitTS() From 2a4426dcfda606da8e427885f3869e8286b204a6 Mon Sep 17 00:00:00 2001 From: Ping Yu Date: Fri, 12 Jun 2026 19:57:49 +0800 Subject: [PATCH 2/3] address comments Signed-off-by: Ping Yu --- internal/client/client_collapse.go | 6 +++++- internal/client/client_test.go | 22 ++++++++++++++++++---- txnkv/txnlock/lock_resolver.go | 3 +++ 3 files changed, 26 insertions(+), 5 deletions(-) diff --git a/internal/client/client_collapse.go b/internal/client/client_collapse.go index 4e753e92a0..922552c488 100644 --- a/internal/client/client_collapse.go +++ b/internal/client/client_collapse.go @@ -130,9 +130,13 @@ func (r reqCollapse) tryCollapseRequest(ctx context.Context, addr string, req *t func resolveLockCollapseKey(req *tikvrpc.Request) string { resolveLock := req.ResolveLock() - return strconv.FormatUint(req.RegionId, 10) + "-" + + key := strconv.FormatUint(req.RegionId, 10) + "-" + strconv.FormatUint(resolveLock.StartVersion, 10) + "-" + strconv.FormatBool(resolveLock.GetIsAsync()) + if resolveLock.GetIsTxnFile() { + key += "-f" + } + return key } func (r reqCollapse) collapse(ctx context.Context, key string, sf *singleflight.Group, diff --git a/internal/client/client_test.go b/internal/client/client_test.go index 4a28f347fd..f99d4ed363 100644 --- a/internal/client/client_test.go +++ b/internal/client/client_test.go @@ -227,13 +227,14 @@ func (c *chanClient) SendRequestAsync(ctx context.Context, addr string, req *tik } func TestCollapseResolveLock(t *testing.T) { - buildResolveLockReq := func(regionID uint64, startTS uint64, commitTS uint64, keys [][]byte, isAsync bool) *tikvrpc.Request { + buildResolveLockReq := func(regionID uint64, startTS uint64, commitTS uint64, keys [][]byte, isAsync bool, isTxnFile bool) *tikvrpc.Request { region := &metapb.Region{Id: regionID} req := tikvrpc.NewRequest(tikvrpc.CmdResolveLock, &kvrpcpb.ResolveLockRequest{ StartVersion: startTS, CommitVersion: commitTS, Keys: keys, IsAsync: isAsync, + IsTxnFile: isTxnFile, }) tikvrpc.SetContextNoAttach(req, region, nil) return req @@ -253,7 +254,7 @@ func TestCollapseResolveLock(t *testing.T) { ctx := context.Background() // Collapse ResolveLock. - resolveLockReq := buildResolveLockReq(1, 10, 20, nil, false) + resolveLockReq := buildResolveLockReq(1, 10, 20, nil, false, false) wg.Add(1) go client.SendRequest(ctx, "", resolveLockReq, time.Second) go client.SendRequest(ctx, "", resolveLockReq, time.Second) @@ -268,7 +269,7 @@ func TestCollapseResolveLock(t *testing.T) { } // Collapse async ResolveLock separately. - asyncResolveLockReq := buildResolveLockReq(1, 10, 20, nil, true) + asyncResolveLockReq := buildResolveLockReq(1, 10, 20, nil, true, false) wg.Add(1) go client.SendRequest(ctx, "", asyncResolveLockReq, time.Second) go client.SendRequest(ctx, "", asyncResolveLockReq, time.Second) @@ -292,8 +293,21 @@ func TestCollapseResolveLock(t *testing.T) { <-reqCh } + // Don't collapse regular ResolveLock with txn-file ResolveLock. + txnFileResolveLockReq := buildResolveLockReq(1, 10, 20, nil, false, true) + require.NotEqual(t, resolveLockCollapseKey(resolveLockReq), resolveLockCollapseKey(txnFileResolveLockReq)) + require.Equal(t, resolveLockCollapseKey(resolveLockReq)+"-f", resolveLockCollapseKey(txnFileResolveLockReq)) + wg.Add(1) + go client.SendRequest(ctx, "", resolveLockReq, time.Second) + go client.SendRequest(ctx, "", txnFileResolveLockReq, time.Second) + time.Sleep(300 * time.Millisecond) + wg.Done() + for i := 0; i < 2; i++ { + <-reqCh + } + // Don't collapse ResolveLockLite. - resolveLockLiteReq := buildResolveLockReq(1, 10, 20, [][]byte{[]byte("foo")}, false) + resolveLockLiteReq := buildResolveLockReq(1, 10, 20, [][]byte{[]byte("foo")}, false, false) wg.Add(1) go client.SendRequest(ctx, "", resolveLockLiteReq, time.Second) go client.SendRequest(ctx, "", resolveLockLiteReq, time.Second) diff --git a/txnkv/txnlock/lock_resolver.go b/txnkv/txnlock/lock_resolver.go index d60a3b5223..92321ad634 100644 --- a/txnkv/txnlock/lock_resolver.go +++ b/txnkv/txnlock/lock_resolver.go @@ -342,6 +342,9 @@ func (lr *LockResolver) BatchResolveLocks(bo *retry.Backoffer, locks []*Lock, lo resolveData, err := lr.checkAllSecondaries(bo, l, &status) if err == nil { txnInfos[l.TxnID] = resolveData.commitTs + if l.IsTxnFile { + txnFileIDs[l.TxnID] = true + } continue } if _, ok := errors.Cause(err).(*nonAsyncCommitLock); ok { From 37d99e22069e0df5c56a536d39d23b65f63f1e6a Mon Sep 17 00:00:00 2001 From: Ping Yu Date: Mon, 15 Jun 2026 10:07:40 +0800 Subject: [PATCH 3/3] revert IsTxnFile from collapse key Signed-off-by: Ping Yu --- internal/client/client_collapse.go | 7 ++----- internal/client/client_test.go | 22 ++++------------------ 2 files changed, 6 insertions(+), 23 deletions(-) diff --git a/internal/client/client_collapse.go b/internal/client/client_collapse.go index 922552c488..9875389ffa 100644 --- a/internal/client/client_collapse.go +++ b/internal/client/client_collapse.go @@ -129,14 +129,11 @@ func (r reqCollapse) tryCollapseRequest(ctx context.Context, addr string, req *t } func resolveLockCollapseKey(req *tikvrpc.Request) string { + // IsTxnFile is implied by StartVersion, so it does not need to be part of the collapse key. resolveLock := req.ResolveLock() - key := strconv.FormatUint(req.RegionId, 10) + "-" + + return strconv.FormatUint(req.RegionId, 10) + "-" + strconv.FormatUint(resolveLock.StartVersion, 10) + "-" + strconv.FormatBool(resolveLock.GetIsAsync()) - if resolveLock.GetIsTxnFile() { - key += "-f" - } - return key } func (r reqCollapse) collapse(ctx context.Context, key string, sf *singleflight.Group, diff --git a/internal/client/client_test.go b/internal/client/client_test.go index f99d4ed363..4a28f347fd 100644 --- a/internal/client/client_test.go +++ b/internal/client/client_test.go @@ -227,14 +227,13 @@ func (c *chanClient) SendRequestAsync(ctx context.Context, addr string, req *tik } func TestCollapseResolveLock(t *testing.T) { - buildResolveLockReq := func(regionID uint64, startTS uint64, commitTS uint64, keys [][]byte, isAsync bool, isTxnFile bool) *tikvrpc.Request { + buildResolveLockReq := func(regionID uint64, startTS uint64, commitTS uint64, keys [][]byte, isAsync bool) *tikvrpc.Request { region := &metapb.Region{Id: regionID} req := tikvrpc.NewRequest(tikvrpc.CmdResolveLock, &kvrpcpb.ResolveLockRequest{ StartVersion: startTS, CommitVersion: commitTS, Keys: keys, IsAsync: isAsync, - IsTxnFile: isTxnFile, }) tikvrpc.SetContextNoAttach(req, region, nil) return req @@ -254,7 +253,7 @@ func TestCollapseResolveLock(t *testing.T) { ctx := context.Background() // Collapse ResolveLock. - resolveLockReq := buildResolveLockReq(1, 10, 20, nil, false, false) + resolveLockReq := buildResolveLockReq(1, 10, 20, nil, false) wg.Add(1) go client.SendRequest(ctx, "", resolveLockReq, time.Second) go client.SendRequest(ctx, "", resolveLockReq, time.Second) @@ -269,7 +268,7 @@ func TestCollapseResolveLock(t *testing.T) { } // Collapse async ResolveLock separately. - asyncResolveLockReq := buildResolveLockReq(1, 10, 20, nil, true, false) + asyncResolveLockReq := buildResolveLockReq(1, 10, 20, nil, true) wg.Add(1) go client.SendRequest(ctx, "", asyncResolveLockReq, time.Second) go client.SendRequest(ctx, "", asyncResolveLockReq, time.Second) @@ -293,21 +292,8 @@ func TestCollapseResolveLock(t *testing.T) { <-reqCh } - // Don't collapse regular ResolveLock with txn-file ResolveLock. - txnFileResolveLockReq := buildResolveLockReq(1, 10, 20, nil, false, true) - require.NotEqual(t, resolveLockCollapseKey(resolveLockReq), resolveLockCollapseKey(txnFileResolveLockReq)) - require.Equal(t, resolveLockCollapseKey(resolveLockReq)+"-f", resolveLockCollapseKey(txnFileResolveLockReq)) - wg.Add(1) - go client.SendRequest(ctx, "", resolveLockReq, time.Second) - go client.SendRequest(ctx, "", txnFileResolveLockReq, time.Second) - time.Sleep(300 * time.Millisecond) - wg.Done() - for i := 0; i < 2; i++ { - <-reqCh - } - // Don't collapse ResolveLockLite. - resolveLockLiteReq := buildResolveLockReq(1, 10, 20, [][]byte{[]byte("foo")}, false, false) + resolveLockLiteReq := buildResolveLockReq(1, 10, 20, [][]byte{[]byte("foo")}, false) wg.Add(1) go client.SendRequest(ctx, "", resolveLockLiteReq, time.Second) go client.SendRequest(ctx, "", resolveLockLiteReq, time.Second)