2021-01-19 14:44:03 +08:00
|
|
|
package masterservice
|
|
|
|
|
|
|
|
import (
|
|
|
|
"context"
|
2021-03-06 17:47:11 +08:00
|
|
|
"errors"
|
2021-01-25 18:33:10 +08:00
|
|
|
"fmt"
|
2021-01-19 14:44:03 +08:00
|
|
|
"math/rand"
|
|
|
|
"sync"
|
|
|
|
"sync/atomic"
|
|
|
|
"time"
|
|
|
|
|
2021-02-27 10:11:52 +08:00
|
|
|
"go.etcd.io/etcd/clientv3"
|
|
|
|
"go.uber.org/zap"
|
|
|
|
|
2021-02-24 17:12:06 +08:00
|
|
|
"github.com/zilliztech/milvus-distributed/internal/allocator"
|
2021-01-19 14:44:03 +08:00
|
|
|
etcdkv "github.com/zilliztech/milvus-distributed/internal/kv/etcd"
|
2021-02-24 16:25:40 +08:00
|
|
|
"github.com/zilliztech/milvus-distributed/internal/log"
|
2021-01-20 09:36:50 +08:00
|
|
|
ms "github.com/zilliztech/milvus-distributed/internal/msgstream"
|
2021-03-06 17:47:11 +08:00
|
|
|
"github.com/zilliztech/milvus-distributed/internal/tso"
|
|
|
|
"github.com/zilliztech/milvus-distributed/internal/types"
|
|
|
|
"github.com/zilliztech/milvus-distributed/internal/util/retry"
|
|
|
|
"github.com/zilliztech/milvus-distributed/internal/util/tsoutil"
|
|
|
|
"github.com/zilliztech/milvus-distributed/internal/util/typeutil"
|
|
|
|
|
2021-01-19 14:44:03 +08:00
|
|
|
"github.com/zilliztech/milvus-distributed/internal/proto/commonpb"
|
2021-01-21 10:01:29 +08:00
|
|
|
"github.com/zilliztech/milvus-distributed/internal/proto/datapb"
|
2021-01-24 20:26:35 +08:00
|
|
|
"github.com/zilliztech/milvus-distributed/internal/proto/indexpb"
|
2021-03-12 14:22:09 +08:00
|
|
|
"github.com/zilliztech/milvus-distributed/internal/proto/internalpb"
|
2021-01-19 14:44:03 +08:00
|
|
|
"github.com/zilliztech/milvus-distributed/internal/proto/masterpb"
|
|
|
|
"github.com/zilliztech/milvus-distributed/internal/proto/milvuspb"
|
2021-01-24 20:26:35 +08:00
|
|
|
"github.com/zilliztech/milvus-distributed/internal/proto/proxypb"
|
2021-02-05 14:09:55 +08:00
|
|
|
"github.com/zilliztech/milvus-distributed/internal/proto/querypb"
|
2021-01-19 14:44:03 +08:00
|
|
|
)
|
|
|
|
|
2021-03-12 14:22:09 +08:00
|
|
|
// internalpb -> internalpb
|
2021-01-19 14:44:03 +08:00
|
|
|
// proxypb(proxy_service)
|
|
|
|
// querypb(query_service)
|
|
|
|
// datapb(data_service)
|
|
|
|
// indexpb(index_service)
|
2021-01-22 09:36:18 +08:00
|
|
|
// milvuspb -> milvuspb
|
2021-01-20 09:36:50 +08:00
|
|
|
// masterpb2 -> masterpb (master_service)
|
2021-01-19 14:44:03 +08:00
|
|
|
|
|
|
|
// ------------------ struct -----------------------
|
|
|
|
|
|
|
|
// master core
|
|
|
|
type Core struct {
|
2021-01-23 17:56:57 +08:00
|
|
|
/*
|
|
|
|
ProxyServiceClient Interface:
|
|
|
|
get proxy service time tick channel,InvalidateCollectionMetaCache
|
2021-01-19 14:44:03 +08:00
|
|
|
|
2021-01-23 17:56:57 +08:00
|
|
|
DataService Interface:
|
|
|
|
Segment States Channel, from DataService, if create new segment, data service should put the segment id into this channel, and let the master add the segment id to the collection meta
|
|
|
|
Segment Flush Watcher, monitor if segment has flushed into disk
|
2021-01-19 14:44:03 +08:00
|
|
|
|
2021-03-04 10:35:28 +08:00
|
|
|
IndexService Interface
|
2021-02-23 18:08:17 +08:00
|
|
|
IndexService Sch, tell index service to build index
|
2021-01-23 17:56:57 +08:00
|
|
|
*/
|
2021-01-19 14:44:03 +08:00
|
|
|
|
|
|
|
MetaTable *metaTable
|
|
|
|
//id allocator
|
2021-02-24 17:12:06 +08:00
|
|
|
idAllocator *allocator.GlobalIDAllocator
|
2021-01-19 14:44:03 +08:00
|
|
|
//tso allocator
|
2021-02-24 17:12:06 +08:00
|
|
|
tsoAllocator *tso.GlobalTSOAllocator
|
2021-01-19 14:44:03 +08:00
|
|
|
|
|
|
|
//inner members
|
|
|
|
ctx context.Context
|
|
|
|
cancel context.CancelFunc
|
|
|
|
etcdCli *clientv3.Client
|
|
|
|
kvBase *etcdkv.EtcdKV
|
|
|
|
metaKV *etcdkv.EtcdKV
|
|
|
|
|
2021-01-20 09:36:50 +08:00
|
|
|
//setMsgStreams, receive time tick from proxy service time tick channel
|
|
|
|
ProxyTimeTickChan chan typeutil.Timestamp
|
2021-01-19 14:44:03 +08:00
|
|
|
|
2021-01-20 09:36:50 +08:00
|
|
|
//setMsgStreams, send time tick into dd channel and time tick channel
|
2021-01-19 14:44:03 +08:00
|
|
|
SendTimeTick func(t typeutil.Timestamp) error
|
|
|
|
|
2021-01-20 09:36:50 +08:00
|
|
|
//setMsgStreams, send create collection into dd channel
|
2021-03-12 14:22:09 +08:00
|
|
|
DdCreateCollectionReq func(req *internalpb.CreateCollectionRequest) error
|
2021-01-19 14:44:03 +08:00
|
|
|
|
2021-01-20 09:36:50 +08:00
|
|
|
//setMsgStreams, send drop collection into dd channel, and notify the proxy to delete this collection
|
2021-03-12 14:22:09 +08:00
|
|
|
DdDropCollectionReq func(req *internalpb.DropCollectionRequest) error
|
2021-01-19 14:44:03 +08:00
|
|
|
|
2021-01-20 09:36:50 +08:00
|
|
|
//setMsgStreams, send create partition into dd channel
|
2021-03-12 14:22:09 +08:00
|
|
|
DdCreatePartitionReq func(req *internalpb.CreatePartitionRequest) error
|
2021-01-19 14:44:03 +08:00
|
|
|
|
2021-01-20 09:36:50 +08:00
|
|
|
//setMsgStreams, send drop partition into dd channel
|
2021-03-12 14:22:09 +08:00
|
|
|
DdDropPartitionReq func(req *internalpb.DropPartitionRequest) error
|
2021-01-19 14:44:03 +08:00
|
|
|
|
2021-01-21 10:01:29 +08:00
|
|
|
//setMsgStreams segment channel, receive segment info from data service, if master create segment
|
|
|
|
DataServiceSegmentChan chan *datapb.SegmentInfo
|
|
|
|
|
2021-01-22 15:41:54 +08:00
|
|
|
//setMsgStreams ,if segment flush completed, data node would put segment id into msg stream
|
|
|
|
DataNodeSegmentFlushCompletedChan chan typeutil.UniqueID
|
|
|
|
|
2021-02-05 14:09:55 +08:00
|
|
|
//get binlog file path from data service,
|
2021-01-21 10:01:29 +08:00
|
|
|
GetBinlogFilePathsFromDataServiceReq func(segID typeutil.UniqueID, fieldID typeutil.UniqueID) ([]string, error)
|
2021-03-08 15:46:51 +08:00
|
|
|
GetNumRowsReq func(segID typeutil.UniqueID) (int64, error)
|
2021-01-21 10:01:29 +08:00
|
|
|
|
2021-02-05 14:09:55 +08:00
|
|
|
//call index builder's client to build index, return build id
|
2021-02-03 11:52:19 +08:00
|
|
|
BuildIndexReq func(binlog []string, typeParams []*commonpb.KeyValuePair, indexParams []*commonpb.KeyValuePair, indexID typeutil.UniqueID, indexName string) (typeutil.UniqueID, error)
|
2021-02-20 15:38:44 +08:00
|
|
|
DropIndexReq func(indexID typeutil.UniqueID) error
|
2021-01-21 10:01:29 +08:00
|
|
|
|
2021-02-05 14:09:55 +08:00
|
|
|
//proxy service interface, notify proxy service to drop collection
|
2021-01-24 20:26:35 +08:00
|
|
|
InvalidateCollectionMetaCache func(ts typeutil.Timestamp, dbName string, collectionName string) error
|
2021-01-23 10:12:41 +08:00
|
|
|
|
2021-02-05 14:09:55 +08:00
|
|
|
//query service interface, notify query service to release collection
|
|
|
|
ReleaseCollection func(ts typeutil.Timestamp, dbID typeutil.UniqueID, collectionID typeutil.UniqueID) error
|
|
|
|
|
2021-01-21 10:01:29 +08:00
|
|
|
// put create index task into this chan
|
|
|
|
indexTaskQueue chan *CreateIndexTask
|
|
|
|
|
2021-01-19 14:44:03 +08:00
|
|
|
//dd request scheduler
|
|
|
|
ddReqQueue chan reqTask //dd request will be push into this chan
|
|
|
|
lastDdTimeStamp typeutil.Timestamp
|
|
|
|
|
|
|
|
//time tick loop
|
|
|
|
lastTimeTick typeutil.Timestamp
|
|
|
|
|
|
|
|
//states code
|
|
|
|
stateCode atomic.Value
|
|
|
|
|
|
|
|
//call once
|
|
|
|
initOnce sync.Once
|
|
|
|
startOnce sync.Once
|
2021-02-23 11:40:30 +08:00
|
|
|
//isInit atomic.Value
|
2021-02-08 14:30:54 +08:00
|
|
|
|
|
|
|
msFactory ms.Factory
|
2021-01-19 14:44:03 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
// --------------------- function --------------------------
|
|
|
|
|
2021-02-08 14:30:54 +08:00
|
|
|
func NewCore(c context.Context, factory ms.Factory) (*Core, error) {
|
2021-01-19 14:44:03 +08:00
|
|
|
ctx, cancel := context.WithCancel(c)
|
|
|
|
rand.Seed(time.Now().UnixNano())
|
|
|
|
core := &Core{
|
2021-02-08 14:30:54 +08:00
|
|
|
ctx: ctx,
|
|
|
|
cancel: cancel,
|
|
|
|
msFactory: factory,
|
2021-01-19 14:44:03 +08:00
|
|
|
}
|
2021-03-12 14:22:09 +08:00
|
|
|
core.UpdateStateCode(internalpb.StateCode_Abnormal)
|
2021-01-19 14:44:03 +08:00
|
|
|
return core, nil
|
|
|
|
}
|
|
|
|
|
2021-03-12 14:22:09 +08:00
|
|
|
func (c *Core) UpdateStateCode(code internalpb.StateCode) {
|
2021-02-23 11:40:30 +08:00
|
|
|
c.stateCode.Store(code)
|
|
|
|
}
|
|
|
|
|
2021-01-19 14:44:03 +08:00
|
|
|
func (c *Core) checkInit() error {
|
|
|
|
if c.MetaTable == nil {
|
2021-03-05 10:15:27 +08:00
|
|
|
return errors.New("MetaTable is nil")
|
2021-01-19 14:44:03 +08:00
|
|
|
}
|
|
|
|
if c.idAllocator == nil {
|
2021-03-05 10:15:27 +08:00
|
|
|
return errors.New("idAllocator is nil")
|
2021-01-19 14:44:03 +08:00
|
|
|
}
|
|
|
|
if c.tsoAllocator == nil {
|
2021-03-05 10:15:27 +08:00
|
|
|
return errors.New("tsoAllocator is nil")
|
2021-01-19 14:44:03 +08:00
|
|
|
}
|
|
|
|
if c.etcdCli == nil {
|
2021-03-05 10:15:27 +08:00
|
|
|
return errors.New("etcdCli is nil")
|
2021-01-19 14:44:03 +08:00
|
|
|
}
|
|
|
|
if c.metaKV == nil {
|
2021-03-05 10:15:27 +08:00
|
|
|
return errors.New("metaKV is nil")
|
2021-01-19 14:44:03 +08:00
|
|
|
}
|
|
|
|
if c.kvBase == nil {
|
2021-03-05 10:15:27 +08:00
|
|
|
return errors.New("kvBase is nil")
|
2021-01-19 14:44:03 +08:00
|
|
|
}
|
|
|
|
if c.ProxyTimeTickChan == nil {
|
2021-03-05 10:15:27 +08:00
|
|
|
return errors.New("ProxyTimeTickChan is nil")
|
2021-01-19 14:44:03 +08:00
|
|
|
}
|
|
|
|
if c.ddReqQueue == nil {
|
2021-03-05 10:15:27 +08:00
|
|
|
return errors.New("ddReqQueue is nil")
|
2021-01-19 14:44:03 +08:00
|
|
|
}
|
|
|
|
if c.DdCreateCollectionReq == nil {
|
2021-03-05 10:15:27 +08:00
|
|
|
return errors.New("DdCreateCollectionReq is nil")
|
2021-01-19 14:44:03 +08:00
|
|
|
}
|
|
|
|
if c.DdDropCollectionReq == nil {
|
2021-03-05 10:15:27 +08:00
|
|
|
return errors.New("DdDropCollectionReq is nil")
|
2021-01-19 14:44:03 +08:00
|
|
|
}
|
|
|
|
if c.DdCreatePartitionReq == nil {
|
2021-03-05 10:15:27 +08:00
|
|
|
return errors.New("DdCreatePartitionReq is nil")
|
2021-01-19 14:44:03 +08:00
|
|
|
}
|
|
|
|
if c.DdDropPartitionReq == nil {
|
2021-03-05 10:15:27 +08:00
|
|
|
return errors.New("DdDropPartitionReq is nil")
|
2021-01-19 14:44:03 +08:00
|
|
|
}
|
2021-01-21 10:01:29 +08:00
|
|
|
if c.DataServiceSegmentChan == nil {
|
2021-03-05 10:15:27 +08:00
|
|
|
return errors.New("DataServiceSegmentChan is nil")
|
2021-01-21 10:01:29 +08:00
|
|
|
}
|
|
|
|
if c.GetBinlogFilePathsFromDataServiceReq == nil {
|
2021-03-05 10:15:27 +08:00
|
|
|
return errors.New("GetBinlogFilePathsFromDataServiceReq is nil")
|
2021-01-21 10:01:29 +08:00
|
|
|
}
|
2021-03-08 15:46:51 +08:00
|
|
|
if c.GetNumRowsReq == nil {
|
|
|
|
return errors.New("GetNumRowsReq is nil")
|
|
|
|
}
|
2021-01-21 10:01:29 +08:00
|
|
|
if c.BuildIndexReq == nil {
|
2021-03-05 10:15:27 +08:00
|
|
|
return errors.New("BuildIndexReq is nil")
|
2021-01-21 10:01:29 +08:00
|
|
|
}
|
2021-02-20 15:38:44 +08:00
|
|
|
if c.DropIndexReq == nil {
|
2021-03-05 10:15:27 +08:00
|
|
|
return errors.New("DropIndexReq is nil")
|
2021-02-20 15:38:44 +08:00
|
|
|
}
|
2021-01-23 10:12:41 +08:00
|
|
|
if c.InvalidateCollectionMetaCache == nil {
|
2021-03-05 10:15:27 +08:00
|
|
|
return errors.New("InvalidateCollectionMetaCache is nil")
|
2021-01-23 10:12:41 +08:00
|
|
|
}
|
2021-01-21 10:01:29 +08:00
|
|
|
if c.indexTaskQueue == nil {
|
2021-03-05 10:15:27 +08:00
|
|
|
return errors.New("indexTaskQueue is nil")
|
2021-01-21 10:01:29 +08:00
|
|
|
}
|
2021-01-22 15:41:54 +08:00
|
|
|
if c.DataNodeSegmentFlushCompletedChan == nil {
|
2021-03-05 10:15:27 +08:00
|
|
|
return errors.New("DataNodeSegmentFlushCompletedChan is nil")
|
2021-01-22 15:41:54 +08:00
|
|
|
}
|
2021-02-05 14:09:55 +08:00
|
|
|
if c.ReleaseCollection == nil {
|
2021-03-05 10:15:27 +08:00
|
|
|
return errors.New("ReleaseCollection is nil")
|
2021-02-05 14:09:55 +08:00
|
|
|
}
|
|
|
|
|
2021-02-27 10:11:52 +08:00
|
|
|
log.Debug("master", zap.Int64("node id", int64(Params.NodeID)))
|
|
|
|
log.Debug("master", zap.String("dd channel name", Params.DdChannel))
|
|
|
|
log.Debug("master", zap.String("time tick channel name", Params.TimeTickChannel))
|
2021-01-19 14:44:03 +08:00
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
|
|
|
func (c *Core) startDdScheduler() {
|
|
|
|
for {
|
|
|
|
select {
|
|
|
|
case <-c.ctx.Done():
|
2021-02-27 10:11:52 +08:00
|
|
|
log.Debug("close dd scheduler, exit task execution loop")
|
2021-01-19 14:44:03 +08:00
|
|
|
return
|
|
|
|
case task, ok := <-c.ddReqQueue:
|
|
|
|
if !ok {
|
2021-02-27 10:11:52 +08:00
|
|
|
log.Debug("dd chan is closed, exit task execution loop")
|
2021-01-19 14:44:03 +08:00
|
|
|
return
|
|
|
|
}
|
|
|
|
ts, err := task.Ts()
|
|
|
|
if err != nil {
|
|
|
|
task.Notify(err)
|
|
|
|
break
|
|
|
|
}
|
2021-01-23 10:12:41 +08:00
|
|
|
if !task.IgnoreTimeStamp() && ts <= c.lastDdTimeStamp {
|
2021-03-05 10:15:27 +08:00
|
|
|
task.Notify(fmt.Errorf("input timestamp = %d, last dd time stamp = %d", ts, c.lastDdTimeStamp))
|
2021-01-19 14:44:03 +08:00
|
|
|
break
|
|
|
|
}
|
|
|
|
err = task.Execute()
|
|
|
|
task.Notify(err)
|
2021-01-23 10:12:41 +08:00
|
|
|
if ts > c.lastDdTimeStamp {
|
|
|
|
c.lastDdTimeStamp = ts
|
|
|
|
}
|
2021-01-19 14:44:03 +08:00
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
func (c *Core) startTimeTickLoop() {
|
|
|
|
for {
|
|
|
|
select {
|
|
|
|
case <-c.ctx.Done():
|
2021-02-27 10:11:52 +08:00
|
|
|
log.Debug("close master time tick loop")
|
2021-01-19 14:44:03 +08:00
|
|
|
return
|
|
|
|
case tt, ok := <-c.ProxyTimeTickChan:
|
|
|
|
if !ok {
|
2021-02-27 10:11:52 +08:00
|
|
|
log.Warn("proxyTimeTickStream is closed, exit time tick loop")
|
2021-01-19 14:44:03 +08:00
|
|
|
return
|
|
|
|
}
|
|
|
|
if tt <= c.lastTimeTick {
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Warn("master time tick go back", zap.Uint64("last time tick", c.lastTimeTick), zap.Uint64("input time tick ", tt))
|
2021-01-19 14:44:03 +08:00
|
|
|
}
|
|
|
|
if err := c.SendTimeTick(tt); err != nil {
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Warn("master send time tick into dd and time_tick channel failed", zap.String("error", err.Error()))
|
2021-01-19 14:44:03 +08:00
|
|
|
}
|
|
|
|
c.lastTimeTick = tt
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2021-01-21 10:01:29 +08:00
|
|
|
//data service send segment info to master when create segment
|
|
|
|
func (c *Core) startDataServiceSegmentLoop() {
|
|
|
|
for {
|
|
|
|
select {
|
|
|
|
case <-c.ctx.Done():
|
2021-02-27 10:11:52 +08:00
|
|
|
log.Debug("close data service segment loop")
|
2021-01-21 10:01:29 +08:00
|
|
|
return
|
|
|
|
case seg, ok := <-c.DataServiceSegmentChan:
|
|
|
|
if !ok {
|
2021-02-27 10:11:52 +08:00
|
|
|
log.Debug("data service segment is closed, exit loop")
|
2021-01-21 10:01:29 +08:00
|
|
|
return
|
|
|
|
}
|
|
|
|
if seg == nil {
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Warn("segment from data service is nil")
|
2021-01-21 10:01:29 +08:00
|
|
|
} else if err := c.MetaTable.AddSegment(seg); err != nil {
|
|
|
|
//what if master add segment failed, but data service success?
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Warn("add segment info meta table failed ", zap.String("error", err.Error()))
|
2021-01-26 19:24:09 +08:00
|
|
|
} else {
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("add segment", zap.Int64("collection id", seg.CollectionID), zap.Int64("partition id", seg.PartitionID), zap.Int64("segment id", seg.SegmentID))
|
2021-01-21 10:01:29 +08:00
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
//create index loop
|
|
|
|
func (c *Core) startCreateIndexLoop() {
|
|
|
|
for {
|
|
|
|
select {
|
|
|
|
case <-c.ctx.Done():
|
2021-02-27 10:11:52 +08:00
|
|
|
log.Debug("close create index loop")
|
2021-01-21 10:01:29 +08:00
|
|
|
return
|
|
|
|
case t, ok := <-c.indexTaskQueue:
|
|
|
|
if !ok {
|
2021-02-27 10:11:52 +08:00
|
|
|
log.Debug("index task chan has closed, exit loop")
|
2021-01-21 10:01:29 +08:00
|
|
|
return
|
|
|
|
}
|
|
|
|
if err := t.BuildIndex(); err != nil {
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Warn("create index failed", zap.String("error", err.Error()))
|
2021-01-26 19:24:09 +08:00
|
|
|
} else {
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("create index", zap.String("index name", t.indexName), zap.String("field name", t.fieldSchema.Name), zap.Int64("segment id", t.segmentID))
|
2021-01-21 10:01:29 +08:00
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2021-01-22 15:41:54 +08:00
|
|
|
func (c *Core) startSegmentFlushCompletedLoop() {
|
|
|
|
for {
|
|
|
|
select {
|
|
|
|
case <-c.ctx.Done():
|
2021-02-27 10:11:52 +08:00
|
|
|
log.Debug("close segment flush completed loop")
|
2021-01-22 15:41:54 +08:00
|
|
|
return
|
|
|
|
case seg, ok := <-c.DataNodeSegmentFlushCompletedChan:
|
|
|
|
if !ok {
|
2021-02-27 10:11:52 +08:00
|
|
|
log.Debug("data node segment flush completed chan has colsed, exit loop")
|
2021-01-22 15:41:54 +08:00
|
|
|
}
|
2021-03-10 14:45:35 +08:00
|
|
|
log.Debug("flush segment", zap.Int64("id", seg))
|
2021-02-02 10:09:10 +08:00
|
|
|
coll, err := c.MetaTable.GetCollectionBySegmentID(seg)
|
2021-01-22 15:41:54 +08:00
|
|
|
if err != nil {
|
2021-03-10 14:45:35 +08:00
|
|
|
log.Warn("GetCollectionBySegmentID error", zap.Error(err))
|
2021-02-20 15:38:44 +08:00
|
|
|
break
|
2021-01-22 15:41:54 +08:00
|
|
|
}
|
2021-03-10 14:45:35 +08:00
|
|
|
err = c.MetaTable.AddFlushedSegment(seg)
|
|
|
|
if err != nil {
|
|
|
|
log.Warn("AddFlushedSegment error", zap.Error(err))
|
|
|
|
}
|
2021-02-09 13:11:55 +08:00
|
|
|
for _, f := range coll.FieldIndexes {
|
2021-02-11 08:41:59 +08:00
|
|
|
idxInfo, err := c.MetaTable.GetIndexByID(f.IndexID)
|
|
|
|
if err != nil {
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Warn("index not found", zap.Int64("index id", f.IndexID))
|
2021-02-20 15:38:44 +08:00
|
|
|
continue
|
2021-02-11 08:41:59 +08:00
|
|
|
}
|
|
|
|
|
2021-02-02 10:09:10 +08:00
|
|
|
fieldSch, err := GetFieldSchemaByID(coll, f.FiledID)
|
|
|
|
if err == nil {
|
|
|
|
t := &CreateIndexTask{
|
|
|
|
core: c,
|
|
|
|
segmentID: seg,
|
2021-02-11 08:41:59 +08:00
|
|
|
indexName: idxInfo.IndexName,
|
|
|
|
indexID: idxInfo.IndexID,
|
2021-02-02 10:09:10 +08:00
|
|
|
fieldSchema: fieldSch,
|
2021-03-10 14:45:35 +08:00
|
|
|
indexParams: idxInfo.IndexParams,
|
2021-02-02 10:09:10 +08:00
|
|
|
}
|
|
|
|
c.indexTaskQueue <- t
|
2021-01-22 15:41:54 +08:00
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2021-01-27 16:38:18 +08:00
|
|
|
func (c *Core) tsLoop() {
|
2021-02-24 17:12:06 +08:00
|
|
|
tsoTicker := time.NewTicker(tso.UpdateTimestampStep)
|
2021-01-27 16:38:18 +08:00
|
|
|
defer tsoTicker.Stop()
|
|
|
|
ctx, cancel := context.WithCancel(c.ctx)
|
|
|
|
defer cancel()
|
|
|
|
for {
|
|
|
|
select {
|
|
|
|
case <-tsoTicker.C:
|
|
|
|
if err := c.tsoAllocator.UpdateTSO(); err != nil {
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Warn("failed to update timestamp", zap.String("error", err.Error()))
|
2021-01-27 16:38:18 +08:00
|
|
|
return
|
|
|
|
}
|
|
|
|
if err := c.idAllocator.UpdateID(); err != nil {
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Warn("failed to update id", zap.String("error", err.Error()))
|
2021-01-27 16:38:18 +08:00
|
|
|
return
|
|
|
|
}
|
|
|
|
case <-ctx.Done():
|
|
|
|
// Server is closed and it should return nil.
|
2021-02-27 10:11:52 +08:00
|
|
|
log.Debug("tsLoop is closed")
|
2021-01-27 16:38:18 +08:00
|
|
|
return
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
2021-01-20 09:36:50 +08:00
|
|
|
func (c *Core) setMsgStreams() error {
|
2021-01-24 20:26:35 +08:00
|
|
|
if Params.PulsarAddress == "" {
|
2021-03-05 10:15:27 +08:00
|
|
|
return errors.New("PulsarAddress is empty")
|
2021-01-24 20:26:35 +08:00
|
|
|
}
|
|
|
|
if Params.MsgChannelSubName == "" {
|
2021-03-05 10:15:27 +08:00
|
|
|
return errors.New("MsgChannelSubName is emptyr")
|
2021-01-24 20:26:35 +08:00
|
|
|
}
|
|
|
|
|
2021-01-20 09:36:50 +08:00
|
|
|
//proxy time tick stream,
|
2021-01-24 20:26:35 +08:00
|
|
|
if Params.ProxyTimeTickChannel == "" {
|
2021-03-05 10:15:27 +08:00
|
|
|
return errors.New("ProxyTimeTickChannel is empty")
|
2021-01-24 20:26:35 +08:00
|
|
|
}
|
2021-02-04 14:37:12 +08:00
|
|
|
|
2021-02-08 14:30:54 +08:00
|
|
|
var err error
|
|
|
|
m := map[string]interface{}{
|
|
|
|
"PulsarAddress": Params.PulsarAddress,
|
|
|
|
"ReceiveBufSize": 1024,
|
|
|
|
"PulsarBufSize": 1024}
|
|
|
|
err = c.msFactory.SetParams(m)
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
|
|
|
proxyTimeTickStream, _ := c.msFactory.NewMsgStream(c.ctx)
|
2021-02-04 14:37:12 +08:00
|
|
|
proxyTimeTickStream.AsConsumer([]string{Params.ProxyTimeTickChannel}, Params.MsgChannelSubName)
|
2021-03-05 18:16:50 +08:00
|
|
|
log.Debug("master AsConsumer: " + Params.ProxyTimeTickChannel + " : " + Params.MsgChannelSubName)
|
2021-01-20 09:36:50 +08:00
|
|
|
proxyTimeTickStream.Start()
|
|
|
|
|
|
|
|
// master time tick channel
|
2021-01-24 20:26:35 +08:00
|
|
|
if Params.TimeTickChannel == "" {
|
2021-03-05 10:15:27 +08:00
|
|
|
return errors.New("TimeTickChannel is empty")
|
2021-01-24 20:26:35 +08:00
|
|
|
}
|
2021-02-08 14:30:54 +08:00
|
|
|
timeTickStream, _ := c.msFactory.NewMsgStream(c.ctx)
|
2021-02-04 14:37:12 +08:00
|
|
|
timeTickStream.AsProducer([]string{Params.TimeTickChannel})
|
2021-03-05 18:16:50 +08:00
|
|
|
log.Debug("masterservice AsProducer: " + Params.TimeTickChannel)
|
2021-01-20 09:36:50 +08:00
|
|
|
|
|
|
|
// master dd channel
|
2021-01-24 20:26:35 +08:00
|
|
|
if Params.DdChannel == "" {
|
2021-03-05 10:15:27 +08:00
|
|
|
return errors.New("DdChannel is empty")
|
2021-01-24 20:26:35 +08:00
|
|
|
}
|
2021-02-08 14:30:54 +08:00
|
|
|
ddStream, _ := c.msFactory.NewMsgStream(c.ctx)
|
2021-02-04 14:37:12 +08:00
|
|
|
ddStream.AsProducer([]string{Params.DdChannel})
|
2021-03-05 18:16:50 +08:00
|
|
|
log.Debug("masterservice AsProducer: " + Params.DdChannel)
|
2021-01-20 09:36:50 +08:00
|
|
|
|
|
|
|
c.SendTimeTick = func(t typeutil.Timestamp) error {
|
|
|
|
msgPack := ms.MsgPack{}
|
|
|
|
baseMsg := ms.BaseMsg{
|
|
|
|
BeginTimestamp: t,
|
|
|
|
EndTimestamp: t,
|
|
|
|
HashValues: []uint32{0},
|
|
|
|
}
|
2021-03-12 14:22:09 +08:00
|
|
|
timeTickResult := internalpb.TimeTickMsg{
|
2021-01-20 09:36:50 +08:00
|
|
|
Base: &commonpb.MsgBase{
|
2021-03-10 14:45:35 +08:00
|
|
|
MsgType: commonpb.MsgType_TimeTick,
|
2021-01-20 09:36:50 +08:00
|
|
|
MsgID: 0,
|
|
|
|
Timestamp: t,
|
|
|
|
SourceID: int64(Params.NodeID),
|
|
|
|
},
|
|
|
|
}
|
|
|
|
timeTickMsg := &ms.TimeTickMsg{
|
|
|
|
BaseMsg: baseMsg,
|
|
|
|
TimeTickMsg: timeTickResult,
|
|
|
|
}
|
|
|
|
msgPack.Msgs = append(msgPack.Msgs, timeTickMsg)
|
2021-02-24 09:48:17 +08:00
|
|
|
if err := timeTickStream.Broadcast(c.ctx, &msgPack); err != nil {
|
2021-01-20 09:36:50 +08:00
|
|
|
return err
|
|
|
|
}
|
2021-02-24 09:48:17 +08:00
|
|
|
if err := ddStream.Broadcast(c.ctx, &msgPack); err != nil {
|
2021-01-20 09:36:50 +08:00
|
|
|
return err
|
|
|
|
}
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2021-03-12 14:22:09 +08:00
|
|
|
c.DdCreateCollectionReq = func(req *internalpb.CreateCollectionRequest) error {
|
2021-01-20 09:36:50 +08:00
|
|
|
msgPack := ms.MsgPack{}
|
|
|
|
baseMsg := ms.BaseMsg{
|
|
|
|
BeginTimestamp: req.Base.Timestamp,
|
|
|
|
EndTimestamp: req.Base.Timestamp,
|
|
|
|
HashValues: []uint32{0},
|
|
|
|
}
|
|
|
|
collMsg := &ms.CreateCollectionMsg{
|
|
|
|
BaseMsg: baseMsg,
|
|
|
|
CreateCollectionRequest: *req,
|
|
|
|
}
|
|
|
|
msgPack.Msgs = append(msgPack.Msgs, collMsg)
|
2021-02-24 09:48:17 +08:00
|
|
|
if err := ddStream.Broadcast(c.ctx, &msgPack); err != nil {
|
2021-01-20 09:36:50 +08:00
|
|
|
return err
|
|
|
|
}
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2021-03-12 14:22:09 +08:00
|
|
|
c.DdDropCollectionReq = func(req *internalpb.DropCollectionRequest) error {
|
2021-01-20 09:36:50 +08:00
|
|
|
msgPack := ms.MsgPack{}
|
|
|
|
baseMsg := ms.BaseMsg{
|
|
|
|
BeginTimestamp: req.Base.Timestamp,
|
|
|
|
EndTimestamp: req.Base.Timestamp,
|
|
|
|
HashValues: []uint32{0},
|
|
|
|
}
|
|
|
|
collMsg := &ms.DropCollectionMsg{
|
|
|
|
BaseMsg: baseMsg,
|
|
|
|
DropCollectionRequest: *req,
|
|
|
|
}
|
|
|
|
msgPack.Msgs = append(msgPack.Msgs, collMsg)
|
2021-02-24 09:48:17 +08:00
|
|
|
if err := ddStream.Broadcast(c.ctx, &msgPack); err != nil {
|
2021-01-20 09:36:50 +08:00
|
|
|
return err
|
|
|
|
}
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2021-03-12 14:22:09 +08:00
|
|
|
c.DdCreatePartitionReq = func(req *internalpb.CreatePartitionRequest) error {
|
2021-01-20 09:36:50 +08:00
|
|
|
msgPack := ms.MsgPack{}
|
|
|
|
baseMsg := ms.BaseMsg{
|
|
|
|
BeginTimestamp: req.Base.Timestamp,
|
|
|
|
EndTimestamp: req.Base.Timestamp,
|
|
|
|
HashValues: []uint32{0},
|
|
|
|
}
|
|
|
|
collMsg := &ms.CreatePartitionMsg{
|
|
|
|
BaseMsg: baseMsg,
|
|
|
|
CreatePartitionRequest: *req,
|
|
|
|
}
|
|
|
|
msgPack.Msgs = append(msgPack.Msgs, collMsg)
|
2021-02-24 09:48:17 +08:00
|
|
|
if err := ddStream.Broadcast(c.ctx, &msgPack); err != nil {
|
2021-01-20 09:36:50 +08:00
|
|
|
return err
|
|
|
|
}
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2021-03-12 14:22:09 +08:00
|
|
|
c.DdDropPartitionReq = func(req *internalpb.DropPartitionRequest) error {
|
2021-01-20 09:36:50 +08:00
|
|
|
msgPack := ms.MsgPack{}
|
|
|
|
baseMsg := ms.BaseMsg{
|
|
|
|
BeginTimestamp: req.Base.Timestamp,
|
|
|
|
EndTimestamp: req.Base.Timestamp,
|
|
|
|
HashValues: []uint32{0},
|
|
|
|
}
|
|
|
|
collMsg := &ms.DropPartitionMsg{
|
|
|
|
BaseMsg: baseMsg,
|
|
|
|
DropPartitionRequest: *req,
|
|
|
|
}
|
|
|
|
msgPack.Msgs = append(msgPack.Msgs, collMsg)
|
2021-02-24 09:48:17 +08:00
|
|
|
if err := ddStream.Broadcast(c.ctx, &msgPack); err != nil {
|
2021-01-20 09:36:50 +08:00
|
|
|
return err
|
|
|
|
}
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
|
|
|
// receive time tick from msg stream
|
2021-01-23 17:56:57 +08:00
|
|
|
c.ProxyTimeTickChan = make(chan typeutil.Timestamp, 1024)
|
2021-01-20 09:36:50 +08:00
|
|
|
go func() {
|
|
|
|
for {
|
|
|
|
select {
|
|
|
|
case <-c.ctx.Done():
|
|
|
|
return
|
|
|
|
case ttmsgs, ok := <-proxyTimeTickStream.Chan():
|
|
|
|
if !ok {
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Warn("proxy time tick msg stream closed")
|
2021-01-20 09:36:50 +08:00
|
|
|
return
|
|
|
|
}
|
|
|
|
if len(ttmsgs.Msgs) > 0 {
|
|
|
|
for _, ttm := range ttmsgs.Msgs {
|
|
|
|
ttmsg, ok := ttm.(*ms.TimeTickMsg)
|
|
|
|
if !ok {
|
|
|
|
continue
|
|
|
|
}
|
|
|
|
c.ProxyTimeTickChan <- ttmsg.Base.Timestamp
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}()
|
|
|
|
|
|
|
|
//segment channel, data service create segment,or data node flush segment will put msg in this channel
|
2021-01-24 20:26:35 +08:00
|
|
|
if Params.DataServiceSegmentChannel == "" {
|
2021-03-05 10:15:27 +08:00
|
|
|
return errors.New("DataServiceSegmentChannel is empty")
|
2021-01-24 20:26:35 +08:00
|
|
|
}
|
2021-02-08 14:30:54 +08:00
|
|
|
dataServiceStream, _ := c.msFactory.NewMsgStream(c.ctx)
|
2021-02-04 14:37:12 +08:00
|
|
|
dataServiceStream.AsConsumer([]string{Params.DataServiceSegmentChannel}, Params.MsgChannelSubName)
|
2021-03-05 18:16:50 +08:00
|
|
|
log.Debug("master AsConsumer: " + Params.DataServiceSegmentChannel + " : " + Params.MsgChannelSubName)
|
2021-01-20 09:36:50 +08:00
|
|
|
dataServiceStream.Start()
|
2021-01-21 10:01:29 +08:00
|
|
|
c.DataServiceSegmentChan = make(chan *datapb.SegmentInfo, 1024)
|
2021-01-23 17:56:57 +08:00
|
|
|
c.DataNodeSegmentFlushCompletedChan = make(chan typeutil.UniqueID, 1024)
|
2021-01-20 09:36:50 +08:00
|
|
|
|
|
|
|
// receive segment info from msg stream
|
|
|
|
go func() {
|
|
|
|
for {
|
|
|
|
select {
|
|
|
|
case <-c.ctx.Done():
|
|
|
|
return
|
|
|
|
case segMsg, ok := <-dataServiceStream.Chan():
|
|
|
|
if !ok {
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Warn("data service segment msg closed")
|
2021-01-20 09:36:50 +08:00
|
|
|
}
|
|
|
|
if len(segMsg.Msgs) > 0 {
|
|
|
|
for _, segm := range segMsg.Msgs {
|
|
|
|
segInfoMsg, ok := segm.(*ms.SegmentInfoMsg)
|
|
|
|
if ok {
|
2021-01-21 10:01:29 +08:00
|
|
|
c.DataServiceSegmentChan <- segInfoMsg.Segment
|
2021-01-23 10:12:41 +08:00
|
|
|
} else {
|
2021-01-23 10:40:34 +08:00
|
|
|
flushMsg, ok := segm.(*ms.FlushCompletedMsg)
|
2021-01-23 10:12:41 +08:00
|
|
|
if ok {
|
|
|
|
c.DataNodeSegmentFlushCompletedChan <- flushMsg.SegmentFlushCompletedMsg.SegmentID
|
|
|
|
} else {
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("receive unexpected msg from data service stream", zap.Stringer("segment", segInfoMsg.SegmentMsg.Segment))
|
2021-01-23 10:12:41 +08:00
|
|
|
}
|
2021-01-20 09:36:50 +08:00
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}()
|
|
|
|
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2021-03-06 17:47:11 +08:00
|
|
|
func (c *Core) SetProxyService(ctx context.Context, s types.ProxyService) error {
|
2021-02-26 17:44:24 +08:00
|
|
|
rsp, err := s.GetTimeTickChannel(ctx)
|
2021-01-24 20:26:35 +08:00
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
2021-02-04 19:34:35 +08:00
|
|
|
Params.ProxyTimeTickChannel = rsp.Value
|
2021-02-27 10:11:52 +08:00
|
|
|
log.Debug("proxy time tick", zap.String("channel name", Params.ProxyTimeTickChannel))
|
2021-01-24 20:26:35 +08:00
|
|
|
|
|
|
|
c.InvalidateCollectionMetaCache = func(ts typeutil.Timestamp, dbName string, collectionName string) error {
|
2021-02-26 17:44:24 +08:00
|
|
|
status, _ := s.InvalidateCollectionMetaCache(ctx, &proxypb.InvalidateCollMetaCacheRequest{
|
2021-01-24 20:26:35 +08:00
|
|
|
Base: &commonpb.MsgBase{
|
|
|
|
MsgType: 0, //TODO,MsgType
|
|
|
|
MsgID: 0,
|
|
|
|
Timestamp: ts,
|
|
|
|
SourceID: int64(Params.NodeID),
|
|
|
|
},
|
|
|
|
DbName: dbName,
|
|
|
|
CollectionName: collectionName,
|
|
|
|
})
|
2021-02-04 19:34:35 +08:00
|
|
|
if status == nil {
|
|
|
|
return errors.New("invalidate collection metacache resp is nil")
|
|
|
|
}
|
2021-03-10 22:06:22 +08:00
|
|
|
if status.ErrorCode != commonpb.ErrorCode_Success {
|
2021-02-04 19:34:35 +08:00
|
|
|
return errors.New(status.Reason)
|
2021-01-24 20:26:35 +08:00
|
|
|
}
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2021-03-06 17:47:11 +08:00
|
|
|
func (c *Core) SetDataService(ctx context.Context, s types.DataService) error {
|
2021-02-26 17:44:24 +08:00
|
|
|
rsp, err := s.GetSegmentInfoChannel(ctx)
|
2021-01-24 20:26:35 +08:00
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
2021-02-04 19:34:35 +08:00
|
|
|
Params.DataServiceSegmentChannel = rsp.Value
|
2021-02-27 10:11:52 +08:00
|
|
|
log.Debug("data service segment", zap.String("channel name", Params.DataServiceSegmentChannel))
|
2021-01-29 17:08:31 +08:00
|
|
|
|
2021-01-24 20:26:35 +08:00
|
|
|
c.GetBinlogFilePathsFromDataServiceReq = func(segID typeutil.UniqueID, fieldID typeutil.UniqueID) ([]string, error) {
|
|
|
|
ts, err := c.tsoAllocator.Alloc(1)
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
2021-03-12 14:22:09 +08:00
|
|
|
binlog, err := s.GetInsertBinlogPaths(ctx, &datapb.GetInsertBinlogPathsRequest{
|
2021-01-24 20:26:35 +08:00
|
|
|
Base: &commonpb.MsgBase{
|
2021-03-08 15:46:51 +08:00
|
|
|
MsgType: 0, //TODO, msg type
|
2021-01-24 20:26:35 +08:00
|
|
|
MsgID: 0,
|
|
|
|
Timestamp: ts,
|
|
|
|
SourceID: int64(Params.NodeID),
|
|
|
|
},
|
|
|
|
SegmentID: segID,
|
|
|
|
})
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
2021-03-10 22:06:22 +08:00
|
|
|
if binlog.Status.ErrorCode != commonpb.ErrorCode_Success {
|
2021-03-05 10:15:27 +08:00
|
|
|
return nil, fmt.Errorf("GetInsertBinlogPaths from data service failed, error = %s", binlog.Status.Reason)
|
2021-01-24 20:26:35 +08:00
|
|
|
}
|
|
|
|
for i := range binlog.FieldIDs {
|
|
|
|
if binlog.FieldIDs[i] == fieldID {
|
|
|
|
return binlog.Paths[i].Values, nil
|
|
|
|
}
|
|
|
|
}
|
2021-03-05 10:15:27 +08:00
|
|
|
return nil, fmt.Errorf("binlog file not exist, segment id = %d, field id = %d", segID, fieldID)
|
2021-01-24 20:26:35 +08:00
|
|
|
}
|
2021-03-08 15:46:51 +08:00
|
|
|
|
|
|
|
c.GetNumRowsReq = func(segID typeutil.UniqueID) (int64, error) {
|
|
|
|
ts, err := c.tsoAllocator.Alloc(1)
|
|
|
|
if err != nil {
|
|
|
|
return 0, err
|
|
|
|
}
|
2021-03-12 14:22:09 +08:00
|
|
|
segInfo, err := s.GetSegmentInfo(ctx, &datapb.GetSegmentInfoRequest{
|
2021-03-08 15:46:51 +08:00
|
|
|
Base: &commonpb.MsgBase{
|
|
|
|
MsgType: 0, //TODO, msg type
|
|
|
|
MsgID: 0,
|
|
|
|
Timestamp: ts,
|
|
|
|
SourceID: int64(Params.NodeID),
|
|
|
|
},
|
|
|
|
SegmentIDs: []typeutil.UniqueID{segID},
|
|
|
|
})
|
|
|
|
if err != nil {
|
|
|
|
return 0, err
|
|
|
|
}
|
2021-03-10 22:06:22 +08:00
|
|
|
if segInfo.Status.ErrorCode != commonpb.ErrorCode_Success {
|
2021-03-08 15:46:51 +08:00
|
|
|
return 0, fmt.Errorf("GetSegmentInfo from data service failed, error = %s", segInfo.Status.Reason)
|
|
|
|
}
|
|
|
|
if len(segInfo.Infos) != 1 {
|
|
|
|
log.Debug("get segment info empty")
|
|
|
|
return 0, nil
|
|
|
|
}
|
|
|
|
if segInfo.Infos[0].FlushedTime == 0 {
|
|
|
|
log.Debug("segment id not flushed", zap.Int64("segment id", segID))
|
|
|
|
return 0, nil
|
|
|
|
}
|
|
|
|
return segInfo.Infos[0].NumRows, nil
|
|
|
|
}
|
2021-01-24 20:26:35 +08:00
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2021-03-06 17:47:11 +08:00
|
|
|
func (c *Core) SetIndexService(ctx context.Context, s types.IndexService) error {
|
2021-02-03 11:52:19 +08:00
|
|
|
c.BuildIndexReq = func(binlog []string, typeParams []*commonpb.KeyValuePair, indexParams []*commonpb.KeyValuePair, indexID typeutil.UniqueID, indexName string) (typeutil.UniqueID, error) {
|
2021-02-26 17:44:24 +08:00
|
|
|
rsp, err := s.BuildIndex(ctx, &indexpb.BuildIndexRequest{
|
2021-01-24 20:26:35 +08:00
|
|
|
DataPaths: binlog,
|
|
|
|
TypeParams: typeParams,
|
|
|
|
IndexParams: indexParams,
|
2021-02-03 11:52:19 +08:00
|
|
|
IndexID: indexID,
|
|
|
|
IndexName: indexName,
|
2021-01-24 20:26:35 +08:00
|
|
|
})
|
|
|
|
if err != nil {
|
|
|
|
return 0, err
|
|
|
|
}
|
2021-03-10 22:06:22 +08:00
|
|
|
if rsp.Status.ErrorCode != commonpb.ErrorCode_Success {
|
2021-03-05 10:15:27 +08:00
|
|
|
return 0, fmt.Errorf("BuildIndex from index service failed, error = %s", rsp.Status.Reason)
|
2021-01-24 20:26:35 +08:00
|
|
|
}
|
2021-02-02 19:56:04 +08:00
|
|
|
return rsp.IndexBuildID, nil
|
2021-01-24 20:26:35 +08:00
|
|
|
}
|
2021-02-20 15:38:44 +08:00
|
|
|
|
|
|
|
c.DropIndexReq = func(indexID typeutil.UniqueID) error {
|
2021-02-26 17:44:24 +08:00
|
|
|
rsp, err := s.DropIndex(ctx, &indexpb.DropIndexRequest{
|
2021-02-20 15:38:44 +08:00
|
|
|
IndexID: indexID,
|
|
|
|
})
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
2021-03-10 22:06:22 +08:00
|
|
|
if rsp.ErrorCode != commonpb.ErrorCode_Success {
|
2021-03-05 10:15:27 +08:00
|
|
|
return errors.New(rsp.Reason)
|
2021-02-20 15:38:44 +08:00
|
|
|
}
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2021-01-24 20:26:35 +08:00
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2021-03-06 17:47:11 +08:00
|
|
|
func (c *Core) SetQueryService(ctx context.Context, s types.QueryService) error {
|
2021-02-05 14:09:55 +08:00
|
|
|
c.ReleaseCollection = func(ts typeutil.Timestamp, dbID typeutil.UniqueID, collectionID typeutil.UniqueID) error {
|
|
|
|
req := &querypb.ReleaseCollectionRequest{
|
|
|
|
Base: &commonpb.MsgBase{
|
2021-03-10 14:45:35 +08:00
|
|
|
MsgType: commonpb.MsgType_ReleaseCollection,
|
2021-02-05 14:09:55 +08:00
|
|
|
MsgID: 0, //TODO, msg ID
|
|
|
|
Timestamp: ts,
|
|
|
|
SourceID: int64(Params.NodeID),
|
|
|
|
},
|
|
|
|
DbID: dbID,
|
|
|
|
CollectionID: collectionID,
|
|
|
|
}
|
2021-02-26 17:44:24 +08:00
|
|
|
rsp, err := s.ReleaseCollection(ctx, req)
|
2021-02-05 14:09:55 +08:00
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
2021-03-10 22:06:22 +08:00
|
|
|
if rsp.ErrorCode != commonpb.ErrorCode_Success {
|
2021-03-05 10:15:27 +08:00
|
|
|
return fmt.Errorf("ReleaseCollection from query service failed, error = %s", rsp.Reason)
|
2021-02-05 14:09:55 +08:00
|
|
|
}
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2021-01-24 20:26:35 +08:00
|
|
|
func (c *Core) Init() error {
|
2021-01-19 14:44:03 +08:00
|
|
|
var initError error = nil
|
|
|
|
c.initOnce.Do(func() {
|
2021-02-26 15:17:47 +08:00
|
|
|
connectEtcdFn := func() error {
|
|
|
|
if c.etcdCli, initError = clientv3.New(clientv3.Config{Endpoints: []string{Params.EtcdAddress}, DialTimeout: 5 * time.Second}); initError != nil {
|
|
|
|
return initError
|
|
|
|
}
|
|
|
|
c.metaKV = etcdkv.NewEtcdKV(c.etcdCli, Params.MetaRootPath)
|
|
|
|
if c.MetaTable, initError = NewMetaTable(c.metaKV); initError != nil {
|
|
|
|
return initError
|
|
|
|
}
|
|
|
|
c.kvBase = etcdkv.NewEtcdKV(c.etcdCli, Params.KvRootPath)
|
|
|
|
return nil
|
2021-01-19 14:44:03 +08:00
|
|
|
}
|
2021-02-26 15:17:47 +08:00
|
|
|
err := retry.Retry(200, time.Millisecond*200, connectEtcdFn)
|
|
|
|
if err != nil {
|
2021-01-19 14:44:03 +08:00
|
|
|
return
|
|
|
|
}
|
|
|
|
|
2021-02-24 17:12:06 +08:00
|
|
|
c.idAllocator = allocator.NewGlobalIDAllocator("idTimestamp", tsoutil.NewTSOKVBase([]string{Params.EtcdAddress}, Params.KvRootPath, "gid"))
|
2021-01-19 14:44:03 +08:00
|
|
|
if initError = c.idAllocator.Initialize(); initError != nil {
|
|
|
|
return
|
|
|
|
}
|
2021-02-24 17:12:06 +08:00
|
|
|
c.tsoAllocator = tso.NewGlobalTSOAllocator("timestamp", tsoutil.NewTSOKVBase([]string{Params.EtcdAddress}, Params.KvRootPath, "tso"))
|
2021-01-19 14:44:03 +08:00
|
|
|
if initError = c.tsoAllocator.Initialize(); initError != nil {
|
|
|
|
return
|
|
|
|
}
|
|
|
|
c.ddReqQueue = make(chan reqTask, 1024)
|
2021-01-21 10:01:29 +08:00
|
|
|
c.indexTaskQueue = make(chan *CreateIndexTask, 1024)
|
2021-01-20 09:36:50 +08:00
|
|
|
initError = c.setMsgStreams()
|
2021-01-19 14:44:03 +08:00
|
|
|
})
|
2021-01-26 19:24:09 +08:00
|
|
|
if initError == nil {
|
2021-03-12 14:22:09 +08:00
|
|
|
log.Debug("Master service", zap.String("State Code", internalpb.StateCode_name[int32(internalpb.StateCode_Initializing)]))
|
2021-01-26 19:24:09 +08:00
|
|
|
}
|
2021-01-19 14:44:03 +08:00
|
|
|
return initError
|
|
|
|
}
|
|
|
|
|
|
|
|
func (c *Core) Start() error {
|
|
|
|
if err := c.checkInit(); err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
c.startOnce.Do(func() {
|
|
|
|
go c.startDdScheduler()
|
|
|
|
go c.startTimeTickLoop()
|
2021-01-21 10:01:29 +08:00
|
|
|
go c.startDataServiceSegmentLoop()
|
|
|
|
go c.startCreateIndexLoop()
|
2021-01-22 15:41:54 +08:00
|
|
|
go c.startSegmentFlushCompletedLoop()
|
2021-01-27 16:38:18 +08:00
|
|
|
go c.tsLoop()
|
2021-03-12 14:22:09 +08:00
|
|
|
c.stateCode.Store(internalpb.StateCode_Healthy)
|
2021-01-19 14:44:03 +08:00
|
|
|
})
|
2021-03-12 14:22:09 +08:00
|
|
|
log.Debug("Master service", zap.String("State Code", internalpb.StateCode_name[int32(internalpb.StateCode_Healthy)]))
|
2021-01-19 14:44:03 +08:00
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
|
|
|
func (c *Core) Stop() error {
|
|
|
|
c.cancel()
|
2021-03-12 14:22:09 +08:00
|
|
|
c.stateCode.Store(internalpb.StateCode_Abnormal)
|
2021-01-19 14:44:03 +08:00
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2021-03-12 14:22:09 +08:00
|
|
|
func (c *Core) GetComponentStates(ctx context.Context) (*internalpb.ComponentStates, error) {
|
|
|
|
code := c.stateCode.Load().(internalpb.StateCode)
|
|
|
|
log.Debug("GetComponentStates", zap.String("State Code", internalpb.StateCode_name[int32(code)]))
|
2021-01-26 19:24:09 +08:00
|
|
|
|
2021-03-12 14:22:09 +08:00
|
|
|
return &internalpb.ComponentStates{
|
|
|
|
State: &internalpb.ComponentInfo{
|
2021-01-20 11:02:29 +08:00
|
|
|
NodeID: int64(Params.NodeID),
|
2021-01-23 18:56:08 +08:00
|
|
|
Role: typeutil.MasterServiceRole,
|
2021-01-20 11:02:29 +08:00
|
|
|
StateCode: code,
|
|
|
|
ExtraInfo: nil,
|
2021-01-19 14:44:03 +08:00
|
|
|
},
|
2021-01-26 17:47:38 +08:00
|
|
|
Status: &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
2021-01-26 17:47:38 +08:00
|
|
|
Reason: "",
|
|
|
|
},
|
2021-03-12 14:22:09 +08:00
|
|
|
SubcomponentStates: []*internalpb.ComponentInfo{
|
2021-01-26 17:47:38 +08:00
|
|
|
{
|
|
|
|
NodeID: int64(Params.NodeID),
|
|
|
|
Role: typeutil.MasterServiceRole,
|
|
|
|
StateCode: code,
|
|
|
|
ExtraInfo: nil,
|
|
|
|
},
|
|
|
|
},
|
2021-01-19 14:44:03 +08:00
|
|
|
}, nil
|
|
|
|
}
|
|
|
|
|
2021-02-26 17:44:24 +08:00
|
|
|
func (c *Core) GetTimeTickChannel(ctx context.Context) (*milvuspb.StringResponse, error) {
|
|
|
|
return &milvuspb.StringResponse{
|
|
|
|
Status: &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
2021-02-26 17:44:24 +08:00
|
|
|
Reason: "",
|
|
|
|
},
|
|
|
|
Value: Params.TimeTickChannel,
|
|
|
|
}, nil
|
2021-01-19 14:44:03 +08:00
|
|
|
}
|
|
|
|
|
2021-02-26 17:44:24 +08:00
|
|
|
func (c *Core) GetDdChannel(ctx context.Context) (*milvuspb.StringResponse, error) {
|
|
|
|
return &milvuspb.StringResponse{
|
|
|
|
Status: &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
2021-02-26 17:44:24 +08:00
|
|
|
Reason: "",
|
|
|
|
},
|
|
|
|
Value: Params.DdChannel,
|
|
|
|
}, nil
|
2021-01-19 14:44:03 +08:00
|
|
|
}
|
|
|
|
|
2021-02-26 17:44:24 +08:00
|
|
|
func (c *Core) GetStatisticsChannel(ctx context.Context) (*milvuspb.StringResponse, error) {
|
|
|
|
return &milvuspb.StringResponse{
|
|
|
|
Status: &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
2021-02-26 17:44:24 +08:00
|
|
|
Reason: "",
|
|
|
|
},
|
|
|
|
Value: Params.StatisticsChannel,
|
|
|
|
}, nil
|
2021-01-19 14:44:03 +08:00
|
|
|
}
|
|
|
|
|
2021-02-26 17:44:24 +08:00
|
|
|
func (c *Core) CreateCollection(ctx context.Context, in *milvuspb.CreateCollectionRequest) (*commonpb.Status, error) {
|
2021-03-12 14:22:09 +08:00
|
|
|
code := c.stateCode.Load().(internalpb.StateCode)
|
|
|
|
if code != internalpb.StateCode_Healthy {
|
2021-01-25 18:33:10 +08:00
|
|
|
return &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-03-12 14:22:09 +08:00
|
|
|
Reason: fmt.Sprintf("state code = %s", internalpb.StateCode_name[int32(code)]),
|
2021-01-25 18:33:10 +08:00
|
|
|
}, nil
|
|
|
|
}
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("CreateCollection ", zap.String("name", in.CollectionName))
|
2021-01-19 14:44:03 +08:00
|
|
|
t := &CreateCollectionReqTask{
|
|
|
|
baseReqTask: baseReqTask{
|
|
|
|
cv: make(chan error),
|
|
|
|
core: c,
|
|
|
|
},
|
|
|
|
Req: in,
|
|
|
|
}
|
|
|
|
c.ddReqQueue <- t
|
|
|
|
err := t.WaitToFinish()
|
|
|
|
if err != nil {
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("CreateCollection failed", zap.String("name", in.CollectionName))
|
2021-01-19 14:44:03 +08:00
|
|
|
return &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-01-19 14:44:03 +08:00
|
|
|
Reason: "Create collection failed: " + err.Error(),
|
|
|
|
}, nil
|
|
|
|
}
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("CreateCollection Success", zap.String("name", in.CollectionName))
|
2021-01-19 14:44:03 +08:00
|
|
|
return &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
2021-01-19 14:44:03 +08:00
|
|
|
Reason: "",
|
|
|
|
}, nil
|
|
|
|
}
|
|
|
|
|
2021-02-26 17:44:24 +08:00
|
|
|
func (c *Core) DropCollection(ctx context.Context, in *milvuspb.DropCollectionRequest) (*commonpb.Status, error) {
|
2021-03-12 14:22:09 +08:00
|
|
|
code := c.stateCode.Load().(internalpb.StateCode)
|
|
|
|
if code != internalpb.StateCode_Healthy {
|
2021-01-25 18:33:10 +08:00
|
|
|
return &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-03-12 14:22:09 +08:00
|
|
|
Reason: fmt.Sprintf("state code = %s", internalpb.StateCode_name[int32(code)]),
|
2021-01-25 18:33:10 +08:00
|
|
|
}, nil
|
|
|
|
}
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("DropCollection", zap.String("name", in.CollectionName))
|
2021-01-19 14:44:03 +08:00
|
|
|
t := &DropCollectionReqTask{
|
|
|
|
baseReqTask: baseReqTask{
|
|
|
|
cv: make(chan error),
|
|
|
|
core: c,
|
|
|
|
},
|
|
|
|
Req: in,
|
|
|
|
}
|
|
|
|
c.ddReqQueue <- t
|
|
|
|
err := t.WaitToFinish()
|
|
|
|
if err != nil {
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("DropCollection Failed", zap.String("name", in.CollectionName))
|
2021-01-19 14:44:03 +08:00
|
|
|
return &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-02-05 11:49:13 +08:00
|
|
|
Reason: "Drop collection failed: " + err.Error(),
|
2021-01-19 14:44:03 +08:00
|
|
|
}, nil
|
|
|
|
}
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("DropCollection Success", zap.String("name", in.CollectionName))
|
2021-01-19 14:44:03 +08:00
|
|
|
return &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
2021-01-19 14:44:03 +08:00
|
|
|
Reason: "",
|
|
|
|
}, nil
|
|
|
|
}
|
|
|
|
|
2021-02-26 17:44:24 +08:00
|
|
|
func (c *Core) HasCollection(ctx context.Context, in *milvuspb.HasCollectionRequest) (*milvuspb.BoolResponse, error) {
|
2021-03-12 14:22:09 +08:00
|
|
|
code := c.stateCode.Load().(internalpb.StateCode)
|
|
|
|
if code != internalpb.StateCode_Healthy {
|
2021-01-25 18:33:10 +08:00
|
|
|
return &milvuspb.BoolResponse{
|
|
|
|
Status: &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-03-12 14:22:09 +08:00
|
|
|
Reason: fmt.Sprintf("state code = %s", internalpb.StateCode_name[int32(code)]),
|
2021-01-25 18:33:10 +08:00
|
|
|
},
|
|
|
|
Value: false,
|
|
|
|
}, nil
|
|
|
|
}
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("HasCollection", zap.String("name", in.CollectionName))
|
2021-01-19 14:44:03 +08:00
|
|
|
t := &HasCollectionReqTask{
|
|
|
|
baseReqTask: baseReqTask{
|
|
|
|
cv: make(chan error),
|
|
|
|
core: c,
|
|
|
|
},
|
|
|
|
Req: in,
|
|
|
|
HasCollection: false,
|
|
|
|
}
|
|
|
|
c.ddReqQueue <- t
|
|
|
|
err := t.WaitToFinish()
|
|
|
|
if err != nil {
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("HasCollection Failed", zap.String("name", in.CollectionName))
|
2021-01-19 14:44:03 +08:00
|
|
|
return &milvuspb.BoolResponse{
|
|
|
|
Status: &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-01-19 14:44:03 +08:00
|
|
|
Reason: "Has collection failed: " + err.Error(),
|
|
|
|
},
|
|
|
|
Value: false,
|
|
|
|
}, nil
|
|
|
|
}
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("HasCollection Success", zap.String("name", in.CollectionName))
|
2021-01-19 14:44:03 +08:00
|
|
|
return &milvuspb.BoolResponse{
|
|
|
|
Status: &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
2021-01-19 14:44:03 +08:00
|
|
|
Reason: "",
|
|
|
|
},
|
|
|
|
Value: t.HasCollection,
|
|
|
|
}, nil
|
|
|
|
}
|
|
|
|
|
2021-02-26 17:44:24 +08:00
|
|
|
func (c *Core) DescribeCollection(ctx context.Context, in *milvuspb.DescribeCollectionRequest) (*milvuspb.DescribeCollectionResponse, error) {
|
2021-03-12 14:22:09 +08:00
|
|
|
code := c.stateCode.Load().(internalpb.StateCode)
|
|
|
|
if code != internalpb.StateCode_Healthy {
|
2021-01-25 18:33:10 +08:00
|
|
|
return &milvuspb.DescribeCollectionResponse{
|
|
|
|
Status: &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-03-12 14:22:09 +08:00
|
|
|
Reason: fmt.Sprintf("state code = %s", internalpb.StateCode_name[int32(code)]),
|
2021-01-25 18:33:10 +08:00
|
|
|
},
|
|
|
|
Schema: nil,
|
|
|
|
CollectionID: 0,
|
|
|
|
}, nil
|
|
|
|
}
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("DescribeCollection", zap.String("name", in.CollectionName))
|
2021-01-19 14:44:03 +08:00
|
|
|
t := &DescribeCollectionReqTask{
|
|
|
|
baseReqTask: baseReqTask{
|
|
|
|
cv: make(chan error),
|
|
|
|
core: c,
|
|
|
|
},
|
|
|
|
Req: in,
|
|
|
|
Rsp: &milvuspb.DescribeCollectionResponse{},
|
|
|
|
}
|
|
|
|
c.ddReqQueue <- t
|
|
|
|
err := t.WaitToFinish()
|
|
|
|
if err != nil {
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("DescribeCollection Failed", zap.String("name", in.CollectionName))
|
2021-01-19 14:44:03 +08:00
|
|
|
return &milvuspb.DescribeCollectionResponse{
|
|
|
|
Status: &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-01-19 14:44:03 +08:00
|
|
|
Reason: "describe collection failed: " + err.Error(),
|
|
|
|
},
|
|
|
|
Schema: nil,
|
|
|
|
}, nil
|
|
|
|
}
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("DescribeCollection Success", zap.String("name", in.CollectionName))
|
2021-01-19 14:44:03 +08:00
|
|
|
t.Rsp.Status = &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
2021-01-19 14:44:03 +08:00
|
|
|
Reason: "",
|
|
|
|
}
|
|
|
|
return t.Rsp, nil
|
|
|
|
}
|
|
|
|
|
2021-03-12 14:22:09 +08:00
|
|
|
func (c *Core) ShowCollections(ctx context.Context, in *milvuspb.ShowCollectionsRequest) (*milvuspb.ShowCollectionsResponse, error) {
|
|
|
|
code := c.stateCode.Load().(internalpb.StateCode)
|
|
|
|
if code != internalpb.StateCode_Healthy {
|
|
|
|
return &milvuspb.ShowCollectionsResponse{
|
2021-01-25 18:33:10 +08:00
|
|
|
Status: &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-03-12 14:22:09 +08:00
|
|
|
Reason: fmt.Sprintf("state code = %s", internalpb.StateCode_name[int32(code)]),
|
2021-01-25 18:33:10 +08:00
|
|
|
},
|
|
|
|
CollectionNames: nil,
|
|
|
|
}, nil
|
|
|
|
}
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("ShowCollections", zap.String("dbname", in.DbName))
|
2021-01-19 14:44:03 +08:00
|
|
|
t := &ShowCollectionReqTask{
|
|
|
|
baseReqTask: baseReqTask{
|
|
|
|
cv: make(chan error),
|
|
|
|
core: c,
|
|
|
|
},
|
|
|
|
Req: in,
|
2021-03-12 14:22:09 +08:00
|
|
|
Rsp: &milvuspb.ShowCollectionsResponse{
|
2021-01-19 14:44:03 +08:00
|
|
|
CollectionNames: nil,
|
|
|
|
},
|
|
|
|
}
|
|
|
|
c.ddReqQueue <- t
|
|
|
|
err := t.WaitToFinish()
|
|
|
|
if err != nil {
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("ShowCollections failed", zap.String("dbname", in.DbName))
|
2021-03-12 14:22:09 +08:00
|
|
|
return &milvuspb.ShowCollectionsResponse{
|
2021-01-19 14:44:03 +08:00
|
|
|
CollectionNames: nil,
|
|
|
|
Status: &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-01-19 14:44:03 +08:00
|
|
|
Reason: "ShowCollections failed: " + err.Error(),
|
|
|
|
},
|
|
|
|
}, nil
|
|
|
|
}
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("ShowCollections Success", zap.String("dbname", in.DbName), zap.Strings("collection names", t.Rsp.CollectionNames))
|
2021-01-19 14:44:03 +08:00
|
|
|
t.Rsp.Status = &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
2021-01-19 14:44:03 +08:00
|
|
|
Reason: "",
|
|
|
|
}
|
|
|
|
return t.Rsp, nil
|
|
|
|
}
|
|
|
|
|
2021-02-26 17:44:24 +08:00
|
|
|
func (c *Core) CreatePartition(ctx context.Context, in *milvuspb.CreatePartitionRequest) (*commonpb.Status, error) {
|
2021-03-12 14:22:09 +08:00
|
|
|
code := c.stateCode.Load().(internalpb.StateCode)
|
|
|
|
if code != internalpb.StateCode_Healthy {
|
2021-01-25 18:33:10 +08:00
|
|
|
return &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-03-12 14:22:09 +08:00
|
|
|
Reason: fmt.Sprintf("state code = %s", internalpb.StateCode_name[int32(code)]),
|
2021-01-25 18:33:10 +08:00
|
|
|
}, nil
|
|
|
|
}
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("CreatePartition", zap.String("collection name", in.CollectionName), zap.String("partition name", in.PartitionName))
|
2021-01-19 14:44:03 +08:00
|
|
|
t := &CreatePartitionReqTask{
|
|
|
|
baseReqTask: baseReqTask{
|
|
|
|
cv: make(chan error),
|
|
|
|
core: c,
|
|
|
|
},
|
|
|
|
Req: in,
|
|
|
|
}
|
|
|
|
c.ddReqQueue <- t
|
|
|
|
err := t.WaitToFinish()
|
|
|
|
if err != nil {
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("CreatePartition Failed", zap.String("collection name", in.CollectionName), zap.String("partition name", in.PartitionName))
|
2021-01-19 14:44:03 +08:00
|
|
|
return &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-01-19 14:44:03 +08:00
|
|
|
Reason: "create partition failed: " + err.Error(),
|
|
|
|
}, nil
|
|
|
|
}
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("CreatePartition Success", zap.String("collection name", in.CollectionName), zap.String("partition name", in.PartitionName))
|
2021-01-19 14:44:03 +08:00
|
|
|
return &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
2021-01-19 14:44:03 +08:00
|
|
|
Reason: "",
|
|
|
|
}, nil
|
|
|
|
}
|
|
|
|
|
2021-02-26 17:44:24 +08:00
|
|
|
func (c *Core) DropPartition(ctx context.Context, in *milvuspb.DropPartitionRequest) (*commonpb.Status, error) {
|
2021-03-12 14:22:09 +08:00
|
|
|
code := c.stateCode.Load().(internalpb.StateCode)
|
|
|
|
if code != internalpb.StateCode_Healthy {
|
2021-01-25 18:33:10 +08:00
|
|
|
return &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-03-12 14:22:09 +08:00
|
|
|
Reason: fmt.Sprintf("state code = %s", internalpb.StateCode_name[int32(code)]),
|
2021-01-25 18:33:10 +08:00
|
|
|
}, nil
|
|
|
|
}
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("DropPartition", zap.String("collection name", in.CollectionName), zap.String("partition name", in.PartitionName))
|
2021-01-19 14:44:03 +08:00
|
|
|
t := &DropPartitionReqTask{
|
|
|
|
baseReqTask: baseReqTask{
|
|
|
|
cv: make(chan error),
|
|
|
|
core: c,
|
|
|
|
},
|
|
|
|
Req: in,
|
|
|
|
}
|
|
|
|
c.ddReqQueue <- t
|
|
|
|
err := t.WaitToFinish()
|
|
|
|
if err != nil {
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("DropPartition Failed", zap.String("collection name", in.CollectionName), zap.String("partition name", in.PartitionName))
|
2021-01-19 14:44:03 +08:00
|
|
|
return &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-01-19 14:44:03 +08:00
|
|
|
Reason: "DropPartition failed: " + err.Error(),
|
|
|
|
}, nil
|
|
|
|
}
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("DropPartition Success", zap.String("collection name", in.CollectionName), zap.String("partition name", in.PartitionName))
|
2021-01-19 14:44:03 +08:00
|
|
|
return &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
2021-01-19 14:44:03 +08:00
|
|
|
Reason: "",
|
|
|
|
}, nil
|
|
|
|
}
|
|
|
|
|
2021-02-26 17:44:24 +08:00
|
|
|
func (c *Core) HasPartition(ctx context.Context, in *milvuspb.HasPartitionRequest) (*milvuspb.BoolResponse, error) {
|
2021-03-12 14:22:09 +08:00
|
|
|
code := c.stateCode.Load().(internalpb.StateCode)
|
|
|
|
if code != internalpb.StateCode_Healthy {
|
2021-01-25 18:33:10 +08:00
|
|
|
return &milvuspb.BoolResponse{
|
|
|
|
Status: &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-03-12 14:22:09 +08:00
|
|
|
Reason: fmt.Sprintf("state code = %s", internalpb.StateCode_name[int32(code)]),
|
2021-01-25 18:33:10 +08:00
|
|
|
},
|
|
|
|
Value: false,
|
|
|
|
}, nil
|
|
|
|
}
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("HasPartition", zap.String("collection name", in.CollectionName), zap.String("partition name", in.PartitionName))
|
2021-01-19 14:44:03 +08:00
|
|
|
t := &HasPartitionReqTask{
|
|
|
|
baseReqTask: baseReqTask{
|
|
|
|
cv: make(chan error),
|
|
|
|
core: c,
|
|
|
|
},
|
|
|
|
Req: in,
|
|
|
|
HasPartition: false,
|
|
|
|
}
|
|
|
|
c.ddReqQueue <- t
|
|
|
|
err := t.WaitToFinish()
|
|
|
|
if err != nil {
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("HasPartition Failed", zap.String("collection name", in.CollectionName), zap.String("partition name", in.PartitionName))
|
2021-01-19 14:44:03 +08:00
|
|
|
return &milvuspb.BoolResponse{
|
|
|
|
Status: &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-01-19 14:44:03 +08:00
|
|
|
Reason: "HasPartition failed: " + err.Error(),
|
|
|
|
},
|
|
|
|
Value: false,
|
|
|
|
}, nil
|
|
|
|
}
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("HasPartition Success", zap.String("collection name", in.CollectionName), zap.String("partition name", in.PartitionName))
|
2021-01-19 14:44:03 +08:00
|
|
|
return &milvuspb.BoolResponse{
|
|
|
|
Status: &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
2021-01-19 14:44:03 +08:00
|
|
|
Reason: "",
|
|
|
|
},
|
|
|
|
Value: t.HasPartition,
|
|
|
|
}, nil
|
|
|
|
}
|
|
|
|
|
2021-03-12 14:22:09 +08:00
|
|
|
func (c *Core) ShowPartitions(ctx context.Context, in *milvuspb.ShowPartitionsRequest) (*milvuspb.ShowPartitionsResponse, error) {
|
2021-03-13 11:59:24 +08:00
|
|
|
log.Debug("ShowPartitionRequest received", zap.String("role", Params.RoleName), zap.Int64("msgID", in.Base.MsgID),
|
|
|
|
zap.String("collection", in.CollectionName))
|
2021-03-12 14:22:09 +08:00
|
|
|
code := c.stateCode.Load().(internalpb.StateCode)
|
|
|
|
if code != internalpb.StateCode_Healthy {
|
2021-03-13 11:59:24 +08:00
|
|
|
log.Error("ShowPartitionRequest failed: master is not healthy", zap.String("role", Params.RoleName),
|
|
|
|
zap.Int64("msgID", in.Base.MsgID), zap.String("state", internalpb.StateCode_name[int32(code)]))
|
2021-03-12 14:22:09 +08:00
|
|
|
return &milvuspb.ShowPartitionsResponse{
|
2021-01-25 18:33:10 +08:00
|
|
|
Status: &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-03-13 11:59:24 +08:00
|
|
|
Reason: fmt.Sprintf("master is not healthy, state code = %s", internalpb.StateCode_name[int32(code)]),
|
2021-01-25 18:33:10 +08:00
|
|
|
},
|
|
|
|
PartitionNames: nil,
|
|
|
|
PartitionIDs: nil,
|
|
|
|
}, nil
|
|
|
|
}
|
2021-01-19 14:44:03 +08:00
|
|
|
t := &ShowPartitionReqTask{
|
|
|
|
baseReqTask: baseReqTask{
|
|
|
|
cv: make(chan error),
|
|
|
|
core: c,
|
|
|
|
},
|
|
|
|
Req: in,
|
2021-03-12 14:22:09 +08:00
|
|
|
Rsp: &milvuspb.ShowPartitionsResponse{
|
2021-01-19 14:44:03 +08:00
|
|
|
PartitionNames: nil,
|
|
|
|
Status: nil,
|
|
|
|
},
|
|
|
|
}
|
|
|
|
c.ddReqQueue <- t
|
|
|
|
err := t.WaitToFinish()
|
|
|
|
if err != nil {
|
2021-03-13 11:59:24 +08:00
|
|
|
log.Error("ShowPartitionsRequest failed", zap.String("role", Params.RoleName), zap.Int64("msgID", in.Base.MsgID), zap.Error(err))
|
2021-03-12 14:22:09 +08:00
|
|
|
return &milvuspb.ShowPartitionsResponse{
|
2021-01-19 14:44:03 +08:00
|
|
|
PartitionNames: nil,
|
|
|
|
Status: &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-03-13 11:59:24 +08:00
|
|
|
Reason: err.Error(),
|
2021-01-19 14:44:03 +08:00
|
|
|
},
|
|
|
|
}, nil
|
|
|
|
}
|
2021-03-13 11:59:24 +08:00
|
|
|
log.Debug("ShowPartitions succeed", zap.String("role", Params.RoleName), zap.Int64("msgID", t.Req.Base.MsgID),
|
|
|
|
zap.String("collection name", in.CollectionName), zap.Strings("partition names", t.Rsp.PartitionNames),
|
|
|
|
zap.Int64s("partition ids", t.Rsp.PartitionIDs))
|
2021-01-19 14:44:03 +08:00
|
|
|
t.Rsp.Status = &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
2021-01-19 14:44:03 +08:00
|
|
|
Reason: "",
|
|
|
|
}
|
|
|
|
return t.Rsp, nil
|
|
|
|
}
|
|
|
|
|
2021-02-26 17:44:24 +08:00
|
|
|
func (c *Core) CreateIndex(ctx context.Context, in *milvuspb.CreateIndexRequest) (*commonpb.Status, error) {
|
2021-03-12 14:22:09 +08:00
|
|
|
code := c.stateCode.Load().(internalpb.StateCode)
|
|
|
|
if code != internalpb.StateCode_Healthy {
|
2021-01-25 18:33:10 +08:00
|
|
|
return &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-03-12 14:22:09 +08:00
|
|
|
Reason: fmt.Sprintf("state code = %s", internalpb.StateCode_name[int32(code)]),
|
2021-01-25 18:33:10 +08:00
|
|
|
}, nil
|
|
|
|
}
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("CreateIndex", zap.String("collection name", in.CollectionName), zap.String("field name", in.FieldName))
|
2021-01-21 10:01:29 +08:00
|
|
|
t := &CreateIndexReqTask{
|
|
|
|
baseReqTask: baseReqTask{
|
|
|
|
cv: make(chan error),
|
|
|
|
core: c,
|
|
|
|
},
|
|
|
|
Req: in,
|
|
|
|
}
|
|
|
|
c.ddReqQueue <- t
|
|
|
|
err := t.WaitToFinish()
|
|
|
|
if err != nil {
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("CreateIndex Failed", zap.String("collection name", in.CollectionName), zap.String("field name", in.FieldName))
|
2021-01-21 10:01:29 +08:00
|
|
|
return &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-01-21 10:01:29 +08:00
|
|
|
Reason: "CreateIndex failed, error = " + err.Error(),
|
|
|
|
}, nil
|
|
|
|
}
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("CreateIndex Success", zap.String("collection name", in.CollectionName), zap.String("field name", in.FieldName))
|
2021-01-21 10:01:29 +08:00
|
|
|
return &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
2021-01-21 10:01:29 +08:00
|
|
|
Reason: "",
|
|
|
|
}, nil
|
2021-01-19 14:44:03 +08:00
|
|
|
}
|
|
|
|
|
2021-02-26 17:44:24 +08:00
|
|
|
func (c *Core) DescribeIndex(ctx context.Context, in *milvuspb.DescribeIndexRequest) (*milvuspb.DescribeIndexResponse, error) {
|
2021-03-12 14:22:09 +08:00
|
|
|
code := c.stateCode.Load().(internalpb.StateCode)
|
|
|
|
if code != internalpb.StateCode_Healthy {
|
2021-01-25 18:33:10 +08:00
|
|
|
return &milvuspb.DescribeIndexResponse{
|
|
|
|
Status: &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-03-12 14:22:09 +08:00
|
|
|
Reason: fmt.Sprintf("state code = %s", internalpb.StateCode_name[int32(code)]),
|
2021-01-25 18:33:10 +08:00
|
|
|
},
|
|
|
|
IndexDescriptions: nil,
|
|
|
|
}, nil
|
|
|
|
}
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("DescribeIndex", zap.String("collection name", in.CollectionName), zap.String("field name", in.FieldName))
|
2021-01-21 10:01:29 +08:00
|
|
|
t := &DescribeIndexReqTask{
|
|
|
|
baseReqTask: baseReqTask{
|
|
|
|
cv: make(chan error),
|
|
|
|
core: c,
|
|
|
|
},
|
|
|
|
Req: in,
|
|
|
|
Rsp: &milvuspb.DescribeIndexResponse{
|
|
|
|
Status: nil,
|
|
|
|
IndexDescriptions: nil,
|
|
|
|
},
|
|
|
|
}
|
|
|
|
c.ddReqQueue <- t
|
|
|
|
err := t.WaitToFinish()
|
|
|
|
if err != nil {
|
|
|
|
return &milvuspb.DescribeIndexResponse{
|
|
|
|
Status: &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-01-21 10:01:29 +08:00
|
|
|
Reason: "DescribeIndex failed, error = " + err.Error(),
|
|
|
|
},
|
|
|
|
IndexDescriptions: nil,
|
|
|
|
}, nil
|
|
|
|
}
|
2021-02-24 16:25:40 +08:00
|
|
|
idxNames := make([]string, 0, len(t.Rsp.IndexDescriptions))
|
|
|
|
for _, i := range t.Rsp.IndexDescriptions {
|
|
|
|
idxNames = append(idxNames, i.IndexName)
|
|
|
|
}
|
|
|
|
log.Debug("DescribeIndex Success", zap.String("collection name", in.CollectionName), zap.String("field name", in.FieldName), zap.Strings("index names", idxNames))
|
2021-03-05 20:41:34 +08:00
|
|
|
if len(t.Rsp.IndexDescriptions) == 0 {
|
|
|
|
t.Rsp.Status = &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_IndexNotExist,
|
2021-03-05 20:41:34 +08:00
|
|
|
Reason: "index not exist",
|
|
|
|
}
|
|
|
|
} else {
|
|
|
|
t.Rsp.Status = &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
2021-03-05 20:41:34 +08:00
|
|
|
Reason: "",
|
|
|
|
}
|
2021-01-21 10:01:29 +08:00
|
|
|
}
|
|
|
|
return t.Rsp, nil
|
2021-01-19 14:44:03 +08:00
|
|
|
}
|
|
|
|
|
2021-02-26 17:44:24 +08:00
|
|
|
func (c *Core) DropIndex(ctx context.Context, in *milvuspb.DropIndexRequest) (*commonpb.Status, error) {
|
2021-03-12 14:22:09 +08:00
|
|
|
code := c.stateCode.Load().(internalpb.StateCode)
|
|
|
|
if code != internalpb.StateCode_Healthy {
|
2021-02-20 15:38:44 +08:00
|
|
|
return &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-03-12 14:22:09 +08:00
|
|
|
Reason: fmt.Sprintf("state code = %s", internalpb.StateCode_name[int32(code)]),
|
2021-02-20 15:38:44 +08:00
|
|
|
}, nil
|
|
|
|
}
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("DropIndex", zap.String("collection name", in.CollectionName), zap.String("field name", in.FieldName), zap.String("index name", in.IndexName))
|
2021-02-20 15:38:44 +08:00
|
|
|
t := &DropIndexReqTask{
|
|
|
|
baseReqTask: baseReqTask{
|
|
|
|
cv: make(chan error),
|
|
|
|
core: c,
|
|
|
|
},
|
|
|
|
Req: in,
|
|
|
|
}
|
|
|
|
c.ddReqQueue <- t
|
|
|
|
err := t.WaitToFinish()
|
|
|
|
if err != nil {
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("DropIndex Failed", zap.String("collection name", in.CollectionName), zap.String("field name", in.FieldName), zap.String("index name", in.IndexName))
|
2021-02-20 15:38:44 +08:00
|
|
|
return &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-02-26 11:07:25 +08:00
|
|
|
Reason: "DropIndex failed, error = " + err.Error(),
|
2021-02-20 15:38:44 +08:00
|
|
|
}, nil
|
|
|
|
}
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("DropIndex Success", zap.String("collection name", in.CollectionName), zap.String("field name", in.FieldName), zap.String("index name", in.IndexName))
|
2021-02-20 15:38:44 +08:00
|
|
|
return &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
2021-02-20 15:38:44 +08:00
|
|
|
Reason: "",
|
|
|
|
}, nil
|
|
|
|
}
|
|
|
|
|
2021-02-26 17:44:24 +08:00
|
|
|
func (c *Core) DescribeSegment(ctx context.Context, in *milvuspb.DescribeSegmentRequest) (*milvuspb.DescribeSegmentResponse, error) {
|
2021-03-12 14:22:09 +08:00
|
|
|
code := c.stateCode.Load().(internalpb.StateCode)
|
|
|
|
if code != internalpb.StateCode_Healthy {
|
2021-01-25 18:33:10 +08:00
|
|
|
return &milvuspb.DescribeSegmentResponse{
|
|
|
|
Status: &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-03-12 14:22:09 +08:00
|
|
|
Reason: fmt.Sprintf("state code = %s", internalpb.StateCode_name[int32(code)]),
|
2021-01-25 18:33:10 +08:00
|
|
|
},
|
|
|
|
IndexID: 0,
|
|
|
|
}, nil
|
|
|
|
}
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("DescribeSegment", zap.Int64("collection id", in.CollectionID), zap.Int64("segment id", in.SegmentID))
|
2021-01-21 10:01:29 +08:00
|
|
|
t := &DescribeSegmentReqTask{
|
|
|
|
baseReqTask: baseReqTask{
|
|
|
|
cv: make(chan error),
|
|
|
|
core: c,
|
|
|
|
},
|
|
|
|
Req: in,
|
|
|
|
Rsp: &milvuspb.DescribeSegmentResponse{
|
|
|
|
Status: nil,
|
|
|
|
IndexID: 0,
|
|
|
|
},
|
|
|
|
}
|
|
|
|
c.ddReqQueue <- t
|
|
|
|
err := t.WaitToFinish()
|
|
|
|
if err != nil {
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("DescribeSegment Failed", zap.Int64("collection id", in.CollectionID), zap.Int64("segment id", in.SegmentID))
|
2021-01-21 10:01:29 +08:00
|
|
|
return &milvuspb.DescribeSegmentResponse{
|
|
|
|
Status: &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-01-21 10:01:29 +08:00
|
|
|
Reason: "DescribeSegment failed, error = " + err.Error(),
|
|
|
|
},
|
|
|
|
IndexID: 0,
|
|
|
|
}, nil
|
|
|
|
}
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("DescribeSegment Success", zap.Int64("collection id", in.CollectionID), zap.Int64("segment id", in.SegmentID))
|
2021-01-21 10:01:29 +08:00
|
|
|
t.Rsp.Status = &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
2021-01-21 10:01:29 +08:00
|
|
|
Reason: "",
|
|
|
|
}
|
|
|
|
return t.Rsp, nil
|
2021-01-19 14:44:03 +08:00
|
|
|
}
|
|
|
|
|
2021-03-12 14:22:09 +08:00
|
|
|
func (c *Core) ShowSegments(ctx context.Context, in *milvuspb.ShowSegmentsRequest) (*milvuspb.ShowSegmentsResponse, error) {
|
|
|
|
code := c.stateCode.Load().(internalpb.StateCode)
|
|
|
|
if code != internalpb.StateCode_Healthy {
|
|
|
|
return &milvuspb.ShowSegmentsResponse{
|
2021-01-25 18:33:10 +08:00
|
|
|
Status: &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-03-12 14:22:09 +08:00
|
|
|
Reason: fmt.Sprintf("state code = %s", internalpb.StateCode_name[int32(code)]),
|
2021-01-25 18:33:10 +08:00
|
|
|
},
|
|
|
|
SegmentIDs: nil,
|
|
|
|
}, nil
|
|
|
|
}
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("ShowSegments", zap.Int64("collection id", in.CollectionID), zap.Int64("partition id", in.PartitionID))
|
2021-01-21 10:01:29 +08:00
|
|
|
t := &ShowSegmentReqTask{
|
|
|
|
baseReqTask: baseReqTask{
|
|
|
|
cv: make(chan error),
|
|
|
|
core: c,
|
|
|
|
},
|
|
|
|
Req: in,
|
2021-03-12 14:22:09 +08:00
|
|
|
Rsp: &milvuspb.ShowSegmentsResponse{
|
2021-01-21 10:01:29 +08:00
|
|
|
Status: nil,
|
|
|
|
SegmentIDs: nil,
|
|
|
|
},
|
|
|
|
}
|
|
|
|
c.ddReqQueue <- t
|
|
|
|
err := t.WaitToFinish()
|
|
|
|
if err != nil {
|
2021-03-12 14:22:09 +08:00
|
|
|
return &milvuspb.ShowSegmentsResponse{
|
2021-01-21 10:01:29 +08:00
|
|
|
Status: &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-01-21 10:01:29 +08:00
|
|
|
Reason: "ShowSegments failed, error: " + err.Error(),
|
|
|
|
},
|
|
|
|
SegmentIDs: nil,
|
|
|
|
}, nil
|
|
|
|
}
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("ShowSegments Success", zap.Int64("collection id", in.CollectionID), zap.Int64("partition id", in.PartitionID), zap.Int64s("segments ids", t.Rsp.SegmentIDs))
|
2021-01-21 10:01:29 +08:00
|
|
|
t.Rsp.Status = &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
2021-01-21 10:01:29 +08:00
|
|
|
Reason: "",
|
|
|
|
}
|
|
|
|
return t.Rsp, nil
|
2021-01-19 14:44:03 +08:00
|
|
|
}
|
|
|
|
|
2021-03-12 14:22:09 +08:00
|
|
|
func (c *Core) AllocTimestamp(ctx context.Context, in *masterpb.AllocTimestampRequest) (*masterpb.AllocTimestampResponse, error) {
|
2021-01-19 14:44:03 +08:00
|
|
|
ts, err := c.tsoAllocator.Alloc(in.Count)
|
|
|
|
if err != nil {
|
2021-03-12 14:22:09 +08:00
|
|
|
return &masterpb.AllocTimestampResponse{
|
2021-01-19 14:44:03 +08:00
|
|
|
Status: &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-01-19 14:44:03 +08:00
|
|
|
Reason: "AllocTimestamp failed: " + err.Error(),
|
|
|
|
},
|
|
|
|
Timestamp: 0,
|
|
|
|
Count: 0,
|
|
|
|
}, nil
|
|
|
|
}
|
2021-02-02 10:58:39 +08:00
|
|
|
// log.Printf("AllocTimestamp : %d", ts)
|
2021-03-12 14:22:09 +08:00
|
|
|
return &masterpb.AllocTimestampResponse{
|
2021-01-19 14:44:03 +08:00
|
|
|
Status: &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
2021-01-19 14:44:03 +08:00
|
|
|
Reason: "",
|
|
|
|
},
|
|
|
|
Timestamp: ts,
|
|
|
|
Count: in.Count,
|
|
|
|
}, nil
|
|
|
|
}
|
|
|
|
|
2021-03-12 14:22:09 +08:00
|
|
|
func (c *Core) AllocID(ctx context.Context, in *masterpb.AllocIDRequest) (*masterpb.AllocIDResponse, error) {
|
2021-01-19 14:44:03 +08:00
|
|
|
start, _, err := c.idAllocator.Alloc(in.Count)
|
|
|
|
if err != nil {
|
2021-03-12 14:22:09 +08:00
|
|
|
return &masterpb.AllocIDResponse{
|
2021-01-19 14:44:03 +08:00
|
|
|
Status: &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-01-19 14:44:03 +08:00
|
|
|
Reason: "AllocID failed: " + err.Error(),
|
|
|
|
},
|
|
|
|
ID: 0,
|
|
|
|
Count: in.Count,
|
|
|
|
}, nil
|
|
|
|
}
|
2021-02-24 16:25:40 +08:00
|
|
|
log.Debug("AllocID", zap.Int64("id start", start), zap.Uint32("count", in.Count))
|
2021-03-12 14:22:09 +08:00
|
|
|
return &masterpb.AllocIDResponse{
|
2021-01-19 14:44:03 +08:00
|
|
|
Status: &commonpb.Status{
|
2021-03-10 22:06:22 +08:00
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
2021-01-19 14:44:03 +08:00
|
|
|
Reason: "",
|
|
|
|
},
|
|
|
|
ID: start,
|
|
|
|
Count: in.Count,
|
|
|
|
}, nil
|
|
|
|
}
|