From 625c0206a37a96157a0f002633d564123be3f809 Mon Sep 17 00:00:00 2001 From: Yuhao Zhang Date: Fri, 17 Jul 2026 13:04:04 +0800 Subject: [PATCH] resource-manager: price remote read bytes separately Signed-off-by: Yuhao Zhang --- client/resource_group/controller/config.go | 11 ++ client/resource_group/controller/model.go | 23 +++- .../resource_group/controller/model_test.go | 115 ++++++++++++++++++ client/resource_group/controller/testutil.go | 12 +- pkg/mcs/resourcemanager/server/config.go | 3 + pkg/mcs/resourcemanager/server/config_test.go | 17 +++ pkg/mcs/resourcemanager/server/manager.go | 4 + .../resourcemanager/server/manager_test.go | 24 +++- .../resourcemanager/resource_manager_test.go | 2 + 9 files changed, 202 insertions(+), 9 deletions(-) diff --git a/client/resource_group/controller/config.go b/client/resource_group/controller/config.go index 3a1302c9b3d..de2b8ec1f2a 100644 --- a/client/resource_group/controller/config.go +++ b/client/resource_group/controller/config.go @@ -59,6 +59,8 @@ const ( defaultWritePerBatchBaseCost = 1 // 1 RU = 64 KiB read bytes defaultReadCostPerByte = 1. / (64 * 1024) + // Remote read bytes cost half as much as local read bytes by default. + defaultRemoteReadCostFactor = 0.5 // 1 RU = 1 KiB written bytes defaultWriteCostPerByte = 1. / 1024 // 1 RU = 3 millisecond CPU time @@ -176,6 +178,9 @@ type RequestUnitConfig struct { ReadPerBatchBaseCost float64 `toml:"read-per-batch-base-cost" json:"read-per-batch-base-cost"` // ReadCostPerByte is the cost for each byte read. It's 1 RU = 64 KiB by default. ReadCostPerByte float64 `toml:"read-cost-per-byte" json:"read-cost-per-byte"` + // ReadCostPerByteRemote is the cost for each byte processed by a remote coprocessor. + // Nil means half of ReadCostPerByte. A pointer preserves an explicitly configured zero. + ReadCostPerByteRemote *float64 `toml:"read-cost-per-byte-remote" json:"read-cost-per-byte-remote,omitempty"` // WriteBaseCost is the base cost for a write request. No matter how many bytes read/written or // the CPU times taken for a request, this cost is inevitable. WriteBaseCost float64 `toml:"write-base-cost" json:"write-base-cost"` @@ -209,6 +214,7 @@ type RUConfig struct { ReadBaseCost RequestUnit ReadPerBatchBaseCost RequestUnit ReadBytesCost RequestUnit + RemoteReadBytesCost RequestUnit WriteBaseCost RequestUnit WritePerBatchBaseCost RequestUnit WriteBytesCost RequestUnit @@ -232,10 +238,15 @@ func DefaultRUConfig() *RUConfig { // GenerateRUConfig generates the configuration by the given request unit configuration. func GenerateRUConfig(config *Config) *RUConfig { + remoteReadCostPerByte := config.RequestUnit.ReadCostPerByte * defaultRemoteReadCostFactor + if config.RequestUnit.ReadCostPerByteRemote != nil { + remoteReadCostPerByte = *config.RequestUnit.ReadCostPerByteRemote + } return &RUConfig{ ReadBaseCost: RequestUnit(config.RequestUnit.ReadBaseCost), ReadPerBatchBaseCost: RequestUnit(config.RequestUnit.ReadPerBatchBaseCost), ReadBytesCost: RequestUnit(config.RequestUnit.ReadCostPerByte), + RemoteReadBytesCost: RequestUnit(remoteReadCostPerByte), WriteBaseCost: RequestUnit(config.RequestUnit.WriteBaseCost), WritePerBatchBaseCost: RequestUnit(config.RequestUnit.WritePerBatchBaseCost), WriteBytesCost: RequestUnit(config.RequestUnit.WriteCostPerByte), diff --git a/client/resource_group/controller/model.go b/client/resource_group/controller/model.go index 68edc1a1380..926e2311207 100644 --- a/client/resource_group/controller/model.go +++ b/client/resource_group/controller/model.go @@ -118,6 +118,20 @@ type ResponseInfo interface { ResponseSize() uint64 } +type remoteReadBytesProvider interface { + // RemoteReadBytes returns the factual subset of ReadBytes that was processed + // by a remote coprocessor. + RemoteReadBytes() uint64 +} + +func getRemoteReadBytes(res ResponseInfo, totalReadBytes uint64) uint64 { + provider, ok := res.(remoteReadBytesProvider) + if !ok { + return 0 + } + return min(provider.RemoteReadBytes(), totalReadBytes) +} + // ResourceCalculator is used to calculate the resource consumption of a request. type ResourceCalculator interface { // Trickle is used to calculate the resource consumption periodically rather than on the request path. @@ -213,9 +227,12 @@ func (kc *KVCalculator) calculateReadCost(consumption *rmpb.Consumption, res Res if consumption == nil { return } - readBytes := float64(res.ReadBytes()) - consumption.ReadBytes += readBytes - consumption.RRU += float64(kc.ReadBytesCost) * readBytes + totalReadBytes := res.ReadBytes() + remoteBytes := getRemoteReadBytes(res, totalReadBytes) + localReadBytes := totalReadBytes - remoteBytes + consumption.ReadBytes += float64(totalReadBytes) + consumption.RRU += float64(kc.ReadBytesCost)*float64(localReadBytes) + + float64(kc.RemoteReadBytesCost)*float64(remoteBytes) } func (kc *KVCalculator) calculateCPUCost(consumption *rmpb.Consumption, res ResponseInfo) { diff --git a/client/resource_group/controller/model_test.go b/client/resource_group/controller/model_test.go index 6f31ed5e64e..ff7a5c213bf 100644 --- a/client/resource_group/controller/model_test.go +++ b/client/resource_group/controller/model_test.go @@ -16,6 +16,7 @@ package controller import ( "testing" + "time" "github.com/stretchr/testify/require" @@ -58,6 +59,26 @@ func (*copRequestInfoWithoutPrediction) IsCop() bool { return true } +type legacyResponseInfo struct { + readBytes uint64 +} + +func (res *legacyResponseInfo) ReadBytes() uint64 { + return res.readBytes +} + +func (*legacyResponseInfo) KVCPU() time.Duration { + return 0 +} + +func (*legacyResponseInfo) Succeed() bool { + return true +} + +func (res *legacyResponseInfo) ResponseSize() uint64 { + return res.readBytes +} + func TestGetRUValueFromConsumption(t *testing.T) { // Positive test case re := require.New(t) @@ -261,6 +282,100 @@ func TestRequestInfoMissingPredictionProviderReturnsZeroHint(t *testing.T) { re.Zero(bytesForEst) } +func TestGenerateRUConfigRemoteReadCost(t *testing.T) { + re := require.New(t) + + config := DefaultConfig() + config.RequestUnit.ReadCostPerByte = 2 + config.RequestUnit.ReadCostPerByteRemote = nil + ruConfig := GenerateRUConfig(config) + re.InDelta(1.0, float64(ruConfig.RemoteReadBytesCost), 1e-7) + + explicitZero := 0.0 + config.RequestUnit.ReadCostPerByteRemote = &explicitZero + ruConfig = GenerateRUConfig(config) + re.Zero(ruConfig.RemoteReadBytesCost) +} + +func TestCalculateReadCostSplitsLocalAndRemoteBytes(t *testing.T) { + ruConfig := DefaultRUConfig() + ruConfig.ReadBytesCost = 2 + ruConfig.RemoteReadBytesCost = 0.5 + kvCalc := newKVCalculator(ruConfig) + + testCases := []struct { + name string + response ResponseInfo + expectedRRU float64 + }{ + { + name: "legacy response charges all bytes at normal rate", + response: &legacyResponseInfo{readBytes: 100}, + expectedRRU: 200, + }, + { + name: "remote subset uses remote rate", + response: &TestResponseInfo{ + readBytes: 100, + remoteReadBytes: 40, + succeed: true, + }, + expectedRRU: 140, + }, + { + name: "remote subset is clamped to total bytes", + response: &TestResponseInfo{ + readBytes: 100, + remoteReadBytes: 200, + succeed: true, + }, + expectedRRU: 50, + }, + } + + for _, testCase := range testCases { + t.Run(testCase.name, func(t *testing.T) { + re := require.New(t) + consumption := &rmpb.Consumption{} + kvCalc.calculateReadCost(consumption, testCase.response) + re.Equal(float64(100), consumption.ReadBytes) + re.InDelta(testCase.expectedRRU, consumption.RRU, 1e-7) + }) + } +} + +func TestRemoteReadCostPreservesPagingSettlementFormula(t *testing.T) { + re := require.New(t) + ruConfig := DefaultRUConfig() + ruConfig.ReadBytesCost = 2 + ruConfig.RemoteReadBytesCost = 0.5 + ruConfig.CPUMsCost = 3 + kvCalc := newKVCalculator(ruConfig) + req := &TestRequestInfo{ + isCop: true, + predictedReadBytes: 80, + } + resp := &TestResponseInfo{ + readBytes: 100, + remoteReadBytes: 40, + kvCPU: 10 * time.Millisecond, + succeed: true, + } + + consumption := &rmpb.Consumption{} + kvCalc.BeforeKVRequest(consumption, req) + kvCalc.AfterKVRequest(consumption, req, resp) + + baseCost := float64(ruConfig.ReadBaseCost) + + float64(ruConfig.ReadPerBatchBaseCost)*defaultAvgBatchProportion + expectedByteCost := 60*float64(ruConfig.ReadBytesCost) + + 40*float64(ruConfig.RemoteReadBytesCost) + expectedCPUCost := 10 * float64(ruConfig.CPUMsCost) + re.InDelta(baseCost+expectedByteCost+expectedCPUCost, consumption.RRU, 1e-7) + re.Equal(float64(100), consumption.ReadBytes) + re.Equal(float64(10), consumption.TotalCpuTimeMs) +} + func TestReportedConsumptionStripsPagingPrecharge(t *testing.T) { re := require.New(t) cfg := DefaultRUConfig() diff --git a/client/resource_group/controller/testutil.go b/client/resource_group/controller/testutil.go index 2284273ec61..96fd703b67d 100644 --- a/client/resource_group/controller/testutil.go +++ b/client/resource_group/controller/testutil.go @@ -84,9 +84,10 @@ func (tri *TestRequestInfo) IsCop() bool { // TestResponseInfo is used to test the response info interface. type TestResponseInfo struct { - readBytes uint64 - kvCPU time.Duration - succeed bool + readBytes uint64 + remoteReadBytes uint64 + kvCPU time.Duration + succeed bool } // NewTestResponseInfo creates a new TestResponseInfo. @@ -103,6 +104,11 @@ func (tri *TestResponseInfo) ReadBytes() uint64 { return tri.readBytes } +// RemoteReadBytes implements the optional remoteReadBytesProvider interface. +func (tri *TestResponseInfo) RemoteReadBytes() uint64 { + return tri.remoteReadBytes +} + // KVCPU implements the ResponseInfo interface. func (tri *TestResponseInfo) KVCPU() time.Duration { return tri.kvCPU diff --git a/pkg/mcs/resourcemanager/server/config.go b/pkg/mcs/resourcemanager/server/config.go index d3ff108cdd6..af153a0dbef 100644 --- a/pkg/mcs/resourcemanager/server/config.go +++ b/pkg/mcs/resourcemanager/server/config.go @@ -225,6 +225,9 @@ type RequestUnitConfig struct { ReadPerBatchBaseCost float64 `toml:"read-per-batch-base-cost" json:"read-per-batch-base-cost"` // ReadCostPerByte is the cost for each byte read. It's 1 RU = 64 KiB by default. ReadCostPerByte float64 `toml:"read-cost-per-byte" json:"read-cost-per-byte"` + // ReadCostPerByteRemote is the cost for each byte processed by a remote coprocessor. + // Nil means the client controller should use half of ReadCostPerByte. + ReadCostPerByteRemote *float64 `toml:"read-cost-per-byte-remote" json:"read-cost-per-byte-remote,omitempty"` // WriteBaseCost is the base cost for a write request. No matter how many bytes read/written or // the CPU times taken for a request, this cost is inevitable. WriteBaseCost float64 `toml:"write-base-cost" json:"write-base-cost"` diff --git a/pkg/mcs/resourcemanager/server/config_test.go b/pkg/mcs/resourcemanager/server/config_test.go index afe22356def..a1853c9ceab 100644 --- a/pkg/mcs/resourcemanager/server/config_test.go +++ b/pkg/mcs/resourcemanager/server/config_test.go @@ -50,5 +50,22 @@ read-cpu-ms-cost = 5.0 re.LessOrEqual(math.Abs(cfg.Controller.RequestUnit.WriteCostPerByte-4), 1e-7) re.LessOrEqual(math.Abs(cfg.Controller.RequestUnit.WriteBaseCost-3), 1e-7) re.LessOrEqual(math.Abs(cfg.Controller.RequestUnit.ReadCostPerByte-2), 1e-7) + re.Nil(cfg.Controller.RequestUnit.ReadCostPerByteRemote) re.LessOrEqual(math.Abs(cfg.Controller.RequestUnit.ReadBaseCost-1), 1e-7) } + +func TestControllerConfigRemoteReadCostOverride(t *testing.T) { + re := require.New(t) + cfgData := ` +[controller.request-unit] +read-cost-per-byte = 2.0 +read-cost-per-byte-remote = 0.25 +` + cfg := NewConfig() + meta, err := toml.Decode(cfgData, &cfg) + re.NoError(err) + re.NoError(cfg.Adjust(&meta)) + + re.NotNil(cfg.Controller.RequestUnit.ReadCostPerByteRemote) + re.InDelta(0.25, *cfg.Controller.RequestUnit.ReadCostPerByteRemote, 1e-7) +} diff --git a/pkg/mcs/resourcemanager/server/manager.go b/pkg/mcs/resourcemanager/server/manager.go index 4f811a0392e..e079b190537 100644 --- a/pkg/mcs/resourcemanager/server/manager.go +++ b/pkg/mcs/resourcemanager/server/manager.go @@ -427,6 +427,10 @@ func cloneControllerConfig(cfg *ControllerConfig) *ControllerConfig { return nil } cloned := *cfg + if cfg.RequestUnit.ReadCostPerByteRemote != nil { + remoteReadCost := *cfg.RequestUnit.ReadCostPerByteRemote + cloned.RequestUnit.ReadCostPerByteRemote = &remoteReadCost + } cloned.RUVersionPolicy = cfg.RUVersionPolicy.Clone() return &cloned } diff --git a/pkg/mcs/resourcemanager/server/manager_test.go b/pkg/mcs/resourcemanager/server/manager_test.go index 16ecac456d1..c0000435519 100644 --- a/pkg/mcs/resourcemanager/server/manager_test.go +++ b/pkg/mcs/resourcemanager/server/manager_test.go @@ -269,16 +269,20 @@ func TestManagerControllerConfigSnapshots(t *testing.T) { re := require.New(t) m := prepareManager() + remoteReadCost := 0.25 m.controllerConfig = &ControllerConfig{ RequestUnit: RequestUnitConfig{ - ReadBaseCost: 0.5, + ReadBaseCost: 0.5, + ReadCostPerByteRemote: &remoteReadCost, }, } snapshot := m.GetControllerConfig() snapshot.RequestUnit.ReadBaseCost = 1.5 + *snapshot.RequestUnit.ReadCostPerByteRemote = 0.75 re.InDelta(0.5, m.controllerConfig.RequestUnit.ReadBaseCost, 0.00001) + re.InDelta(0.25, *m.controllerConfig.RequestUnit.ReadCostPerByteRemote, 0.00001) }) t.Run("publishes_new_snapshot_after_successful_update", func(t *testing.T) { @@ -307,17 +311,31 @@ func TestManagerControllerConfigSnapshots(t *testing.T) { Storage: storage.NewStorageWithMemoryBackend(), err: expectedErr, } + remoteReadCost := 0.25 m.controllerConfig = &ControllerConfig{ RequestUnit: RequestUnitConfig{ - ReadBaseCost: 0.5, + ReadBaseCost: 0.5, + ReadCostPerByteRemote: &remoteReadCost, }, } previous := m.controllerConfig - err := m.UpdateControllerConfigItem("request-unit.read-base-cost", 1.5) + err := m.UpdateControllerConfigItem("request-unit.read-cost-per-byte-remote", 0.75) re.ErrorIs(err, expectedErr) re.Same(previous, m.controllerConfig) re.InDelta(0.5, m.controllerConfig.RequestUnit.ReadBaseCost, 0.00001) + re.InDelta(0.25, *m.controllerConfig.RequestUnit.ReadCostPerByteRemote, 0.00001) + }) + + t.Run("updates_optional_remote_read_cost", func(t *testing.T) { + re := require.New(t) + + m := prepareManager() + m.controllerConfig = &ControllerConfig{} + + re.NoError(m.UpdateControllerConfigItem("request-unit.read-cost-per-byte-remote", 0.25)) + re.NotNil(m.controllerConfig.RequestUnit.ReadCostPerByteRemote) + re.InDelta(0.25, *m.controllerConfig.RequestUnit.ReadCostPerByteRemote, 0.00001) }) } diff --git a/tests/integrations/mcs/resourcemanager/resource_manager_test.go b/tests/integrations/mcs/resourcemanager/resource_manager_test.go index d30a5cb3726..e5c3c5c2da5 100644 --- a/tests/integrations/mcs/resourcemanager/resource_manager_test.go +++ b/tests/integrations/mcs/resourcemanager/resource_manager_test.go @@ -1791,6 +1791,7 @@ func (suite *resourceManagerClientTestSuite) TestLoadRequestUnitConfig() { expectedConfig := controller.DefaultRUConfig() re.Equal(expectedConfig.ReadBaseCost, config.ReadBaseCost) re.Equal(expectedConfig.ReadBytesCost, config.ReadBytesCost) + re.Equal(expectedConfig.RemoteReadBytesCost, config.RemoteReadBytesCost) re.Equal(expectedConfig.WriteBaseCost, config.WriteBaseCost) re.Equal(expectedConfig.WriteBytesCost, config.WriteBytesCost) re.Equal(expectedConfig.CPUMsCost, config.CPUMsCost) @@ -1811,6 +1812,7 @@ func (suite *resourceManagerClientTestSuite) TestLoadRequestUnitConfig() { expectedConfig = controller.GenerateRUConfig(controllerConfig) re.Equal(expectedConfig.ReadBaseCost, config.ReadBaseCost) re.Equal(expectedConfig.ReadBytesCost, config.ReadBytesCost) + re.Equal(expectedConfig.RemoteReadBytesCost, config.RemoteReadBytesCost) re.Equal(expectedConfig.WriteBaseCost, config.WriteBaseCost) re.Equal(expectedConfig.WriteBytesCost, config.WriteBytesCost) re.Equal(expectedConfig.CPUMsCost, config.CPUMsCost)