mirror of
https://gitee.com/milvus-io/milvus.git
synced 2025-12-06 17:18:35 +08:00
Remove sending newSegmentMsg to rootcoord (#6301)
Signed-off-by: sunby <bingyi.sun@zilliz.com>
This commit is contained in:
parent
972c44d82c
commit
836a45ec26
@ -71,14 +71,13 @@ type Server struct {
|
||||
serverLoopWg sync.WaitGroup
|
||||
isServing ServerState
|
||||
|
||||
kvClient *etcdkv.EtcdKV
|
||||
meta *meta
|
||||
segmentInfoStream msgstream.MsgStream
|
||||
segmentManager Manager
|
||||
allocator allocator
|
||||
cluster *cluster
|
||||
rootCoordClient types.RootCoord
|
||||
ddChannelName string
|
||||
kvClient *etcdkv.EtcdKV
|
||||
meta *meta
|
||||
segmentManager Manager
|
||||
allocator allocator
|
||||
cluster *cluster
|
||||
rootCoordClient types.RootCoord
|
||||
ddChannelName string
|
||||
|
||||
flushCh chan UniqueID
|
||||
msFactory msgstream.Factory
|
||||
@ -155,10 +154,6 @@ func (s *Server) Start() error {
|
||||
return err
|
||||
}
|
||||
|
||||
if err = s.initSegmentInfoChannel(); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
s.allocator = newRootCoordAllocator(s.ctx, s.rootCoordClient)
|
||||
|
||||
s.startSegmentManager()
|
||||
@ -231,20 +226,7 @@ func (s *Server) loadDataNodes() []*datapb.DataNodeInfo {
|
||||
}
|
||||
|
||||
func (s *Server) startSegmentManager() {
|
||||
helper := createNewSegmentHelper(s.segmentInfoStream)
|
||||
s.segmentManager = newSegmentManager(s.meta, s.allocator, withAllocHelper(helper))
|
||||
}
|
||||
|
||||
func (s *Server) initSegmentInfoChannel() error {
|
||||
var err error
|
||||
s.segmentInfoStream, err = s.msFactory.NewMsgStream(s.ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
s.segmentInfoStream.AsProducer([]string{Params.SegmentInfoChannelName})
|
||||
log.Debug("DataCoord AsProducer: " + Params.SegmentInfoChannelName)
|
||||
s.segmentInfoStream.Start()
|
||||
return nil
|
||||
s.segmentManager = newSegmentManager(s.meta, s.allocator)
|
||||
}
|
||||
|
||||
func (s *Server) initMeta() error {
|
||||
@ -501,7 +483,6 @@ func (s *Server) Stop() error {
|
||||
log.Debug("DataCoord server shutdown")
|
||||
atomic.StoreInt64(&s.isServing, ServerStateStopped)
|
||||
s.cluster.releaseSessions()
|
||||
s.segmentInfoStream.Close()
|
||||
s.stopServerLoop()
|
||||
return nil
|
||||
}
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user