milvus/internal/querynode/flow_graph_key2seg_node.go
bigsheeper f02bd8c8f5 Rename query node package, implement watchDmChannel
Signed-off-by: bigsheeper <yihao.dai@zilliz.com>
2021-01-16 10:12:14 +08:00

53 lines
1.4 KiB
Go

package querynode
type key2SegNode struct {
baseNode
key2SegMsg key2SegMsg
}
func (ksNode *key2SegNode) Name() string {
return "ksNode"
}
func (ksNode *key2SegNode) Operate(in []*Msg) []*Msg {
return in
}
func newKey2SegNode() *key2SegNode {
maxQueueLength := Params.FlowGraphMaxQueueLength
maxParallelism := Params.FlowGraphMaxParallelism
baseNode := baseNode{}
baseNode.SetMaxQueueLength(maxQueueLength)
baseNode.SetMaxParallelism(maxParallelism)
return &key2SegNode{
baseNode: baseNode,
}
}
/************************************** util functions ***************************************/
// Function `GetSegmentByEntityId` should return entityIDs, timestamps and segmentIDs
//func (node *QueryNode) GetKey2Segments() (*[]int64, *[]uint64, *[]int64) {
// var entityIDs = make([]int64, 0)
// var timestamps = make([]uint64, 0)
// var segmentIDs = make([]int64, 0)
//
// var key2SegMsg = node.messageClient.Key2SegMsg
// for _, msg := range key2SegMsg {
// if msg.SegmentID == nil {
// segmentIDs = append(segmentIDs, -1)
// entityIDs = append(entityIDs, msg.Uid)
// timestamps = append(timestamps, msg.Timestamp)
// } else {
// for _, segmentID := range msg.SegmentID {
// segmentIDs = append(segmentIDs, segmentID)
// entityIDs = append(entityIDs, msg.Uid)
// timestamps = append(timestamps, msg.Timestamp)
// }
// }
// }
//
// return &entityIDs, &timestamps, &segmentIDs
//}