2021-10-15 10:03:25 +00:00
|
|
|
// Licensed to the LF AI & Data foundation under one
|
|
|
|
// or more contributor license agreements. See the NOTICE file
|
|
|
|
// distributed with this work for additional information
|
|
|
|
// regarding copyright ownership. The ASF licenses this file
|
|
|
|
// to you under the Apache License, Version 2.0 (the
|
|
|
|
// "License"); you may not use this file except in compliance
|
2021-04-19 07:16:33 +00:00
|
|
|
// with the License. You may obtain a copy of the License at
|
|
|
|
//
|
2021-10-15 10:03:25 +00:00
|
|
|
// http://www.apache.org/licenses/LICENSE-2.0
|
2021-04-19 07:16:33 +00:00
|
|
|
//
|
2021-10-15 10:03:25 +00:00
|
|
|
// Unless required by applicable law or agreed to in writing, software
|
|
|
|
// distributed under the License is distributed on an "AS IS" BASIS,
|
|
|
|
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|
|
|
// See the License for the specific language governing permissions and
|
|
|
|
// limitations under the License.
|
2021-04-19 07:16:33 +00:00
|
|
|
|
2021-05-08 07:24:12 +00:00
|
|
|
// Package datanode implements data persistence logic.
|
|
|
|
//
|
2021-06-11 09:53:37 +00:00
|
|
|
// Data node persists insert logs into persistent storage like minIO/S3.
|
2021-01-19 03:37:16 +00:00
|
|
|
package datanode
|
|
|
|
|
|
|
|
import (
|
|
|
|
"context"
|
2021-03-23 10:50:13 +00:00
|
|
|
"errors"
|
2021-02-03 09:30:10 +00:00
|
|
|
"fmt"
|
2021-01-19 03:37:16 +00:00
|
|
|
"io"
|
2021-03-08 07:25:55 +00:00
|
|
|
"math/rand"
|
2021-11-01 02:19:55 +00:00
|
|
|
"path"
|
2021-05-21 11:28:52 +00:00
|
|
|
"strconv"
|
2021-08-13 02:50:09 +00:00
|
|
|
"strings"
|
2021-05-25 07:35:37 +00:00
|
|
|
"sync"
|
2021-02-04 12:31:23 +00:00
|
|
|
"sync/atomic"
|
2021-11-22 08:23:17 +00:00
|
|
|
"syscall"
|
2021-01-24 13:20:11 +00:00
|
|
|
"time"
|
2021-01-19 03:37:16 +00:00
|
|
|
|
2021-08-13 02:50:09 +00:00
|
|
|
"github.com/golang/protobuf/proto"
|
2021-12-23 10:39:11 +00:00
|
|
|
"github.com/milvus-io/milvus/internal/common"
|
2021-11-01 02:19:55 +00:00
|
|
|
"github.com/milvus-io/milvus/internal/kv"
|
2021-08-13 02:50:09 +00:00
|
|
|
etcdkv "github.com/milvus-io/milvus/internal/kv/etcd"
|
2021-11-08 11:49:07 +00:00
|
|
|
miniokv "github.com/milvus-io/milvus/internal/kv/minio"
|
2021-04-22 06:45:57 +00:00
|
|
|
"github.com/milvus-io/milvus/internal/log"
|
2021-08-13 02:50:09 +00:00
|
|
|
"github.com/milvus-io/milvus/internal/logutil"
|
2021-06-21 10:08:15 +00:00
|
|
|
"github.com/milvus-io/milvus/internal/metrics"
|
2021-04-22 06:45:57 +00:00
|
|
|
"github.com/milvus-io/milvus/internal/msgstream"
|
|
|
|
"github.com/milvus-io/milvus/internal/proto/commonpb"
|
|
|
|
"github.com/milvus-io/milvus/internal/proto/datapb"
|
|
|
|
"github.com/milvus-io/milvus/internal/proto/internalpb"
|
|
|
|
"github.com/milvus-io/milvus/internal/proto/milvuspb"
|
2021-06-22 08:14:09 +00:00
|
|
|
"github.com/milvus-io/milvus/internal/proto/rootcoordpb"
|
2021-12-23 10:39:11 +00:00
|
|
|
"github.com/milvus-io/milvus/internal/types"
|
|
|
|
"github.com/milvus-io/milvus/internal/util/metricsinfo"
|
|
|
|
"github.com/milvus-io/milvus/internal/util/paramtable"
|
|
|
|
"github.com/milvus-io/milvus/internal/util/retry"
|
|
|
|
"github.com/milvus-io/milvus/internal/util/sessionutil"
|
|
|
|
"github.com/milvus-io/milvus/internal/util/typeutil"
|
|
|
|
v3rpc "go.etcd.io/etcd/api/v3/v3rpc/rpctypes"
|
|
|
|
clientv3 "go.etcd.io/etcd/client/v3"
|
|
|
|
"go.uber.org/zap"
|
2021-01-19 03:37:16 +00:00
|
|
|
)
|
|
|
|
|
2021-01-24 13:20:11 +00:00
|
|
|
const (
|
2021-10-23 08:57:35 +00:00
|
|
|
// RPCConnectionTimeout is used to set the timeout for rpc request
|
2021-01-24 13:20:11 +00:00
|
|
|
RPCConnectionTimeout = 30 * time.Second
|
2021-06-21 10:08:15 +00:00
|
|
|
|
2021-10-23 08:57:35 +00:00
|
|
|
// MetricRequestsTotal is used to count the num of total requests
|
2021-06-21 10:08:15 +00:00
|
|
|
MetricRequestsTotal = "total"
|
|
|
|
|
2021-10-23 08:57:35 +00:00
|
|
|
// MetricRequestsSuccess is used to count the num of successful requests
|
2021-06-21 10:08:15 +00:00
|
|
|
MetricRequestsSuccess = "success"
|
2021-08-13 02:50:09 +00:00
|
|
|
|
2021-10-23 08:57:35 +00:00
|
|
|
// ConnectEtcdMaxRetryTime is used to limit the max retry time for connection etcd
|
2021-08-13 02:50:09 +00:00
|
|
|
ConnectEtcdMaxRetryTime = 1000
|
2021-01-24 13:20:11 +00:00
|
|
|
)
|
2021-01-22 01:36:40 +00:00
|
|
|
|
2021-09-09 08:40:00 +00:00
|
|
|
const illegalRequestErrStr = "Illegal request"
|
|
|
|
|
2021-10-03 11:46:05 +00:00
|
|
|
// makes sure DataNode implements types.DataNode
|
|
|
|
var _ types.DataNode = (*DataNode)(nil)
|
|
|
|
|
2021-12-23 10:39:11 +00:00
|
|
|
var Params paramtable.GlobalParamTable
|
|
|
|
|
2021-06-29 09:34:13 +00:00
|
|
|
// DataNode communicates with outside services and unioun all
|
|
|
|
// services in datanode package.
|
2021-05-08 07:24:12 +00:00
|
|
|
//
|
2021-06-29 09:34:13 +00:00
|
|
|
// DataNode implements `types.Component`, `types.DataNode` interfaces.
|
2021-09-10 06:46:00 +00:00
|
|
|
// `rootCoord` is a grpc client of root coordinator.
|
|
|
|
// `dataCoord` is a grpc client of data service.
|
2021-06-29 09:34:13 +00:00
|
|
|
// `NodeID` is unique to each datanode.
|
2021-05-25 07:35:37 +00:00
|
|
|
// `State` is current statement of this data node, indicating whether it's healthy.
|
2021-05-08 07:24:12 +00:00
|
|
|
//
|
2021-06-29 09:34:13 +00:00
|
|
|
// `vchan2SyncService` is a map of vchannlName to dataSyncService, so that datanode
|
|
|
|
// has ability to scale flowgraph.
|
|
|
|
// `vchan2FlushCh` holds flush-signal channels for every flowgraph.
|
|
|
|
// `clearSignal` is a signal channel for releasing the flowgraph resources.
|
|
|
|
// `segmentCache` stores all flushing and flushed segments.
|
2021-03-05 12:41:34 +00:00
|
|
|
type DataNode struct {
|
2021-06-29 09:34:13 +00:00
|
|
|
ctx context.Context
|
|
|
|
cancel context.CancelFunc
|
|
|
|
NodeID UniqueID
|
|
|
|
Role string
|
|
|
|
State atomic.Value // internalpb.StateCode_Initializing
|
2021-01-24 13:20:11 +00:00
|
|
|
|
2021-11-11 12:56:49 +00:00
|
|
|
// TODO struct
|
2021-05-28 08:47:29 +00:00
|
|
|
chanMut sync.RWMutex
|
2021-06-07 03:25:37 +00:00
|
|
|
vchan2SyncService map[string]*dataSyncService // vchannel name
|
2021-10-18 04:34:34 +00:00
|
|
|
vchan2FlushChs map[string]chan flushMsg // vchannel name to flush channels
|
2021-09-28 10:22:16 +00:00
|
|
|
|
2021-11-25 01:43:15 +00:00
|
|
|
clearSignal chan string // vchannel name
|
2021-11-08 11:49:07 +00:00
|
|
|
segmentCache *Cache
|
|
|
|
compactionExecutor *compactionExecutor
|
2021-01-24 13:20:11 +00:00
|
|
|
|
2021-06-21 10:22:13 +00:00
|
|
|
rootCoord types.RootCoord
|
|
|
|
dataCoord types.DataCoord
|
2021-01-24 13:20:11 +00:00
|
|
|
|
2021-11-01 02:19:55 +00:00
|
|
|
session *sessionutil.Session
|
|
|
|
watchKv kv.MetaKv
|
2021-11-08 11:49:07 +00:00
|
|
|
blobKv kv.BaseKV
|
2021-05-21 11:28:52 +00:00
|
|
|
|
2021-03-05 12:41:34 +00:00
|
|
|
closer io.Closer
|
2021-02-08 06:30:54 +00:00
|
|
|
|
2021-03-05 12:41:34 +00:00
|
|
|
msFactory msgstream.Factory
|
|
|
|
}
|
2021-01-24 13:20:11 +00:00
|
|
|
|
2021-05-08 07:24:12 +00:00
|
|
|
// NewDataNode will return a DataNode with abnormal state.
|
2021-02-08 06:30:54 +00:00
|
|
|
func NewDataNode(ctx context.Context, factory msgstream.Factory) *DataNode {
|
2021-03-08 07:25:55 +00:00
|
|
|
rand.Seed(time.Now().UnixNano())
|
2021-02-01 03:44:02 +00:00
|
|
|
ctx2, cancel2 := context.WithCancel(ctx)
|
2021-01-19 03:37:16 +00:00
|
|
|
node := &DataNode{
|
2021-06-29 09:34:13 +00:00
|
|
|
ctx: ctx2,
|
|
|
|
cancel: cancel2,
|
|
|
|
Role: typeutil.DataNodeRole,
|
2021-02-04 09:47:19 +00:00
|
|
|
|
2021-11-08 11:49:07 +00:00
|
|
|
rootCoord: nil,
|
|
|
|
dataCoord: nil,
|
|
|
|
msFactory: factory,
|
|
|
|
segmentCache: newCache(),
|
|
|
|
compactionExecutor: newCompactionExecutor(),
|
2021-05-25 07:35:37 +00:00
|
|
|
|
|
|
|
vchan2SyncService: make(map[string]*dataSyncService),
|
2021-10-18 04:34:34 +00:00
|
|
|
vchan2FlushChs: make(map[string]chan flushMsg),
|
2021-11-25 01:43:15 +00:00
|
|
|
clearSignal: make(chan string, 100),
|
2021-01-19 03:37:16 +00:00
|
|
|
}
|
2021-03-12 06:22:09 +00:00
|
|
|
node.UpdateStateCode(internalpb.StateCode_Abnormal)
|
2021-01-19 03:37:16 +00:00
|
|
|
return node
|
|
|
|
}
|
|
|
|
|
2021-10-09 03:28:57 +00:00
|
|
|
// SetRootCoord sets RootCoord's grpc client, error is returned if repeatedly set.
|
2021-10-09 02:10:59 +00:00
|
|
|
func (node *DataNode) SetRootCoord(rc types.RootCoord) error {
|
2021-02-03 09:30:10 +00:00
|
|
|
switch {
|
2021-06-21 09:28:03 +00:00
|
|
|
case rc == nil, node.rootCoord != nil:
|
2021-02-03 09:30:10 +00:00
|
|
|
return errors.New("Nil parameter or repeatly set")
|
|
|
|
default:
|
2021-06-21 09:28:03 +00:00
|
|
|
node.rootCoord = rc
|
2021-02-03 09:30:10 +00:00
|
|
|
return nil
|
|
|
|
}
|
2021-01-26 06:46:54 +00:00
|
|
|
}
|
|
|
|
|
2021-10-09 03:28:57 +00:00
|
|
|
// SetDataCoord sets data service's grpc client, error is returned if repeatedly set.
|
2021-10-09 02:10:59 +00:00
|
|
|
func (node *DataNode) SetDataCoord(ds types.DataCoord) error {
|
2021-02-03 09:30:10 +00:00
|
|
|
switch {
|
2021-06-21 10:22:13 +00:00
|
|
|
case ds == nil, node.dataCoord != nil:
|
2021-02-03 09:30:10 +00:00
|
|
|
return errors.New("Nil parameter or repeatly set")
|
|
|
|
default:
|
2021-06-21 10:22:13 +00:00
|
|
|
node.dataCoord = ds
|
2021-02-03 09:30:10 +00:00
|
|
|
return nil
|
|
|
|
}
|
2021-01-26 06:46:54 +00:00
|
|
|
}
|
|
|
|
|
2021-10-09 02:10:59 +00:00
|
|
|
// SetNodeID set node id for DataNode
|
|
|
|
func (node *DataNode) SetNodeID(id UniqueID) {
|
|
|
|
node.NodeID = id
|
|
|
|
}
|
|
|
|
|
2021-06-29 09:34:13 +00:00
|
|
|
// Register register datanode to etcd
|
2021-05-25 07:06:05 +00:00
|
|
|
func (node *DataNode) Register() error {
|
2021-12-15 03:47:10 +00:00
|
|
|
node.session.Register()
|
|
|
|
|
2021-10-14 08:40:35 +00:00
|
|
|
// Start liveness check
|
|
|
|
go node.session.LivenessCheck(node.ctx, func() {
|
2021-10-30 02:24:38 +00:00
|
|
|
log.Error("Data Node disconnected from etcd, process will exit", zap.Int64("Server Id", node.session.ServerID))
|
|
|
|
if err := node.Stop(); err != nil {
|
|
|
|
log.Fatal("failed to stop server", zap.Error(err))
|
|
|
|
}
|
2021-11-22 08:23:17 +00:00
|
|
|
// manually send signal to starter goroutine
|
|
|
|
syscall.Kill(syscall.Getpid(), syscall.SIGINT)
|
2021-10-14 08:40:35 +00:00
|
|
|
})
|
2021-06-16 11:03:57 +00:00
|
|
|
|
2021-12-15 03:47:10 +00:00
|
|
|
return nil
|
|
|
|
}
|
2021-09-17 13:32:47 +00:00
|
|
|
|
2021-12-15 03:47:10 +00:00
|
|
|
func (node *DataNode) initSession() error {
|
2021-12-23 10:39:11 +00:00
|
|
|
node.session = sessionutil.NewSession(node.ctx, Params.DataNodeCfg.MetaRootPath, Params.DataNodeCfg.EtcdEndpoints)
|
2021-12-15 03:47:10 +00:00
|
|
|
if node.session == nil {
|
|
|
|
return errors.New("failed to initialize session")
|
|
|
|
}
|
2021-12-23 10:39:11 +00:00
|
|
|
node.session.Init(typeutil.DataNodeRole, Params.DataNodeCfg.IP+":"+strconv.Itoa(Params.DataNodeCfg.Port), false)
|
|
|
|
Params.DataNodeCfg.NodeID = node.session.ServerID
|
2021-12-15 03:47:10 +00:00
|
|
|
node.NodeID = node.session.ServerID
|
2021-12-23 10:39:11 +00:00
|
|
|
Params.BaseParams.SetLogger(Params.DataNodeCfg.NodeID)
|
2021-05-25 07:06:05 +00:00
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2021-11-01 02:19:55 +00:00
|
|
|
// Init function does nothing now.
|
2021-01-22 01:36:40 +00:00
|
|
|
func (node *DataNode) Init() error {
|
2021-06-08 11:25:37 +00:00
|
|
|
log.Debug("DataNode Init",
|
2021-12-23 10:39:11 +00:00
|
|
|
zap.String("TimeTickChannelName", Params.DataNodeCfg.TimeTickChannelName),
|
2021-06-08 11:25:37 +00:00
|
|
|
)
|
2021-12-15 03:47:10 +00:00
|
|
|
if err := node.initSession(); err != nil {
|
|
|
|
log.Error("DataNode init session failed", zap.Error(err))
|
|
|
|
return err
|
|
|
|
}
|
2021-12-24 14:22:17 +00:00
|
|
|
Params.DataNodeCfg.Refresh()
|
2021-12-15 03:47:10 +00:00
|
|
|
|
|
|
|
log.Debug("DataNode Init",
|
2021-12-23 10:39:11 +00:00
|
|
|
zap.String("MsgChannelSubName", Params.DataNodeCfg.MsgChannelSubName),
|
2021-12-15 03:47:10 +00:00
|
|
|
)
|
2021-01-24 13:20:11 +00:00
|
|
|
|
2021-05-25 07:35:37 +00:00
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2021-08-13 02:50:09 +00:00
|
|
|
// StartWatchChannels start loop to watch channel allocation status via kv(etcd for now)
|
|
|
|
func (node *DataNode) StartWatchChannels(ctx context.Context) {
|
|
|
|
defer logutil.LogPanic()
|
|
|
|
// REF MEP#7 watch path should be [prefix]/channel/{node_id}/{channel_name}
|
2021-12-23 10:39:11 +00:00
|
|
|
watchPrefix := path.Join(Params.DataNodeCfg.ChannelWatchSubPath, fmt.Sprintf("%d", node.NodeID))
|
2021-11-01 02:19:55 +00:00
|
|
|
evtChan := node.watchKv.WatchWithPrefix(watchPrefix)
|
2021-10-13 09:02:33 +00:00
|
|
|
// after watch, first check all exists nodes first
|
2021-10-15 09:02:37 +00:00
|
|
|
err := node.checkWatchedList()
|
|
|
|
if err != nil {
|
|
|
|
log.Warn("StartWatchChannels failed", zap.Error(err))
|
|
|
|
return
|
|
|
|
}
|
2021-08-13 02:50:09 +00:00
|
|
|
for {
|
|
|
|
select {
|
|
|
|
case <-ctx.Done():
|
|
|
|
log.Debug("watch etcd loop quit")
|
|
|
|
return
|
|
|
|
case event := <-evtChan:
|
2021-10-13 09:02:33 +00:00
|
|
|
if event.Canceled { // event canceled
|
|
|
|
log.Warn("watch channel canceled", zap.Error(event.Err()))
|
|
|
|
// https://github.com/etcd-io/etcd/issues/8980
|
|
|
|
if event.Err() == v3rpc.ErrCompacted {
|
|
|
|
go node.StartWatchChannels(ctx)
|
|
|
|
return
|
|
|
|
}
|
2021-08-13 02:50:09 +00:00
|
|
|
// if watch loop return due to event canceled, the datanode is not functional anymore
|
|
|
|
// stop the datanode and wait for restart
|
2021-10-07 14:16:56 +00:00
|
|
|
err := node.Stop()
|
|
|
|
if err != nil {
|
|
|
|
log.Warn("node stop failed", zap.Error(err))
|
|
|
|
}
|
2021-08-13 02:50:09 +00:00
|
|
|
return
|
|
|
|
}
|
|
|
|
for _, evt := range event.Events {
|
|
|
|
go node.handleChannelEvt(evt)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2021-10-13 09:02:33 +00:00
|
|
|
// checkWatchedList list all nodes under [prefix]/channel/{node_id} and make sure all nodeds are watched
|
|
|
|
// serves the corner case for etcd connection lost and missing some events
|
|
|
|
func (node *DataNode) checkWatchedList() error {
|
|
|
|
// REF MEP#7 watch path should be [prefix]/channel/{node_id}/{channel_name}
|
2021-12-23 10:39:11 +00:00
|
|
|
prefix := path.Join(Params.DataNodeCfg.ChannelWatchSubPath, fmt.Sprintf("%d", node.NodeID))
|
2021-11-01 02:19:55 +00:00
|
|
|
keys, values, err := node.watchKv.LoadWithPrefix(prefix)
|
2021-10-13 09:02:33 +00:00
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
for i, val := range values {
|
|
|
|
node.handleWatchInfo(keys[i], []byte(val))
|
|
|
|
}
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2021-09-10 06:46:00 +00:00
|
|
|
// handleChannelEvt handles event from kv watch event
|
2021-08-13 02:50:09 +00:00
|
|
|
func (node *DataNode) handleChannelEvt(evt *clientv3.Event) {
|
|
|
|
switch evt.Type {
|
|
|
|
case clientv3.EventTypePut: // datacoord shall put channels needs to be watched here
|
2021-10-13 09:02:33 +00:00
|
|
|
node.handleWatchInfo(string(evt.Kv.Key), evt.Kv.Value)
|
2021-08-13 02:50:09 +00:00
|
|
|
case clientv3.EventTypeDelete:
|
|
|
|
// guaranteed there is no "/" in channel name
|
|
|
|
parts := strings.Split(string(evt.Kv.Key), "/")
|
2021-12-02 08:39:33 +00:00
|
|
|
vchanName := parts[len(parts)-1]
|
2021-12-23 10:39:11 +00:00
|
|
|
log.Warn("handle channel delete event", zap.Int64("node id", Params.DataNodeCfg.NodeID), zap.String("vchannel", vchanName))
|
2021-12-02 08:39:33 +00:00
|
|
|
node.ReleaseDataSyncService(vchanName)
|
2021-08-13 02:50:09 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2021-10-13 09:02:33 +00:00
|
|
|
func (node *DataNode) handleWatchInfo(key string, data []byte) {
|
|
|
|
watchInfo := datapb.ChannelWatchInfo{}
|
|
|
|
err := proto.Unmarshal(data, &watchInfo)
|
|
|
|
if err != nil {
|
|
|
|
log.Warn("fail to parse ChannelWatchInfo", zap.String("key", key), zap.Error(err))
|
|
|
|
return
|
|
|
|
}
|
|
|
|
if watchInfo.State == datapb.ChannelWatchState_Complete {
|
|
|
|
return
|
|
|
|
}
|
|
|
|
if watchInfo.Vchan == nil {
|
|
|
|
log.Warn("found ChannelWatchInfo with nil VChannelInfo", zap.String("key", key))
|
|
|
|
return
|
|
|
|
}
|
|
|
|
err = node.NewDataSyncService(watchInfo.Vchan)
|
|
|
|
if err != nil {
|
|
|
|
log.Warn("fail to create DataSyncService", zap.String("key", key), zap.Error(err))
|
|
|
|
return
|
|
|
|
}
|
|
|
|
watchInfo.State = datapb.ChannelWatchState_Complete
|
|
|
|
v, err := proto.Marshal(&watchInfo)
|
|
|
|
if err != nil {
|
|
|
|
log.Warn("fail to Marshal watchInfo", zap.String("key", key), zap.Error(err))
|
|
|
|
return
|
|
|
|
}
|
2021-12-23 10:39:11 +00:00
|
|
|
k := path.Join(Params.DataNodeCfg.ChannelWatchSubPath, fmt.Sprintf("%d", node.NodeID), watchInfo.GetVchan().GetChannelName())
|
2021-11-01 02:19:55 +00:00
|
|
|
err = node.watchKv.Save(k, string(v))
|
2021-10-13 09:02:33 +00:00
|
|
|
if err != nil {
|
|
|
|
log.Warn("fail to change WatchState to complete", zap.String("key", key), zap.Error(err))
|
|
|
|
node.ReleaseDataSyncService(key)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2021-05-28 06:54:31 +00:00
|
|
|
// NewDataSyncService adds a new dataSyncService for new dmlVchannel and starts dataSyncService.
|
2021-06-04 01:57:54 +00:00
|
|
|
func (node *DataNode) NewDataSyncService(vchan *datapb.VchannelInfo) error {
|
2021-05-28 08:47:29 +00:00
|
|
|
node.chanMut.Lock()
|
|
|
|
defer node.chanMut.Unlock()
|
2021-06-04 01:57:54 +00:00
|
|
|
if _, ok := node.vchan2SyncService[vchan.GetChannelName()]; ok {
|
2021-05-25 07:35:37 +00:00
|
|
|
return nil
|
2021-05-08 07:24:12 +00:00
|
|
|
}
|
|
|
|
|
2021-10-14 02:24:33 +00:00
|
|
|
replica, err := newReplica(node.ctx, node.rootCoord, vchan.CollectionID)
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
2021-05-29 10:04:30 +00:00
|
|
|
|
2021-06-21 09:28:03 +00:00
|
|
|
var alloc allocatorInterface = newAllocator(node.rootCoord)
|
2021-02-03 09:30:10 +00:00
|
|
|
|
2021-06-11 01:24:52 +00:00
|
|
|
log.Debug("Received Vchannel Info",
|
|
|
|
zap.Int("Unflushed Segment Number", len(vchan.GetUnflushedSegments())),
|
|
|
|
zap.Int("Flushed Segment Number", len(vchan.GetFlushedSegments())),
|
|
|
|
)
|
|
|
|
|
2021-10-18 04:34:34 +00:00
|
|
|
flushCh := make(chan flushMsg, 100)
|
2021-09-28 10:22:16 +00:00
|
|
|
|
2021-12-02 08:39:33 +00:00
|
|
|
dataSyncService, err := newDataSyncService(node.ctx, flushCh, replica, alloc, node.msFactory, vchan, node.clearSignal, node.dataCoord, node.segmentCache, node.blobKv, node.compactionExecutor)
|
2021-06-19 07:18:06 +00:00
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
2021-06-29 09:34:13 +00:00
|
|
|
|
2021-06-04 01:57:54 +00:00
|
|
|
node.vchan2SyncService[vchan.GetChannelName()] = dataSyncService
|
2021-10-18 04:34:34 +00:00
|
|
|
node.vchan2FlushChs[vchan.GetChannelName()] = flushCh
|
2021-05-25 07:35:37 +00:00
|
|
|
|
2021-06-29 09:34:13 +00:00
|
|
|
log.Info("Start New dataSyncService",
|
|
|
|
zap.Int64("Collection ID", vchan.GetCollectionID()),
|
|
|
|
zap.String("Vchannel name", vchan.GetChannelName()),
|
|
|
|
)
|
2021-09-29 12:19:59 +00:00
|
|
|
dataSyncService.start()
|
2021-05-25 07:35:37 +00:00
|
|
|
|
2021-01-24 13:20:11 +00:00
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2021-06-07 03:25:37 +00:00
|
|
|
// BackGroundGC runs in background to release datanode resources
|
2021-11-25 01:43:15 +00:00
|
|
|
func (node *DataNode) BackGroundGC(vChannelCh <-chan string) {
|
2021-06-07 03:25:37 +00:00
|
|
|
log.Info("DataNode Background GC Start")
|
|
|
|
for {
|
|
|
|
select {
|
2021-11-25 01:43:15 +00:00
|
|
|
case vChan := <-vChannelCh:
|
|
|
|
log.Info("GC flowgraph", zap.String("vChan", vChan))
|
|
|
|
node.ReleaseDataSyncService(vChan)
|
2021-06-07 03:25:37 +00:00
|
|
|
case <-node.ctx.Done():
|
2021-06-09 09:31:48 +00:00
|
|
|
log.Info("DataNode ctx done")
|
2021-06-07 03:25:37 +00:00
|
|
|
return
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
// ReleaseDataSyncService release flowgraph resources for a vchanName
|
|
|
|
func (node *DataNode) ReleaseDataSyncService(vchanName string) {
|
|
|
|
log.Info("Release flowgraph resources begin", zap.String("Vchannel", vchanName))
|
|
|
|
|
|
|
|
node.chanMut.Lock()
|
2021-06-11 09:53:37 +00:00
|
|
|
defer node.chanMut.Unlock()
|
2021-06-07 03:25:37 +00:00
|
|
|
if dss, ok := node.vchan2SyncService[vchanName]; ok {
|
|
|
|
dss.close()
|
|
|
|
}
|
|
|
|
|
|
|
|
delete(node.vchan2SyncService, vchanName)
|
2021-09-28 10:22:16 +00:00
|
|
|
delete(node.vchan2FlushChs, vchanName)
|
2021-06-07 03:25:37 +00:00
|
|
|
|
|
|
|
log.Debug("Release flowgraph resources end", zap.String("Vchannel", vchanName))
|
|
|
|
}
|
|
|
|
|
2021-09-24 12:44:05 +00:00
|
|
|
// FilterThreshold is the start time ouf DataNode
|
2021-06-07 05:58:37 +00:00
|
|
|
var FilterThreshold Timestamp
|
|
|
|
|
2021-05-28 06:54:31 +00:00
|
|
|
// Start will update DataNode state to HEALTHY
|
2021-01-24 13:20:11 +00:00
|
|
|
func (node *DataNode) Start() error {
|
2021-06-07 05:58:37 +00:00
|
|
|
|
2021-06-22 08:14:09 +00:00
|
|
|
rep, err := node.rootCoord.AllocTimestamp(node.ctx, &rootcoordpb.AllocTimestampRequest{
|
2021-06-07 05:58:37 +00:00
|
|
|
Base: &commonpb.MsgBase{
|
|
|
|
MsgType: commonpb.MsgType_RequestTSO,
|
|
|
|
MsgID: 0,
|
|
|
|
Timestamp: 0,
|
|
|
|
SourceID: node.NodeID,
|
|
|
|
},
|
|
|
|
Count: 1,
|
|
|
|
})
|
2021-08-13 02:50:09 +00:00
|
|
|
if err != nil {
|
|
|
|
log.Warn("fail to alloc timestamp", zap.Error(err))
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
|
|
|
connectEtcdFn := func() error {
|
2021-12-23 10:39:11 +00:00
|
|
|
etcdKV, err := etcdkv.NewEtcdKV(Params.DataNodeCfg.EtcdEndpoints, Params.DataNodeCfg.MetaRootPath)
|
2021-08-13 02:50:09 +00:00
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
2021-11-01 02:19:55 +00:00
|
|
|
node.watchKv = etcdKV
|
2021-08-13 02:50:09 +00:00
|
|
|
return nil
|
|
|
|
}
|
|
|
|
err = retry.Do(node.ctx, connectEtcdFn, retry.Attempts(ConnectEtcdMaxRetryTime))
|
|
|
|
if err != nil {
|
|
|
|
return errors.New("DataNode fail to connect etcd")
|
|
|
|
}
|
2021-06-07 05:58:37 +00:00
|
|
|
|
2021-11-08 11:49:07 +00:00
|
|
|
option := &miniokv.Option{
|
2021-12-23 10:39:11 +00:00
|
|
|
Address: Params.DataNodeCfg.MinioAddress,
|
|
|
|
AccessKeyID: Params.DataNodeCfg.MinioAccessKeyID,
|
|
|
|
SecretAccessKeyID: Params.DataNodeCfg.MinioSecretAccessKey,
|
|
|
|
UseSSL: Params.DataNodeCfg.MinioUseSSL,
|
2021-11-08 11:49:07 +00:00
|
|
|
CreateBucket: true,
|
2021-12-23 10:39:11 +00:00
|
|
|
BucketName: Params.DataNodeCfg.MinioBucketName,
|
2021-11-08 11:49:07 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
kv, err := miniokv.NewMinIOKV(node.ctx, option)
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
|
|
|
node.blobKv = kv
|
|
|
|
|
2021-06-07 05:58:37 +00:00
|
|
|
if rep.Status.ErrorCode != commonpb.ErrorCode_Success || err != nil {
|
|
|
|
return errors.New("DataNode fail to start")
|
|
|
|
}
|
|
|
|
|
|
|
|
FilterThreshold = rep.GetTimestamp()
|
|
|
|
|
2021-06-07 03:25:37 +00:00
|
|
|
go node.BackGroundGC(node.clearSignal)
|
2021-09-17 13:32:47 +00:00
|
|
|
|
2021-11-08 11:49:07 +00:00
|
|
|
go node.compactionExecutor.start(node.ctx)
|
|
|
|
|
2021-12-15 03:47:10 +00:00
|
|
|
// Start node watch node
|
|
|
|
go node.StartWatchChannels(node.ctx)
|
|
|
|
|
2021-12-23 10:39:11 +00:00
|
|
|
Params.DataNodeCfg.CreatedTime = time.Now()
|
|
|
|
Params.DataNodeCfg.UpdatedTime = time.Now()
|
2021-09-17 13:32:47 +00:00
|
|
|
|
2021-03-12 06:22:09 +00:00
|
|
|
node.UpdateStateCode(internalpb.StateCode_Healthy)
|
2021-02-04 09:47:19 +00:00
|
|
|
return nil
|
2021-01-19 03:37:16 +00:00
|
|
|
}
|
|
|
|
|
2021-05-28 06:54:31 +00:00
|
|
|
// UpdateStateCode updates datanode's state code
|
2021-03-12 06:22:09 +00:00
|
|
|
func (node *DataNode) UpdateStateCode(code internalpb.StateCode) {
|
2021-02-23 03:40:30 +00:00
|
|
|
node.State.Store(code)
|
|
|
|
}
|
|
|
|
|
2021-10-09 02:10:59 +00:00
|
|
|
// GetStateCode return datanode's state code
|
|
|
|
func (node *DataNode) GetStateCode() internalpb.StateCode {
|
|
|
|
return node.State.Load().(internalpb.StateCode)
|
|
|
|
}
|
|
|
|
|
2021-09-01 02:13:15 +00:00
|
|
|
func (node *DataNode) isHealthy() bool {
|
|
|
|
code := node.State.Load().(internalpb.StateCode)
|
|
|
|
return code == internalpb.StateCode_Healthy
|
|
|
|
}
|
|
|
|
|
2021-11-01 03:01:49 +00:00
|
|
|
// WatchDmChannels is not in use
|
2021-03-12 06:22:09 +00:00
|
|
|
func (node *DataNode) WatchDmChannels(ctx context.Context, in *datapb.WatchDmChannelsRequest) (*commonpb.Status, error) {
|
2021-11-01 03:01:49 +00:00
|
|
|
log.Warn("DataNode WatchDmChannels is not in use")
|
2021-02-03 09:30:10 +00:00
|
|
|
|
2021-11-01 03:01:49 +00:00
|
|
|
return &commonpb.Status{
|
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
|
|
|
Reason: "watchDmChannels do nothing",
|
|
|
|
}, nil
|
2021-01-22 01:36:40 +00:00
|
|
|
}
|
|
|
|
|
2021-05-29 10:04:30 +00:00
|
|
|
// GetComponentStates will return current state of DataNode
|
2021-03-12 06:22:09 +00:00
|
|
|
func (node *DataNode) GetComponentStates(ctx context.Context) (*internalpb.ComponentStates, error) {
|
2021-02-26 02:13:36 +00:00
|
|
|
log.Debug("DataNode current state", zap.Any("State", node.State.Load()))
|
2021-11-19 05:57:12 +00:00
|
|
|
nodeID := common.NotRegisteredID
|
|
|
|
if node.session != nil && node.session.Registered() {
|
|
|
|
nodeID = node.session.ServerID
|
|
|
|
}
|
2021-03-12 06:22:09 +00:00
|
|
|
states := &internalpb.ComponentStates{
|
|
|
|
State: &internalpb.ComponentInfo{
|
2021-11-19 05:57:12 +00:00
|
|
|
// NodeID: Params.NodeID, // will race with DataNode.Register()
|
|
|
|
NodeID: nodeID,
|
2021-01-24 13:20:11 +00:00
|
|
|
Role: node.Role,
|
2021-03-12 06:22:09 +00:00
|
|
|
StateCode: node.State.Load().(internalpb.StateCode),
|
2021-01-24 13:20:11 +00:00
|
|
|
},
|
2021-03-12 06:22:09 +00:00
|
|
|
SubcomponentStates: make([]*internalpb.ComponentInfo, 0),
|
2021-03-10 14:06:22 +00:00
|
|
|
Status: &commonpb.Status{ErrorCode: commonpb.ErrorCode_Success},
|
2021-01-24 13:20:11 +00:00
|
|
|
}
|
|
|
|
return states, nil
|
2021-01-22 01:36:40 +00:00
|
|
|
}
|
|
|
|
|
2021-06-07 03:25:37 +00:00
|
|
|
func (node *DataNode) getChannelNamebySegmentID(segID UniqueID) string {
|
2021-05-28 08:47:29 +00:00
|
|
|
node.chanMut.RLock()
|
|
|
|
defer node.chanMut.RUnlock()
|
2021-05-25 07:35:37 +00:00
|
|
|
for name, dataSync := range node.vchan2SyncService {
|
2021-10-25 10:01:50 +00:00
|
|
|
if dataSync.replica.hasSegment(segID, true) {
|
2021-05-25 07:35:37 +00:00
|
|
|
return name
|
|
|
|
}
|
|
|
|
}
|
|
|
|
return ""
|
|
|
|
}
|
|
|
|
|
2021-06-07 03:25:37 +00:00
|
|
|
func (node *DataNode) getChannelNamesbyCollectionID(collID UniqueID) []string {
|
|
|
|
node.chanMut.RLock()
|
|
|
|
defer node.chanMut.RUnlock()
|
|
|
|
|
|
|
|
channels := make([]string, 0, len(node.vchan2SyncService))
|
|
|
|
for name, dataSync := range node.vchan2SyncService {
|
|
|
|
if dataSync.collectionID == collID {
|
|
|
|
channels = append(channels, name)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
return channels
|
|
|
|
}
|
|
|
|
|
2021-05-29 10:04:30 +00:00
|
|
|
// ReadyToFlush tells wether DataNode is ready for flushing
|
|
|
|
func (node *DataNode) ReadyToFlush() error {
|
|
|
|
if node.State.Load().(internalpb.StateCode) != internalpb.StateCode_Healthy {
|
|
|
|
return errors.New("DataNode not in HEALTHY state")
|
|
|
|
}
|
|
|
|
|
|
|
|
node.chanMut.RLock()
|
|
|
|
defer node.chanMut.RUnlock()
|
2021-09-28 10:22:16 +00:00
|
|
|
if len(node.vchan2SyncService) == 0 && len(node.vchan2FlushChs) == 0 {
|
2021-05-29 10:04:30 +00:00
|
|
|
// Healthy but Idle
|
|
|
|
msg := "DataNode HEALTHY but IDLE, please try WatchDmChannels to make it work"
|
2021-06-11 09:53:37 +00:00
|
|
|
log.Warn(msg)
|
2021-05-29 10:04:30 +00:00
|
|
|
return errors.New(msg)
|
|
|
|
}
|
|
|
|
|
2021-09-28 10:22:16 +00:00
|
|
|
if len(node.vchan2SyncService) != len(node.vchan2FlushChs) {
|
2021-05-29 10:04:30 +00:00
|
|
|
// TODO restart
|
|
|
|
msg := "DataNode HEALTHY but abnormal inside, restarting..."
|
2021-06-11 09:53:37 +00:00
|
|
|
log.Warn(msg)
|
2021-05-29 10:04:30 +00:00
|
|
|
return errors.New(msg)
|
|
|
|
}
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2021-05-08 07:24:12 +00:00
|
|
|
// FlushSegments packs flush messages into flowgraph through flushChan.
|
2021-05-28 06:54:31 +00:00
|
|
|
// If DataNode receives a valid segment to flush, new flush message for the segment should be ignored.
|
|
|
|
// So if receiving calls to flush segment A, DataNode should guarantee the segment to be flushed.
|
2021-05-29 10:04:30 +00:00
|
|
|
//
|
2021-10-11 10:50:30 +00:00
|
|
|
// One precondition: The segmentID in req is in ascending order.
|
2021-03-12 06:22:09 +00:00
|
|
|
func (node *DataNode) FlushSegments(ctx context.Context, req *datapb.FlushSegmentsRequest) (*commonpb.Status, error) {
|
2021-06-21 10:08:15 +00:00
|
|
|
metrics.DataNodeFlushSegmentsCounter.WithLabelValues(MetricRequestsTotal).Inc()
|
2021-05-25 07:35:37 +00:00
|
|
|
status := &commonpb.Status{
|
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-01-26 06:46:54 +00:00
|
|
|
}
|
|
|
|
|
2021-05-29 10:04:30 +00:00
|
|
|
if err := node.ReadyToFlush(); err != nil {
|
|
|
|
status.Reason = err.Error()
|
|
|
|
return status, nil
|
|
|
|
}
|
|
|
|
|
2021-09-23 08:03:54 +00:00
|
|
|
log.Debug("Receive FlushSegments req",
|
2021-10-11 10:50:30 +00:00
|
|
|
zap.Int64("collectionID", req.GetCollectionID()), zap.Int("num", len(req.SegmentIDs)),
|
2021-06-29 09:34:13 +00:00
|
|
|
zap.Int64s("segments", req.SegmentIDs),
|
|
|
|
)
|
|
|
|
|
2021-10-20 07:02:36 +00:00
|
|
|
processSegments := func(segmentIDs []UniqueID, flushed bool) bool {
|
|
|
|
noErr := true
|
|
|
|
for _, id := range segmentIDs {
|
|
|
|
chanName := node.getChannelNamebySegmentID(id)
|
|
|
|
if len(chanName) == 0 {
|
|
|
|
log.Warn("FlushSegments failed, cannot find segment in DataNode replica",
|
|
|
|
zap.Int64("collectionID", req.GetCollectionID()), zap.Int64("segmentID", id))
|
|
|
|
|
|
|
|
status.Reason = fmt.Sprintf("DataNode replica not find segment %d!", id)
|
|
|
|
noErr = false
|
|
|
|
continue
|
|
|
|
}
|
2021-05-29 10:04:30 +00:00
|
|
|
|
2021-10-20 07:02:36 +00:00
|
|
|
if node.segmentCache.checkIfCached(id) {
|
|
|
|
// Segment in flushing, ignore
|
|
|
|
log.Info("Segment flushing, ignore the flush request until flush is done.",
|
|
|
|
zap.Int64("collectionID", req.GetCollectionID()), zap.Int64("segmentID", id))
|
2021-10-11 10:50:30 +00:00
|
|
|
|
2021-10-20 07:02:36 +00:00
|
|
|
continue
|
|
|
|
}
|
2021-06-11 09:53:37 +00:00
|
|
|
|
2021-10-20 07:02:36 +00:00
|
|
|
node.segmentCache.Cache(id)
|
2021-06-11 09:53:37 +00:00
|
|
|
|
2021-10-20 07:02:36 +00:00
|
|
|
node.chanMut.RLock()
|
|
|
|
flushChs, ok := node.vchan2FlushChs[chanName]
|
|
|
|
node.chanMut.RUnlock()
|
|
|
|
if !ok {
|
|
|
|
status.Reason = "DataNode abnormal, restarting"
|
|
|
|
log.Error("DataNode abnormal, no flushCh for a vchannel")
|
|
|
|
noErr = false
|
|
|
|
continue
|
|
|
|
}
|
2021-05-25 07:35:37 +00:00
|
|
|
|
2021-10-20 07:02:36 +00:00
|
|
|
flushChs <- flushMsg{
|
|
|
|
msgID: req.Base.MsgID,
|
|
|
|
timestamp: req.Base.Timestamp,
|
|
|
|
segmentID: id,
|
|
|
|
collectionID: req.CollectionID,
|
|
|
|
flushed: flushed,
|
|
|
|
}
|
2021-05-25 07:35:37 +00:00
|
|
|
}
|
2021-10-20 07:02:36 +00:00
|
|
|
log.Debug("Flowgraph flushSegment tasks triggered", zap.Bool("flushed", flushed),
|
|
|
|
zap.Int64("collectionID", req.GetCollectionID()), zap.Int64s("segments", segmentIDs))
|
|
|
|
|
|
|
|
return noErr
|
2021-05-25 07:35:37 +00:00
|
|
|
}
|
2021-08-11 06:24:09 +00:00
|
|
|
|
2021-10-20 07:02:36 +00:00
|
|
|
ok := processSegments(req.GetSegmentIDs(), true)
|
|
|
|
if !ok {
|
|
|
|
return status, nil
|
|
|
|
}
|
|
|
|
ok = processSegments(req.GetMarkSegmentIDs(), false)
|
|
|
|
if !ok {
|
|
|
|
return status, nil
|
|
|
|
}
|
2021-05-25 07:35:37 +00:00
|
|
|
|
|
|
|
status.ErrorCode = commonpb.ErrorCode_Success
|
2021-06-21 10:08:15 +00:00
|
|
|
metrics.DataNodeFlushSegmentsCounter.WithLabelValues(MetricRequestsSuccess).Inc()
|
2021-05-25 07:35:37 +00:00
|
|
|
return status, nil
|
2021-01-22 01:36:40 +00:00
|
|
|
}
|
|
|
|
|
2021-09-24 12:44:05 +00:00
|
|
|
// Stop will release DataNode resources and shutdown datanode
|
2021-01-22 01:36:40 +00:00
|
|
|
func (node *DataNode) Stop() error {
|
2021-12-08 06:19:04 +00:00
|
|
|
// https://github.com/milvus-io/milvus/issues/12282
|
|
|
|
node.UpdateStateCode(internalpb.StateCode_Abnormal)
|
|
|
|
|
2021-02-01 03:44:02 +00:00
|
|
|
node.cancel()
|
2021-01-19 03:37:16 +00:00
|
|
|
|
2021-05-28 08:47:29 +00:00
|
|
|
node.chanMut.RLock()
|
|
|
|
defer node.chanMut.RUnlock()
|
2021-01-19 03:37:16 +00:00
|
|
|
// close services
|
2021-05-25 07:35:37 +00:00
|
|
|
for _, syncService := range node.vchan2SyncService {
|
|
|
|
if syncService != nil {
|
|
|
|
(*syncService).close()
|
|
|
|
}
|
2021-01-19 03:37:16 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
if node.closer != nil {
|
2021-10-07 14:16:56 +00:00
|
|
|
err := node.closer.Close()
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
2021-01-19 03:37:16 +00:00
|
|
|
}
|
2021-11-16 14:31:14 +00:00
|
|
|
|
|
|
|
node.session.Revoke(time.Second)
|
2021-11-26 03:39:16 +00:00
|
|
|
|
2021-01-22 01:36:40 +00:00
|
|
|
return nil
|
2021-01-26 06:46:54 +00:00
|
|
|
}
|
|
|
|
|
2021-09-24 12:44:05 +00:00
|
|
|
// GetTimeTickChannel currently do nothing
|
2021-02-26 09:44:24 +00:00
|
|
|
func (node *DataNode) GetTimeTickChannel(ctx context.Context) (*milvuspb.StringResponse, error) {
|
|
|
|
return &milvuspb.StringResponse{
|
|
|
|
Status: &commonpb.Status{
|
2021-03-10 14:06:22 +00:00
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
2021-02-26 09:44:24 +00:00
|
|
|
Reason: "",
|
|
|
|
},
|
|
|
|
Value: "",
|
|
|
|
}, nil
|
2021-01-26 06:46:54 +00:00
|
|
|
}
|
2021-01-19 03:37:16 +00:00
|
|
|
|
2021-09-24 12:44:05 +00:00
|
|
|
// GetStatisticsChannel currently do nothing
|
2021-02-26 09:44:24 +00:00
|
|
|
func (node *DataNode) GetStatisticsChannel(ctx context.Context) (*milvuspb.StringResponse, error) {
|
|
|
|
return &milvuspb.StringResponse{
|
|
|
|
Status: &commonpb.Status{
|
2021-03-10 14:06:22 +00:00
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
2021-02-26 09:44:24 +00:00
|
|
|
Reason: "",
|
|
|
|
},
|
|
|
|
Value: "",
|
|
|
|
}, nil
|
2021-01-19 03:37:16 +00:00
|
|
|
}
|
2021-09-01 02:13:15 +00:00
|
|
|
|
2021-09-15 13:25:49 +00:00
|
|
|
// GetMetrics return datanode metrics
|
2021-09-01 02:13:15 +00:00
|
|
|
// TODO(dragondriver): cache the Metrics and set a retention to the cache
|
|
|
|
func (node *DataNode) GetMetrics(ctx context.Context, req *milvuspb.GetMetricsRequest) (*milvuspb.GetMetricsResponse, error) {
|
|
|
|
log.Debug("DataNode.GetMetrics",
|
2021-12-23 10:39:11 +00:00
|
|
|
zap.Int64("node_id", Params.DataNodeCfg.NodeID),
|
2021-09-01 02:13:15 +00:00
|
|
|
zap.String("req", req.Request))
|
|
|
|
|
|
|
|
if !node.isHealthy() {
|
|
|
|
log.Warn("DataNode.GetMetrics failed",
|
2021-12-23 10:39:11 +00:00
|
|
|
zap.Int64("node_id", Params.DataNodeCfg.NodeID),
|
2021-09-01 02:13:15 +00:00
|
|
|
zap.String("req", req.Request),
|
2021-12-23 10:39:11 +00:00
|
|
|
zap.Error(errDataNodeIsUnhealthy(Params.DataNodeCfg.NodeID)))
|
2021-09-01 02:13:15 +00:00
|
|
|
|
|
|
|
return &milvuspb.GetMetricsResponse{
|
|
|
|
Status: &commonpb.Status{
|
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
2021-12-23 10:39:11 +00:00
|
|
|
Reason: msgDataNodeIsUnhealthy(Params.DataNodeCfg.NodeID),
|
2021-09-01 02:13:15 +00:00
|
|
|
},
|
|
|
|
Response: "",
|
|
|
|
}, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
metricType, err := metricsinfo.ParseMetricType(req.Request)
|
|
|
|
if err != nil {
|
|
|
|
log.Warn("DataNode.GetMetrics failed to parse metric type",
|
2021-12-23 10:39:11 +00:00
|
|
|
zap.Int64("node_id", Params.DataNodeCfg.NodeID),
|
2021-09-01 02:13:15 +00:00
|
|
|
zap.String("req", req.Request),
|
|
|
|
zap.Error(err))
|
|
|
|
|
|
|
|
return &milvuspb.GetMetricsResponse{
|
|
|
|
Status: &commonpb.Status{
|
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
|
|
|
Reason: err.Error(),
|
|
|
|
},
|
|
|
|
Response: "",
|
|
|
|
}, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
log.Debug("DataNode.GetMetrics",
|
|
|
|
zap.String("metric_type", metricType))
|
|
|
|
|
|
|
|
if metricType == metricsinfo.SystemInfoMetrics {
|
|
|
|
systemInfoMetrics, err := node.getSystemInfoMetrics(ctx, req)
|
|
|
|
|
|
|
|
log.Debug("DataNode.GetMetrics",
|
2021-12-23 10:39:11 +00:00
|
|
|
zap.Int64("node_id", Params.DataNodeCfg.NodeID),
|
2021-09-01 02:13:15 +00:00
|
|
|
zap.String("req", req.Request),
|
|
|
|
zap.String("metric_type", metricType),
|
|
|
|
zap.Any("systemInfoMetrics", systemInfoMetrics), // TODO(dragondriver): necessary? may be very large
|
|
|
|
zap.Error(err))
|
|
|
|
|
2021-12-01 14:17:46 +00:00
|
|
|
return systemInfoMetrics, nil
|
2021-09-01 02:13:15 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
log.Debug("DataNode.GetMetrics failed, request metric type is not implemented yet",
|
2021-12-23 10:39:11 +00:00
|
|
|
zap.Int64("node_id", Params.DataNodeCfg.NodeID),
|
2021-09-01 02:13:15 +00:00
|
|
|
zap.String("req", req.Request),
|
|
|
|
zap.String("metric_type", metricType))
|
|
|
|
|
|
|
|
return &milvuspb.GetMetricsResponse{
|
|
|
|
Status: &commonpb.Status{
|
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
|
|
|
Reason: metricsinfo.MsgUnimplementedMetric,
|
|
|
|
},
|
|
|
|
Response: "",
|
|
|
|
}, nil
|
|
|
|
}
|
2021-11-05 14:25:00 +00:00
|
|
|
|
2021-12-28 01:38:07 +00:00
|
|
|
// Compaction handles compaction request from DataCoord
|
|
|
|
// returns status as long as compaction task enqueued or invalid
|
2021-11-05 14:25:00 +00:00
|
|
|
func (node *DataNode) Compaction(ctx context.Context, req *datapb.CompactionPlan) (*commonpb.Status, error) {
|
2021-11-08 11:49:07 +00:00
|
|
|
status := &commonpb.Status{
|
|
|
|
ErrorCode: commonpb.ErrorCode_UnexpectedError,
|
|
|
|
}
|
|
|
|
|
|
|
|
ds, ok := node.vchan2SyncService[req.GetChannel()]
|
|
|
|
if !ok {
|
|
|
|
log.Warn("illegel compaction plan, channel not in this DataNode", zap.String("channel name", req.GetChannel()))
|
|
|
|
status.Reason = errIllegalCompactionPlan.Error()
|
|
|
|
return status, nil
|
|
|
|
}
|
|
|
|
|
2021-12-02 08:39:33 +00:00
|
|
|
if !node.compactionExecutor.channelValidateForCompaction(req.GetChannel()) {
|
|
|
|
log.Warn("channel of compaction is marked invalid in compaction executor", zap.String("channel name", req.GetChannel()))
|
|
|
|
status.Reason = "channel marked invalid"
|
|
|
|
return status, nil
|
|
|
|
}
|
|
|
|
|
2021-11-08 11:49:07 +00:00
|
|
|
binlogIO := &binlogIO{node.blobKv, ds.idAllocator}
|
|
|
|
task := newCompactionTask(
|
2021-11-17 02:07:13 +00:00
|
|
|
node.ctx,
|
2021-11-08 11:49:07 +00:00
|
|
|
binlogIO, binlogIO,
|
|
|
|
ds.replica,
|
|
|
|
ds.flushManager,
|
|
|
|
ds.idAllocator,
|
|
|
|
node.dataCoord,
|
|
|
|
req,
|
|
|
|
)
|
|
|
|
|
|
|
|
node.compactionExecutor.execute(task)
|
|
|
|
|
|
|
|
return &commonpb.Status{
|
|
|
|
ErrorCode: commonpb.ErrorCode_Success,
|
|
|
|
}, nil
|
2021-11-05 14:25:00 +00:00
|
|
|
}
|