-
Notifications
You must be signed in to change notification settings - Fork 727
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
coordinator: use a healthy region count to start coordinator #7044
Changes from 2 commits
accb06b
46406e7
d161acb
6be6f17
1f16cda
68a1107
c67ed32
67b4081
7bae37c
382816d
f89619a
4eebfe3
e90fa56
cc2d832
664427a
8a0c520
c88260d
32f856f
be7296f
99dc596
0663294
183649a
45d2cc5
6281684
a9d6e3e
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
Original file line number | Diff line number | Diff line change |
---|---|---|
|
@@ -71,8 +71,14 @@ type RegionInfo struct { | |
queryStats *pdpb.QueryStats | ||
flowRoundDivisor uint64 | ||
// buckets is not thread unsafe, it should be accessed by the request `report buckets` with greater version. | ||
buckets unsafe.Pointer | ||
fromHeartbeat bool | ||
buckets unsafe.Pointer | ||
// source is used to indicate region's source, such as FromHeartbeat/FromSync/InDisk. | ||
source RegionSource | ||
} | ||
|
||
// GetRegionSource returns the region source. | ||
func (r *RegionInfo) GetRegionSource() RegionSource { | ||
return r.source | ||
} | ||
|
||
// NewRegionInfo creates RegionInfo with region's meta and leader peer. | ||
|
@@ -171,6 +177,7 @@ func RegionFromHeartbeat(heartbeat *pdpb.RegionHeartbeatRequest, opts ...RegionC | |
interval: heartbeat.GetInterval(), | ||
replicationStatus: heartbeat.GetReplicationStatus(), | ||
queryStats: heartbeat.GetQueryStats(), | ||
source: FromHeartbeat, | ||
} | ||
|
||
for _, opt := range opts { | ||
|
@@ -639,11 +646,6 @@ func (r *RegionInfo) IsFlashbackChanged(l *RegionInfo) bool { | |
return r.meta.FlashbackStartTs != l.meta.FlashbackStartTs || r.meta.IsInFlashback != l.meta.IsInFlashback | ||
} | ||
|
||
// IsFromHeartbeat returns whether the region info is from the region heartbeat. | ||
func (r *RegionInfo) IsFromHeartbeat() bool { | ||
return r.fromHeartbeat | ||
} | ||
|
||
func (r *RegionInfo) isInvolved(startKey, endKey []byte) bool { | ||
return bytes.Compare(r.GetStartKey(), startKey) >= 0 && (len(endKey) == 0 || (len(r.GetEndKey()) > 0 && bytes.Compare(r.GetEndKey(), endKey) <= 0)) | ||
} | ||
|
@@ -683,7 +685,7 @@ func GenerateRegionGuideFunc(enableLog bool) RegionGuideFunc { | |
} | ||
saveKV, saveCache, isNew = true, true, true | ||
} else { | ||
if !origin.IsFromHeartbeat() { | ||
if origin.source == FromSync || origin.source == InDisk { | ||
isNew = true | ||
} | ||
r := region.GetRegionEpoch() | ||
|
@@ -793,6 +795,18 @@ type RegionsInfo struct { | |
pendingPeers map[uint64]*regionTree // storeID -> sub regionTree | ||
} | ||
|
||
// RegionSource is the source of region. | ||
type RegionSource uint32 | ||
|
||
const ( | ||
// InDisk means region is stale. | ||
InDisk RegionSource = iota | ||
// FromSync means region is stale. | ||
FromSync | ||
// FromHeartbeat means region is fresh. | ||
HuSharp marked this conversation as resolved.
Show resolved
Hide resolved
|
||
FromHeartbeat | ||
) | ||
|
||
// NewRegionsInfo creates RegionsInfo with tree, regions, leaders and followers | ||
func NewRegionsInfo() *RegionsInfo { | ||
return &RegionsInfo{ | ||
|
@@ -840,6 +854,8 @@ func (r *RegionsInfo) CheckAndPutRegion(region *RegionInfo) []*RegionInfo { | |
origin, overlaps, rangeChanged := r.setRegionLocked(region, true, ols...) | ||
r.t.Unlock() | ||
r.UpdateSubTree(region, origin, overlaps, rangeChanged) | ||
// InDisk means region is stale. | ||
r.AtomicAddStaleRegionCnt() | ||
return overlaps | ||
} | ||
|
||
|
@@ -857,6 +873,21 @@ func (r *RegionsInfo) PreCheckPutRegion(region *RegionInfo) (*RegionInfo, []*reg | |
return origin, overlaps, err | ||
} | ||
|
||
// GetStaleRegionCnt returns the stale region count. | ||
func (r *RegionsInfo) GetStaleRegionCnt() int64 { | ||
r.t.RLock() | ||
defer r.t.RUnlock() | ||
if r.tree.length() == 0 { | ||
return 0 | ||
} | ||
return r.tree.staleRegionCnt | ||
} | ||
|
||
// AtomicAddStaleRegionCnt atomically adds the stale region count. | ||
func (r *RegionsInfo) AtomicAddStaleRegionCnt() { | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. maybe we can maintain it in the There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I'm not sure what the difference is between putting in There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. changed |
||
r.tree.AtomicAddStaleRegionCnt() | ||
} | ||
|
||
// AtomicCheckAndPutRegion checks if the region is valid to put, if valid then put. | ||
func (r *RegionsInfo) AtomicCheckAndPutRegion(region *RegionInfo) ([]*RegionInfo, error) { | ||
r.t.Lock() | ||
|
@@ -870,6 +901,10 @@ func (r *RegionsInfo) AtomicCheckAndPutRegion(region *RegionInfo) ([]*RegionInfo | |
r.t.Unlock() | ||
return nil, err | ||
} | ||
// If origin is stale, need to sub the stale region count. | ||
if origin != nil && origin.source != FromHeartbeat && region.source == FromHeartbeat { | ||
r.tree.AtomicSubStaleRegionCnt() | ||
} | ||
origin, overlaps, rangeChanged := r.setRegionLocked(region, true, ols...) | ||
r.t.Unlock() | ||
r.UpdateSubTree(region, origin, overlaps, rangeChanged) | ||
|
Original file line number | Diff line number | Diff line change |
---|---|---|
|
@@ -16,6 +16,7 @@ package core | |
import ( | ||
"bytes" | ||
"math/rand" | ||
"sync/atomic" | ||
|
||
"github.com/pingcap/kvproto/pkg/metapb" | ||
"github.com/pingcap/log" | ||
|
@@ -61,6 +62,8 @@ type regionTree struct { | |
totalSize int64 | ||
totalWriteBytesRate float64 | ||
totalWriteKeysRate float64 | ||
|
||
staleRegionCnt int64 | ||
} | ||
|
||
func newRegionTree() *regionTree { | ||
|
@@ -69,6 +72,7 @@ func newRegionTree() *regionTree { | |
totalSize: 0, | ||
totalWriteBytesRate: 0, | ||
totalWriteKeysRate: 0, | ||
staleRegionCnt: 0, | ||
} | ||
} | ||
|
||
|
@@ -342,3 +346,20 @@ func (t *regionTree) TotalWriteRate() (bytesRate, keysRate float64) { | |
} | ||
return t.totalWriteBytesRate, t.totalWriteKeysRate | ||
} | ||
|
||
func (t *regionTree) AtomicAddStaleRegionCnt() { | ||
if t.length() == 0 { | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Why do we need to check the length here? There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Now we do not need to check it after moving |
||
return | ||
} | ||
atomic.AddInt64(&t.staleRegionCnt, 1) | ||
} | ||
|
||
func (t *regionTree) AtomicSubStaleRegionCnt() { | ||
if t.length() == 0 { | ||
return | ||
} | ||
if t.staleRegionCnt == 0 { | ||
HuSharp marked this conversation as resolved.
Show resolved
Hide resolved
|
||
return | ||
} | ||
atomic.AddInt64(&t.staleRegionCnt, -1) | ||
} |
Original file line number | Diff line number | Diff line change |
---|---|---|
|
@@ -17,8 +17,10 @@ package schedule | |
import ( | ||
"time" | ||
|
||
"github.com/pingcap/log" | ||
"github.com/tikv/pd/pkg/core" | ||
"github.com/tikv/pd/pkg/utils/syncutil" | ||
"go.uber.org/zap" | ||
) | ||
|
||
type prepareChecker struct { | ||
|
@@ -47,6 +49,11 @@ func (checker *prepareChecker) check(c *core.BasicCluster) bool { | |
checker.prepared = true | ||
return true | ||
} | ||
if float64(c.GetStaleRegionCnt()) < float64(c.GetTotalRegionCount())*(1-collectFactor) { | ||
log.Info("stale region num is satisfied, skip prepare checker", zap.Int64("stale-region", c.GetStaleRegionCnt()), zap.Int("total-region", c.GetTotalRegionCount())) | ||
checker.prepared = true | ||
return true | ||
} | ||
// The number of active regions should be more than total region of all stores * collectFactor | ||
if float64(c.GetTotalRegionCount())*collectFactor > float64(checker.sum) { | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Do we still need it? There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. maybe it is better to remove them in another pr? This pr just focuses on accelerating coordinator. |
||
return false | ||
|
Original file line number | Diff line number | Diff line change |
---|---|---|
|
@@ -1002,7 +1002,7 @@ func TestRegionSizeChanged(t *testing.T) { | |
core.WithLeader(region.GetPeers()[2]), | ||
core.SetApproximateSize(curMaxMergeSize-1), | ||
core.SetApproximateKeys(curMaxMergeKeys-1), | ||
core.SetFromHeartbeat(true), | ||
core.SetSource(core.FromHeartbeat), | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Is there any test to cover the problem? There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. done |
||
) | ||
cluster.processRegionHeartbeat(region) | ||
regionID := region.GetID() | ||
|
@@ -1012,7 +1012,7 @@ func TestRegionSizeChanged(t *testing.T) { | |
core.WithLeader(region.GetPeers()[2]), | ||
core.SetApproximateSize(curMaxMergeSize+1), | ||
core.SetApproximateKeys(curMaxMergeKeys+1), | ||
core.SetFromHeartbeat(true), | ||
core.SetSource(core.FromHeartbeat), | ||
) | ||
cluster.processRegionHeartbeat(region) | ||
re.False(cluster.regionStats.IsRegionStatsType(regionID, statistics.UndersizedRegion)) | ||
|
@@ -2375,6 +2375,7 @@ func (c *testCluster) LoadRegion(regionID uint64, followerStoreIDs ...uint64) er | |
peer, _ := c.AllocPeer(id) | ||
region.Peers = append(region.Peers, peer) | ||
} | ||
c.core.AtomicAddStaleRegionCnt() | ||
return c.putRegion(core.NewRegionInfo(region, nil)) | ||
} | ||
|
||
|
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
How about FromLoad or FromStorage?
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
update to
FromStorage