mirror of
https://github.com/milvus-io/milvus.git
synced 2026-07-21 02:05:41 +00:00
fix: [cp 2.6] wait for delegator serviceability in load config compliance (#51298)
issue: #51289 pr: #51295 - QueryCoord serviceability: treat delegator-reported non-serviceable leader views as not ready after data readiness checks pass - load config compliance: report live replica serviceability failures before query-visible fallback reasons --------- Signed-off-by: chyezh <chyezh@outlook.com>
This commit is contained in:
@@ -120,16 +120,6 @@ func (s *mixCoordImpl) HandleReplicaLoadConfigCompliance(w http.ResponseWriter,
|
||||
}
|
||||
}
|
||||
|
||||
for _, replica := range internalReplicas {
|
||||
if !replica.IsQueryVisible() {
|
||||
reason := fmt.Sprintf("collection %d: replica %d (rg=%s) is not query visible",
|
||||
collectionID, replica.GetID(), replica.GetResourceGroup())
|
||||
logger.Info("collection has query-invisible replica", zap.String("reason", reason))
|
||||
s.writeComplianceResponse(w, LoadConfigComplianceStateNotReady, reason)
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
// Now that replica count and RG distribution match, verify every replica actually
|
||||
// has a serviceable shard leader for every channel. This live dist check avoids
|
||||
// the stale CollectionObserver-persisted LoadPercentage that can falsely report
|
||||
@@ -141,6 +131,16 @@ func (s *mixCoordImpl) HandleReplicaLoadConfigCompliance(w http.ResponseWriter,
|
||||
return
|
||||
}
|
||||
|
||||
for _, replica := range internalReplicas {
|
||||
if !replica.IsQueryVisible() {
|
||||
reason := fmt.Sprintf("collection %d: replica %d (rg=%s) is not query visible",
|
||||
collectionID, replica.GetID(), replica.GetResourceGroup())
|
||||
logger.Info("collection has query-invisible replica", zap.String("reason", reason))
|
||||
s.writeComplianceResponse(w, LoadConfigComplianceStateNotReady, reason)
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
// Check that physical resources have been released from querynodes no longer
|
||||
// part of any replica. During scale-down a decommissioned replica's querynode may
|
||||
// still hold segments/channels while release is in flight; compliance must wait for
|
||||
|
||||
@@ -236,6 +236,9 @@ func TestHandleReplicaLoadConfigCompliance(t *testing.T) {
|
||||
mocker3 := mockey.Mock((*querycoordv2.Server).GetInternalReplicasByCollection).Return(replicas).Build()
|
||||
defer mocker3.UnPatch()
|
||||
|
||||
mockerSvc := mockey.Mock((*querycoordv2.Server).CheckAllReplicasServiceable).Return(nil).Build()
|
||||
defer mockerSvc.UnPatch()
|
||||
|
||||
req := httptest.NewRequest(http.MethodGet, "/api/v1/replicas/compliance", nil)
|
||||
w := httptest.NewRecorder()
|
||||
|
||||
@@ -471,6 +474,53 @@ func TestHandleReplicaLoadConfigCompliance(t *testing.T) {
|
||||
assert.Contains(t, resp.Reason, "catching up")
|
||||
})
|
||||
|
||||
t.Run("delegator not serviceable reason takes precedence over query-invisible", func(t *testing.T) {
|
||||
paramtable.Get().Save(Params.QueryCoordCfg.ClusterLevelLoadReplicaNumber.Key, "1")
|
||||
paramtable.Get().Save(Params.QueryCoordCfg.ClusterLevelLoadResourceGroups.Key, "rg1")
|
||||
defer paramtable.Get().Reset(Params.QueryCoordCfg.ClusterLevelLoadReplicaNumber.Key)
|
||||
defer paramtable.Get().Reset(Params.QueryCoordCfg.ClusterLevelLoadResourceGroups.Key)
|
||||
defer registerTestBalancer(t, nil)()
|
||||
|
||||
replica := meta.NewReplica(&querypb.Replica{
|
||||
ID: 1,
|
||||
CollectionID: 100,
|
||||
ResourceGroup: "rg1",
|
||||
}, typeutil.NewUniqueSet())
|
||||
mutableReplica := replica.CopyForWrite()
|
||||
mutableReplica.SetQueryInvisible(true)
|
||||
replicas := []*meta.Replica{mutableReplica.IntoReplica()}
|
||||
|
||||
coord := &mixCoordImpl{
|
||||
queryCoordServer: &querycoordv2.Server{},
|
||||
}
|
||||
|
||||
mocker1 := mockey.Mock((*mixCoordImpl).ShowLoadCollections).Return(&querypb.ShowCollectionsResponse{
|
||||
Status: &commonpb.Status{ErrorCode: commonpb.ErrorCode_Success},
|
||||
CollectionIDs: []int64{100},
|
||||
InMemoryPercentages: []int64{100},
|
||||
}, nil).Build()
|
||||
defer mocker1.UnPatch()
|
||||
|
||||
mocker2 := mockey.Mock((*querycoordv2.Server).GetInternalReplicasByCollection).Return(replicas).Build()
|
||||
defer mocker2.UnPatch()
|
||||
|
||||
mocker3 := mockey.Mock((*querycoordv2.Server).CheckAllReplicasServiceable).
|
||||
Return(fmt.Errorf("replica 1 (rg=rg1) channel c1 not serviceable: delegator reported not serviceable")).Build()
|
||||
defer mocker3.UnPatch()
|
||||
|
||||
req := httptest.NewRequest(http.MethodGet, "/api/v1/replicas/compliance", nil)
|
||||
w := httptest.NewRecorder()
|
||||
|
||||
coord.HandleReplicaLoadConfigCompliance(w, req)
|
||||
|
||||
assert.Equal(t, http.StatusOK, w.Code)
|
||||
var resp LoadConfigComplianceResponse
|
||||
assert.NoError(t, json.Unmarshal(w.Body.Bytes(), &resp))
|
||||
assert.Equal(t, LoadConfigComplianceStateNotReady, resp.State)
|
||||
assert.Contains(t, resp.Reason, "delegator reported not serviceable")
|
||||
assert.NotContains(t, resp.Reason, "not query visible")
|
||||
})
|
||||
|
||||
t.Run("query-invisible replica returns NotReady", func(t *testing.T) {
|
||||
paramtable.Get().Save(Params.QueryCoordCfg.ClusterLevelLoadReplicaNumber.Key, "1")
|
||||
paramtable.Get().Save(Params.QueryCoordCfg.ClusterLevelLoadResourceGroups.Key, "rg1")
|
||||
@@ -501,6 +551,9 @@ func TestHandleReplicaLoadConfigCompliance(t *testing.T) {
|
||||
mocker2 := mockey.Mock((*querycoordv2.Server).GetInternalReplicasByCollection).Return(replicas).Build()
|
||||
defer mocker2.UnPatch()
|
||||
|
||||
mocker3 := mockey.Mock((*querycoordv2.Server).CheckAllReplicasServiceable).Return(nil).Build()
|
||||
defer mocker3.UnPatch()
|
||||
|
||||
req := httptest.NewRequest(http.MethodGet, "/api/v1/replicas/compliance", nil)
|
||||
w := httptest.NewRecorder()
|
||||
|
||||
|
||||
@@ -35,6 +35,10 @@ func (s *Server) tryPromoteReadyLoadConfigReplicas(ctx context.Context) {
|
||||
if len(replicas) == 0 {
|
||||
return
|
||||
}
|
||||
// Promote load-config replicas all-or-nothing. Partially exposing ready
|
||||
// replicas can let the query path enter a resource group while the load path
|
||||
// is still moving the remaining replicas, causing the two paths to compete
|
||||
// for resources during the switch.
|
||||
for _, replica := range replicas {
|
||||
if err := s.checkReplicaServiceable(ctx, replica); err != nil {
|
||||
return
|
||||
|
||||
@@ -906,7 +906,11 @@ func (s *Server) checkReplicaServiceable(ctx context.Context, replica *meta.Repl
|
||||
replica.GetID(), replica.GetResourceGroup(), channelName)
|
||||
}
|
||||
if err := utils.CheckDelegatorDataReady(s.nodeMgr, s.targetMgr, leader.View, meta.CurrentTarget); err != nil {
|
||||
return merr.Wrapf(err, "replica %d (rg=%s) channel %s not serviceable", replica.GetID(), replica.GetResourceGroup(), channelName)
|
||||
return merr.Wrapf(err, "replica %d (rg=%s) not serviceable", replica.GetID(), replica.GetResourceGroup())
|
||||
}
|
||||
if !leader.IsServiceable() {
|
||||
err := merr.WrapErrChannelNotAvailable(channelName, "delegator reported not serviceable")
|
||||
return merr.Wrapf(err, "replica %d (rg=%s) not serviceable", replica.GetID(), replica.GetResourceGroup())
|
||||
}
|
||||
}
|
||||
return nil
|
||||
|
||||
@@ -1258,6 +1258,59 @@ func TestCheckAllReplicasServiceable(t *testing.T) {
|
||||
assert.ErrorContains(t, err, "not serviceable")
|
||||
})
|
||||
|
||||
t.Run("leader reported non-serviceable returns error even when data ready", func(t *testing.T) {
|
||||
s := newServer()
|
||||
mocker := mockey.Mock((*meta.ReplicaManager).GetByCollection).Return([]*meta.Replica{replica}).Build()
|
||||
defer mocker.UnPatch()
|
||||
s.nodeMgr.Add(session.NewNodeInfo(session.ImmutableNodeInfo{NodeID: 10, Address: "localhost:10", Hostname: "localhost"}))
|
||||
s.targetMgr.(*meta.MockTargetManager).EXPECT().GetDmChannelsByCollection(mock.Anything, collectionID, meta.CurrentTarget).Return(map[string]*meta.DmChannel{
|
||||
channelName: {VchannelInfo: &datapb.VchannelInfo{CollectionID: collectionID, ChannelName: channelName}},
|
||||
})
|
||||
s.targetMgr.(*meta.MockTargetManager).EXPECT().GetSealedSegmentsByChannel(mock.Anything, collectionID, channelName, meta.CurrentTarget).Return(map[int64]*datapb.SegmentInfo{
|
||||
42: {ID: 42, CollectionID: collectionID},
|
||||
})
|
||||
|
||||
s.dist.ChannelDistManager.Update(10, &meta.DmChannel{
|
||||
VchannelInfo: &datapb.VchannelInfo{CollectionID: collectionID, ChannelName: channelName},
|
||||
Node: 10,
|
||||
Version: 1,
|
||||
View: &meta.LeaderView{
|
||||
ID: 10, CollectionID: collectionID, Channel: channelName,
|
||||
Status: &querypb.LeaderViewStatus{Serviceable: false},
|
||||
Segments: map[int64]*querypb.SegmentDist{42: {NodeID: 10}},
|
||||
},
|
||||
})
|
||||
|
||||
err := s.CheckAllReplicasServiceable(context.Background(), collectionID)
|
||||
assert.ErrorContains(t, err, "reported not serviceable")
|
||||
})
|
||||
|
||||
t.Run("nil leader status returns error even when data ready", func(t *testing.T) {
|
||||
s := newServer()
|
||||
mocker := mockey.Mock((*meta.ReplicaManager).GetByCollection).Return([]*meta.Replica{replica}).Build()
|
||||
defer mocker.UnPatch()
|
||||
s.nodeMgr.Add(session.NewNodeInfo(session.ImmutableNodeInfo{NodeID: 10, Address: "localhost:10", Hostname: "localhost"}))
|
||||
s.targetMgr.(*meta.MockTargetManager).EXPECT().GetDmChannelsByCollection(mock.Anything, collectionID, meta.CurrentTarget).Return(map[string]*meta.DmChannel{
|
||||
channelName: {VchannelInfo: &datapb.VchannelInfo{CollectionID: collectionID, ChannelName: channelName}},
|
||||
})
|
||||
s.targetMgr.(*meta.MockTargetManager).EXPECT().GetSealedSegmentsByChannel(mock.Anything, collectionID, channelName, meta.CurrentTarget).Return(map[int64]*datapb.SegmentInfo{
|
||||
42: {ID: 42, CollectionID: collectionID},
|
||||
})
|
||||
|
||||
s.dist.ChannelDistManager.Update(10, &meta.DmChannel{
|
||||
VchannelInfo: &datapb.VchannelInfo{CollectionID: collectionID, ChannelName: channelName},
|
||||
Node: 10,
|
||||
Version: 1,
|
||||
View: &meta.LeaderView{
|
||||
ID: 10, CollectionID: collectionID, Channel: channelName,
|
||||
Segments: map[int64]*querypb.SegmentDist{42: {NodeID: 10}},
|
||||
},
|
||||
})
|
||||
|
||||
err := s.CheckAllReplicasServiceable(context.Background(), collectionID)
|
||||
assert.ErrorContains(t, err, "reported not serviceable")
|
||||
})
|
||||
|
||||
t.Run("all serviceable returns nil", func(t *testing.T) {
|
||||
s := newServer()
|
||||
mocker := mockey.Mock((*meta.ReplicaManager).GetByCollection).Return([]*meta.Replica{replica}).Build()
|
||||
|
||||
Reference in New Issue
Block a user