Skip to content
Open
Show file tree
Hide file tree
Changes from 5 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
19 changes: 15 additions & 4 deletions pkg/schedule/filter/filters.go
Original file line number Diff line number Diff line change
Expand Up @@ -703,13 +703,17 @@ type ruleLeaderFitFilter struct {

// newRuleLeaderFitFilter creates a filter that ensures after transfer leader with new store,
// the isolation level will not decrease.
func newRuleLeaderFitFilter(scope string, cluster *core.BasicCluster, ruleManager *placement.RuleManager, region *core.RegionInfo, srcLeaderStoreID uint64, allowMoveLeader bool) Filter {
func newRuleLeaderFitFilter(scope string, cluster *core.BasicCluster, ruleManager *placement.RuleManager,
region *core.RegionInfo, oldFit *placement.RegionFit, srcLeaderStoreID uint64, allowMoveLeader bool) Filter {
if oldFit == nil {
oldFit = ruleManager.FitRegion(cluster, region)
}
return &ruleLeaderFitFilter{
scope: scope,
cluster: cluster,
ruleManager: ruleManager,
region: region,
oldFit: ruleManager.FitRegion(cluster, region),
oldFit: oldFit,
srcLeaderStoreID: srcLeaderStoreID,
allowMoveLeader: allowMoveLeader,
}
Expand Down Expand Up @@ -816,9 +820,16 @@ func NewPlacementSafeguard(scope string, conf config.SharedConfigProvider, clust
// NewPlacementLeaderSafeguard creates a filter that ensures after transfer a leader with
// existed peer, the placement restriction will not become worse.
// Note that it only worked when PlacementRules enabled otherwise it will always permit the sourceStore.
func NewPlacementLeaderSafeguard(scope string, conf config.SharedConfigProvider, cluster *core.BasicCluster, ruleManager *placement.RuleManager, region *core.RegionInfo, sourceStore *core.StoreInfo, allowMoveLeader bool) Filter {
func NewPlacementLeaderSafeguard(scope string, conf config.SharedConfigProvider, cluster *core.BasicCluster, ruleManager *placement.RuleManager,
region *core.RegionInfo, sourceStore *core.StoreInfo, allowMoveLeader bool) Filter {
return NewPlacementLeaderSafeguardWithFit(scope, conf, cluster, ruleManager, region, sourceStore, nil, allowMoveLeader)
}

// NewPlacementLeaderSafeguardWithFit creates a placement leader safeguard with a reusable region fit.
func NewPlacementLeaderSafeguardWithFit(scope string, conf config.SharedConfigProvider, cluster *core.BasicCluster, ruleManager *placement.RuleManager,
region *core.RegionInfo, sourceStore *core.StoreInfo, oldFit *placement.RegionFit, allowMoveLeader bool) Filter {
if conf.IsPlacementRulesEnabled() {
return newRuleLeaderFitFilter(scope, cluster, ruleManager, region, sourceStore.GetID(), allowMoveLeader)
return newRuleLeaderFitFilter(scope, cluster, ruleManager, region, oldFit, sourceStore.GetID(), allowMoveLeader)
}
return nil
}
Expand Down
6 changes: 3 additions & 3 deletions pkg/schedule/filter/filters_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -144,12 +144,12 @@ func TestRuleFitFilter(t *testing.T) {
filter := newRuleFitFilter("", testCluster.GetBasicCluster(), testCluster.GetRuleManager(), region, nil, 1)
re.Equal(testCase.sourceRes, filter.Source(testCluster.GetSharedConfig(), testCluster.GetStore(testCase.storeID)).StatusCode)
re.Equal(testCase.targetRes, filter.Target(testCluster.GetSharedConfig(), testCluster.GetStore(testCase.storeID)).StatusCode)
leaderFilter := newRuleLeaderFitFilter("", testCluster.GetBasicCluster(), testCluster.GetRuleManager(), region, 1, true)
leaderFilter := newRuleLeaderFitFilter("", testCluster.GetBasicCluster(), testCluster.GetRuleManager(), region, nil, 1, true)
re.Equal(testCase.targetRes, leaderFilter.Target(testCluster.GetSharedConfig(), testCluster.GetStore(testCase.storeID)).StatusCode)
}

// store-6 is not exist in the peers, so it will not allow transferring leader to store 6.
leaderFilter := newRuleLeaderFitFilter("", testCluster.GetBasicCluster(), testCluster.GetRuleManager(), region, 1, false)
leaderFilter := newRuleLeaderFitFilter("", testCluster.GetBasicCluster(), testCluster.GetRuleManager(), region, nil, 1, false)
re.False(leaderFilter.Target(testCluster.GetSharedConfig(), testCluster.GetStore(6)).IsOK())
}

Expand Down Expand Up @@ -223,7 +223,7 @@ func TestRuleFitFilterWithPlacementRule(t *testing.T) {
{StoreId: 5, Id: 4},
{StoreId: 6, Id: 5},
}}, &metapb.Peer{StoreId: 1, Id: 1})
leaderFilter := newRuleLeaderFitFilter("", testCluster.GetBasicCluster(), testCluster.GetRuleManager(), region, 1, true)
leaderFilter := newRuleLeaderFitFilter("", testCluster.GetBasicCluster(), testCluster.GetRuleManager(), region, nil, 1, true)
re.Equal(plan.StatusText(plan.StatusStoreNotMatchRule), leaderFilter.Target(testCluster.GetSharedConfig(), testCluster.GetStore(6)).String())
}

Expand Down
45 changes: 43 additions & 2 deletions pkg/schedule/placement/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,12 +18,15 @@ import (
"bytes"
"encoding/json"
"time"

"github.com/tikv/pd/pkg/core"
)

// ruleConfig contains rule and rule group configurations.
type ruleConfig struct {
rules map[[2]string]*Rule // {group, id} => Rule
groups map[string]*RuleGroup // id => RuleGroup
rules map[[2]string]*Rule // {group, id} => Rule
groups map[string]*RuleGroup // id => RuleGroup
mayRestrictStoreLoad [2]bool
}

func newRuleConfig() *ruleConfig {
Expand All @@ -35,6 +38,7 @@ func newRuleConfig() *ruleConfig {

// adjust configs for `buildRuleList` and API use.
func (c *ruleConfig) adjust() {
c.mayRestrictStoreLoad = [2]bool{}
// remove all default group configurations.
// if there are rules belong to the group, it will be re-add later.
for id, g := range c.groups {
Expand All @@ -43,6 +47,10 @@ func (c *ruleConfig) adjust() {
}
}
for _, r := range c.rules {
if len(r.LabelConstraints) > 0 {
c.mayRestrictStoreLoad[0] = c.mayRestrictStoreLoad[0] || ruleMayRestrictEngine(r, "", core.EngineTiKV)
c.mayRestrictStoreLoad[1] = c.mayRestrictStoreLoad[1] || ruleMayRestrictEngine(r, core.EngineTiFlash)
}
g := c.groups[r.GroupID]
if g == nil {
// create default group configurations.
Expand All @@ -54,6 +62,39 @@ func (c *ruleConfig) adjust() {
}
}

func ruleMayRestrictEngine(rule *Rule, values ...string) bool {
hasOtherConstraint := false
for i := range rule.LabelConstraints {
if rule.LabelConstraints[i].Key != core.EngineKey {
hasOtherConstraint = true
break
}
}

matched := 0
for _, value := range values {
hasEngineConstraint, matchesValue := false, true
for i := range rule.LabelConstraints {
constraint := &rule.LabelConstraints[i]
if constraint.Key != core.EngineKey {
continue
}
hasEngineConstraint = true
if !constraint.matchValue(value) {
matchesValue = false
break
}
}
if value != "" && !hasEngineConstraint {
matchesValue = false
}
if matchesValue {
matched++
}
}
return matched > 0 && (hasOtherConstraint || matched < len(values))
}

func (c *ruleConfig) getRule(key [2]string) *Rule {
return c.rules[key]
}
Expand Down
14 changes: 8 additions & 6 deletions pkg/schedule/placement/label_constraint.go
Original file line number Diff line number Diff line change
Expand Up @@ -51,17 +51,19 @@ type LabelConstraint struct {

// MatchStore checks if a store matches the constraint.
func (c *LabelConstraint) MatchStore(store *core.StoreInfo) bool {
return c.matchValue(store.GetLabelValue(c.Key))
}

func (c *LabelConstraint) matchValue(value string) bool {
switch c.Op {
case In:
label := store.GetLabelValue(c.Key)
return label != "" && slice.AnyOf(c.Values, func(i int) bool { return c.Values[i] == label })
return value != "" && slice.AnyOf(c.Values, func(i int) bool { return c.Values[i] == value })
case NotIn:
label := store.GetLabelValue(c.Key)
return label == "" || slice.NoneOf(c.Values, func(i int) bool { return c.Values[i] == label })
return value == "" || slice.NoneOf(c.Values, func(i int) bool { return c.Values[i] == value })
case Exists:
return store.GetLabelValue(c.Key) != ""
return value != ""
case NotExists:
return store.GetLabelValue(c.Key) == ""
return value == ""
}
return false
}
Expand Down
11 changes: 11 additions & 0 deletions pkg/schedule/placement/rule_manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -368,6 +368,17 @@ func (m *RuleManager) GetRulesCount() int {
return len(m.ruleConfig.rules)
}

// MayRestrictStoreLoad returns whether placement rules may restrict the load
// statistics population for the given store engine.
func (m *RuleManager) MayRestrictStoreLoad(isTiKV bool) bool {
m.RLock()
defer m.RUnlock()
if isTiKV {
return m.ruleConfig.mayRestrictStoreLoad[0]
}
return m.ruleConfig.mayRestrictStoreLoad[1]
}

// GetGroupsCount returns the number of rule groups.
func (m *RuleManager) GetGroupsCount() int {
m.RLock()
Expand Down
33 changes: 33 additions & 0 deletions pkg/schedule/placement/rule_manager_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,39 @@ func TestDefault2(t *testing.T) {
re.Equal([]string{"zone", "rack", "host"}, rules[1].LocationLabels)
}

func TestMayRestrictStoreLoad(t *testing.T) {
re := require.New(t)
_, manager := newTestManager(t, false)
re.False(manager.MayRestrictStoreLoad(true))
re.False(manager.MayRestrictStoreLoad(false))

setConstraints := func(constraints ...LabelConstraint) {
re.NoError(manager.SetRule(&Rule{
GroupID: "group", ID: "rule", Role: Voter, Count: 1,
LabelConstraints: constraints,
}))
}

setConstraints(LabelConstraint{Key: "zone", Op: In, Values: []string{"z1"}})
re.True(manager.MayRestrictStoreLoad(true))
re.False(manager.MayRestrictStoreLoad(false))

setConstraints(LabelConstraint{Key: core.EngineKey, Op: In, Values: []string{core.EngineTiFlash}})
re.False(manager.MayRestrictStoreLoad(true))
re.False(manager.MayRestrictStoreLoad(false))

setConstraints(
LabelConstraint{Key: core.EngineKey, Op: In, Values: []string{core.EngineTiFlash}},
LabelConstraint{Key: "zone", Op: In, Values: []string{"z1"}},
)
re.False(manager.MayRestrictStoreLoad(true))
re.True(manager.MayRestrictStoreLoad(false))

setConstraints()
re.False(manager.MayRestrictStoreLoad(true))
re.False(manager.MayRestrictStoreLoad(false))
}

func TestAdjustRule(t *testing.T) {
re := require.New(t)
_, manager := newTestManager(t, false)
Expand Down
Loading
Loading