Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions client/go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -44,3 +44,5 @@ require (
gopkg.in/natefinch/lumberjack.v2 v2.2.1 // indirect
gopkg.in/yaml.v3 v3.0.1 // indirect
)

replace github.com/pingcap/kvproto v0.0.0-20260903054228-107095f1d250 => github.com/JmPotato/kvproto v0.0.0-20260920123146-2a45fb4cd2dc
4 changes: 2 additions & 2 deletions client/go.sum
Original file line number Diff line number Diff line change
@@ -1,5 +1,7 @@
github.com/BurntSushi/toml v0.3.1 h1:WXkYYl6Yr3qBf1K79EBnL4mak0OimBfB0XUf9Vl28OQ=
github.com/BurntSushi/toml v0.3.1/go.mod h1:xHWCNGjB5oqiDr8zfno3MHue2Ht5sIBksp03qcyfWMU=
github.com/JmPotato/kvproto v0.0.0-20260920123146-2a45fb4cd2dc h1:7QHq/GlSggwICab7ETDp9ike7ED5XKLkTrfuJMlZasU=
github.com/JmPotato/kvproto v0.0.0-20260920123146-2a45fb4cd2dc/go.mod h1:z6+aAHB7dBkA+LyinEX+48/ImRJ3jag0Hg0c7wkhEvE=
github.com/benbjohnson/clock v1.1.0/go.mod h1:J11/hYXuz8f4ySSvYwY0FKfm+ezbsZBKZxNJlLklBHA=
github.com/benbjohnson/clock v1.3.0 h1:ip6w0uFQkncKQ979AypyG0ER7mqUSBdKLOgAle/AT8A=
github.com/benbjohnson/clock v1.3.0/go.mod h1:J11/hYXuz8f4ySSvYwY0FKfm+ezbsZBKZxNJlLklBHA=
Expand Down Expand Up @@ -53,8 +55,6 @@ github.com/pingcap/errors v0.11.5-0.20211224045212-9687c2b0f87c h1:xpW9bvK+HuuTm
github.com/pingcap/errors v0.11.5-0.20211224045212-9687c2b0f87c/go.mod h1:X2r9ueLEUZgtx2cIogM0v4Zj5uvvzhuuiu7Pn8HzMPg=
github.com/pingcap/failpoint v0.0.0-20240528011301-b51a646c7c86 h1:tdMsjOqUR7YXHoBitzdebTvOjs/swniBTOLy5XiMtuE=
github.com/pingcap/failpoint v0.0.0-20240528011301-b51a646c7c86/go.mod h1:exzhVYca3WRtd6gclGNErRWb1qEgff3LYta0LvRmON4=
github.com/pingcap/kvproto v0.0.0-20260903054228-107095f1d250 h1:6yUryXKVbKpCNdZWL58/OcZj8NPLUA/xsJYXSbsD59w=
github.com/pingcap/kvproto v0.0.0-20260903054228-107095f1d250/go.mod h1:z6+aAHB7dBkA+LyinEX+48/ImRJ3jag0Hg0c7wkhEvE=
github.com/pingcap/log v1.1.1-0.20221110025148-ca232912c9f3 h1:HR/ylkkLmGdSSDaD8IDP+SZrdhV1Kibl9KrHxJ9eciw=
github.com/pingcap/log v1.1.1-0.20221110025148-ca232912c9f3/go.mod h1:DWQW5jICDR7UJh4HtxXSM20Churx4CQL0fwL/SoOSA4=
github.com/pkg/errors v0.8.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0=
Expand Down
2 changes: 2 additions & 0 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -264,3 +264,5 @@ require (
// which will cause several different tests to fail. So this is a temporary workaround to use the old version of `testify`.
// TODO: fix those flasky tests introduced by the behavior change of `Eventually` and `EventuallyWithT` assertions.
replace github.com/stretchr/testify => github.com/stretchr/testify v1.10.0

replace github.com/pingcap/kvproto v0.0.0-20260903054228-107095f1d250 => github.com/JmPotato/kvproto v0.0.0-20260920123146-2a45fb4cd2dc
4 changes: 2 additions & 2 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,8 @@ github.com/AzureAD/microsoft-authentication-library-for-go v1.6.0/go.mod h1:HKpQ
github.com/BurntSushi/toml v0.3.1/go.mod h1:xHWCNGjB5oqiDr8zfno3MHue2Ht5sIBksp03qcyfWMU=
github.com/BurntSushi/toml v1.5.0 h1:W5quZX/G/csjUnuI8SUYlsHs9M38FC7znL0lIO+DvMg=
github.com/BurntSushi/toml v1.5.0/go.mod h1:ukJfTF/6rtPPRCnwkur4qwRxa8vTRFBF0uk2lLoLwho=
github.com/JmPotato/kvproto v0.0.0-20260920123146-2a45fb4cd2dc h1:7QHq/GlSggwICab7ETDp9ike7ED5XKLkTrfuJMlZasU=
github.com/JmPotato/kvproto v0.0.0-20260920123146-2a45fb4cd2dc/go.mod h1:z6+aAHB7dBkA+LyinEX+48/ImRJ3jag0Hg0c7wkhEvE=
github.com/KyleBanks/depth v1.2.1 h1:5h8fQADFrWtarTdtDudMmGsC7GPbOAu6RVB3ffsVFHc=
github.com/KyleBanks/depth v1.2.1/go.mod h1:jzSb9d0L43HxTQfT+oSA1EEp2q+ne2uh6XgeJcm8brE=
github.com/Masterminds/semver v1.5.0 h1:H65muMkzWKEuNDnfl9d70GUjFniHKHRbFPGBuZ3QEww=
Expand Down Expand Up @@ -490,8 +492,6 @@ github.com/pingcap/errors v0.11.5-0.20211224045212-9687c2b0f87c/go.mod h1:X2r9ue
github.com/pingcap/failpoint v0.0.0-20240528011301-b51a646c7c86 h1:tdMsjOqUR7YXHoBitzdebTvOjs/swniBTOLy5XiMtuE=
github.com/pingcap/failpoint v0.0.0-20240528011301-b51a646c7c86/go.mod h1:exzhVYca3WRtd6gclGNErRWb1qEgff3LYta0LvRmON4=
github.com/pingcap/kvproto v0.0.0-20191211054548-3c6b38ea5107/go.mod h1:WWLmULLO7l8IOcQG+t+ItJ3fEcrL5FxF0Wu+HrMy26w=
github.com/pingcap/kvproto v0.0.0-20260903054228-107095f1d250 h1:6yUryXKVbKpCNdZWL58/OcZj8NPLUA/xsJYXSbsD59w=
github.com/pingcap/kvproto v0.0.0-20260903054228-107095f1d250/go.mod h1:z6+aAHB7dBkA+LyinEX+48/ImRJ3jag0Hg0c7wkhEvE=
github.com/pingcap/log v0.0.0-20210625125904-98ed8e2eb1c7/go.mod h1:8AanEdAHATuRurdGxZXBz0At+9avep+ub7U1AGYLIMM=
github.com/pingcap/log v1.1.1-0.20221110025148-ca232912c9f3 h1:HR/ylkkLmGdSSDaD8IDP+SZrdhV1Kibl9KrHxJ9eciw=
github.com/pingcap/log v1.1.1-0.20221110025148-ca232912c9f3/go.mod h1:DWQW5jICDR7UJh4HtxXSM20Churx4CQL0fwL/SoOSA4=
Expand Down
2 changes: 1 addition & 1 deletion pkg/mcs/resourcemanager/server/grpc_service.go
Original file line number Diff line number Diff line change
Expand Up @@ -244,7 +244,7 @@ func (s *Service) AcquireTokenBuckets(stream rmpb.ResourceManager_AcquireTokenBu
continue
}
// Send the consumption to update the metrics.
err = s.manager.dispatchConsumption(req)
err = s.manager.dispatchConsumption(clientUniqueID, req)
if err != nil {
return err
}
Expand Down
1 change: 1 addition & 0 deletions pkg/mcs/resourcemanager/server/keyspace_manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -58,6 +58,7 @@ const (

// consumptionItem is used to send the consumption info to the background metrics flusher.
type consumptionItem struct {
clientUniqueID uint64
keyspaceID uint32
keyspaceName string
resourceGroupName string
Expand Down
11 changes: 10 additions & 1 deletion pkg/mcs/resourcemanager/server/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -111,6 +111,9 @@ type Manager struct {
keyspaceIDLookup map[string]uint32
// metrics is the collection of metrics.
metrics *metrics
// metricsLoopMu serializes backgroundMetricsFlush, since a new leadership
// term can start before the previous flusher exits.
metricsLoopMu sync.Mutex
// ruCollector is used to collect the RU metering data.
ruCollector *ruCollector
}
Expand Down Expand Up @@ -659,13 +662,14 @@ func (m *Manager) getKeyspaceResourceGroupManagers() []*keyspaceResourceGroupMan
return krgms
}

func (m *Manager) dispatchConsumption(req *rmpb.TokenBucketRequest) error {
func (m *Manager) dispatchConsumption(clientUniqueID uint64, req *rmpb.TokenBucketRequest) error {
isBackground := req.GetIsBackground()
isTiFlash := req.GetIsTiflash()
if isBackground && isTiFlash {
return errors.New("background and tiflash cannot be true at the same time")
}
m.consumptionDispatcher <- &consumptionItem{
clientUniqueID: clientUniqueID,
keyspaceID: ExtractKeyspaceID(req.GetKeyspaceId()),
resourceGroupName: req.GetResourceGroupName(),
Consumption: req.GetConsumptionSinceLastRequest(),
Expand Down Expand Up @@ -752,6 +756,10 @@ func (m *Manager) GetKeyspaceIDByName(ctx context.Context, name string) (*rmpb.K
func (m *Manager) backgroundMetricsFlush(ctx context.Context) {
defer logutil.LogPanic()
defer m.wg.Done()
m.metricsLoopMu.Lock()
defer m.metricsLoopMu.Unlock()
// The RU timeline belongs to one leadership term.
defer m.metrics.ruTimeline.reset()
cleanUpTicker := time.NewTicker(metricsCleanupInterval)
defer cleanUpTicker.Stop()
metricsTicker := time.NewTicker(tickPerSecond)
Expand Down Expand Up @@ -832,6 +840,7 @@ func (m *Manager) backgroundMetricsFlush(ctx context.Context) {
}
}
case <-metricsTicker.C:
m.metrics.ruTimeline.flush(time.Now())
// Prevent from holding the lock too long when there're many keyspaces and resource groups.
for _, krgm := range m.getKeyspaceResourceGroupManagers() {
// Conciliate the fill rates.
Expand Down
5 changes: 3 additions & 2 deletions pkg/mcs/resourcemanager/server/manager_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -460,7 +460,7 @@ func checkBackgroundMetricsFlush(ctx context.Context, re *require.Assertions, ma
},
KeyspaceId: keyspaceIDValue,
}
err = manager.dispatchConsumption(req)
err = manager.dispatchConsumption(1, req)
re.NoError(err)

keyspaceID := ExtractKeyspaceID(req.GetKeyspaceId())
Expand All @@ -487,10 +487,11 @@ func TestDispatchConsumptionIncludesOnlyConsumption(t *testing.T) {
KeyspaceId: &rmpb.KeyspaceIDValue{Keyspace: &rmpb.KeyspaceIDValue_Value{Value: 42}},
}

err := m.dispatchConsumption(req)
err := m.dispatchConsumption(7, req)
re.NoError(err)

item := <-m.consumptionDispatcher
re.Equal(uint64(7), item.clientUniqueID)
re.Equal(uint32(42), item.keyspaceID)
re.Equal(req.GetResourceGroupName(), item.resourceGroupName)
re.Equal(req.GetConsumptionSinceLastRequest(), item.Consumption)
Expand Down
3 changes: 3 additions & 0 deletions pkg/mcs/resourcemanager/server/metrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -219,6 +219,7 @@ var (
)

type metrics struct {
ruTimeline *ruTimeline
// record update time of each resource group
consumptionRecordMap map[consumptionRecordKey]time.Time
// max per sec trackers for each keyspace and resource group.
Expand Down Expand Up @@ -273,6 +274,7 @@ func init() {

func newMetrics() *metrics {
return &metrics{
ruTimeline: newRUTimeline(ruSummaryMetrics),
consumptionRecordMap: make(map[consumptionRecordKey]time.Time),
maxPerSecTrackerMap: make(map[trackerKey]*maxPerSecCostTracker),
counterMetricsMap: make(map[metricsKey]*counterMetrics),
Expand Down Expand Up @@ -348,6 +350,7 @@ func (m *metrics) recordConsumption(
if consumption == nil {
return
}
m.ruTimeline.record(consumptionInfo, now)
m.getMaxPerSecTracker(keyspaceID, keyspaceName, groupName).collect(consumption)
m.getCounterMetrics(keyspaceID, keyspaceName, groupName, ruLabelType).add(consumption, controllerConfig, keyspaceID)
m.insertConsumptionRecord(keyspaceID, groupName, ruLabelType, now)
Expand Down
13 changes: 13 additions & 0 deletions pkg/mcs/resourcemanager/server/metrics_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -150,6 +150,19 @@ func TestRecordConsumptionUsesActualRequestUnitCounters(t *testing.T) {
groupName: groupName,
ruType: defaultTypeLabel,
})

// Client RU timelines are merged without changing counter accounting.
now := time.Now()
report := timelineReport(1, now.Unix(), [][2]float64{{12, 8}})
report.keyspaceID, report.keyspaceName = keyspaceID, keyspaceName
report.resourceGroupName = groupName
m.recordConsumption(report, &ControllerConfig{}, now.Add(time.Second))
source := m.ruTimeline.groups[trackerKey{keyspaceID, groupName}].sources[ruSourceKey{client: 1}]
bucket := source.buckets[now.Unix()%int64(len(source.buckets))]
re.Equal(float64(12), bucket.rru)
re.Equal(float64(8), bucket.wru)
re.Equal(float64(12), testutil.ToFloat64(counter.RRUMetrics))
re.Equal(float64(8), testutil.ToFloat64(counter.WRUMetrics))
}

func TestRecordConsumptionKeepsRecordForEmptyConsumption(t *testing.T) {
Expand Down
Loading
Loading