Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
milvus
milvus
提交
028c5cb8
M
milvus
项目概览
milvus
/
milvus
11 个月 前同步成功
通知
261
Star
22476
Fork
2472
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
0
列表
看板
标记
里程碑
合并请求
0
DevOps
流水线
流水线任务
计划
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
M
milvus
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
0
Issue
0
列表
看板
标记
里程碑
合并请求
0
合并请求
0
Pages
DevOps
DevOps
流水线
流水线任务
计划
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
流水线任务
提交
Issue看板
体验新版 GitCode,发现更多精彩内容 >>
未验证
提交
028c5cb8
编写于
4月 02, 2022
作者:
L
Letian Jiang
提交者:
GitHub
4月 02, 2022
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
Modify grpc interface for replica Search/Query in QueryNode (#16326)
Signed-off-by:
N
Letian Jiang
<
letian.jiang@zilliz.com
>
上级
7f7379d5
变更
9
展开全部
隐藏空白更改
内联
并排
Showing
9 changed file
with
242 addition
and
237 deletion
+242
-237
internal/distributed/querynode/client/client.go
internal/distributed/querynode/client/client.go
+6
-5
internal/distributed/querynode/service.go
internal/distributed/querynode/service.go
+7
-6
internal/distributed/querynode/service_test.go
internal/distributed/querynode/service_test.go
+7
-6
internal/proto/query_coord.proto
internal/proto/query_coord.proto
+4
-4
internal/proto/querypb/query_coord.pb.go
internal/proto/querypb/query_coord.pb.go
+205
-205
internal/querycoord/mock_querynode_client_test.go
internal/querycoord/mock_querynode_client_test.go
+2
-2
internal/querynode/impl.go
internal/querynode/impl.go
+2
-2
internal/types/types.go
internal/types/types.go
+4
-3
internal/util/mock/querynode_client.go
internal/util/mock/querynode_client.go
+5
-4
未找到文件。
internal/distributed/querynode/client/client.go
浏览文件 @
028c5cb8
...
...
@@ -20,6 +20,8 @@ import (
"context"
"fmt"
"google.golang.org/grpc"
"github.com/milvus-io/milvus/internal/proto/commonpb"
"github.com/milvus-io/milvus/internal/proto/internalpb"
"github.com/milvus-io/milvus/internal/proto/milvuspb"
...
...
@@ -28,7 +30,6 @@ import (
"github.com/milvus-io/milvus/internal/util/grpcclient"
"github.com/milvus-io/milvus/internal/util/paramtable"
"github.com/milvus-io/milvus/internal/util/typeutil"
"google.golang.org/grpc"
)
var
ClientParams
paramtable
.
GrpcClientConfig
...
...
@@ -242,7 +243,7 @@ func (c *Client) ReleaseSegments(ctx context.Context, req *querypb.ReleaseSegmen
}
// Search performs replica search tasks in QueryNode.
func
(
c
*
Client
)
Search
(
ctx
context
.
Context
,
req
*
querypb
.
SearchRequest
)
(
*
milvus
pb
.
SearchResults
,
error
)
{
func
(
c
*
Client
)
Search
(
ctx
context
.
Context
,
req
*
querypb
.
SearchRequest
)
(
*
internal
pb
.
SearchResults
,
error
)
{
ret
,
err
:=
c
.
grpcClient
.
ReCall
(
ctx
,
func
(
client
interface
{})
(
interface
{},
error
)
{
if
!
funcutil
.
CheckCtxValid
(
ctx
)
{
return
nil
,
ctx
.
Err
()
...
...
@@ -252,11 +253,11 @@ func (c *Client) Search(ctx context.Context, req *querypb.SearchRequest) (*milvu
if
err
!=
nil
||
ret
==
nil
{
return
nil
,
err
}
return
ret
.
(
*
milvus
pb
.
SearchResults
),
err
return
ret
.
(
*
internal
pb
.
SearchResults
),
err
}
// Query performs replica query tasks in QueryNode.
func
(
c
*
Client
)
Query
(
ctx
context
.
Context
,
req
*
querypb
.
QueryRequest
)
(
*
milvuspb
.
Query
Results
,
error
)
{
func
(
c
*
Client
)
Query
(
ctx
context
.
Context
,
req
*
querypb
.
QueryRequest
)
(
*
internalpb
.
Retrieve
Results
,
error
)
{
ret
,
err
:=
c
.
grpcClient
.
ReCall
(
ctx
,
func
(
client
interface
{})
(
interface
{},
error
)
{
if
!
funcutil
.
CheckCtxValid
(
ctx
)
{
return
nil
,
ctx
.
Err
()
...
...
@@ -266,7 +267,7 @@ func (c *Client) Query(ctx context.Context, req *querypb.QueryRequest) (*milvusp
if
err
!=
nil
||
ret
==
nil
{
return
nil
,
err
}
return
ret
.
(
*
milvuspb
.
Query
Results
),
err
return
ret
.
(
*
internalpb
.
Retrieve
Results
),
err
}
// GetSegmentInfo gets the information of the specified segments in QueryNode.
...
...
internal/distributed/querynode/service.go
浏览文件 @
028c5cb8
...
...
@@ -26,6 +26,11 @@ import (
"time"
ot
"github.com/grpc-ecosystem/go-grpc-middleware/tracing/opentracing"
clientv3
"go.etcd.io/etcd/client/v3"
"go.uber.org/zap"
"google.golang.org/grpc"
"google.golang.org/grpc/keepalive"
"github.com/milvus-io/milvus/internal/log"
"github.com/milvus-io/milvus/internal/mq/msgstream"
"github.com/milvus-io/milvus/internal/proto/commonpb"
...
...
@@ -40,10 +45,6 @@ import (
"github.com/milvus-io/milvus/internal/util/retry"
"github.com/milvus-io/milvus/internal/util/trace"
"github.com/milvus-io/milvus/internal/util/typeutil"
clientv3
"go.etcd.io/etcd/client/v3"
"go.uber.org/zap"
"google.golang.org/grpc"
"google.golang.org/grpc/keepalive"
)
var
Params
paramtable
.
GrpcServerConfig
...
...
@@ -305,12 +306,12 @@ func (s *Server) GetSegmentInfo(ctx context.Context, req *querypb.GetSegmentInfo
}
// Search performs search of streaming/historical replica on QueryNode.
func
(
s
*
Server
)
Search
(
ctx
context
.
Context
,
req
*
querypb
.
SearchRequest
)
(
*
milvus
pb
.
SearchResults
,
error
)
{
func
(
s
*
Server
)
Search
(
ctx
context
.
Context
,
req
*
querypb
.
SearchRequest
)
(
*
internal
pb
.
SearchResults
,
error
)
{
return
s
.
querynode
.
Search
(
ctx
,
req
)
}
// Query performs query of streaming/historical replica on QueryNode.
func
(
s
*
Server
)
Query
(
ctx
context
.
Context
,
req
*
querypb
.
QueryRequest
)
(
*
milvuspb
.
Query
Results
,
error
)
{
func
(
s
*
Server
)
Query
(
ctx
context
.
Context
,
req
*
querypb
.
QueryRequest
)
(
*
internalpb
.
Retrieve
Results
,
error
)
{
return
s
.
querynode
.
Query
(
ctx
,
req
)
}
...
...
internal/distributed/querynode/service_test.go
浏览文件 @
028c5cb8
...
...
@@ -23,12 +23,13 @@ import (
"github.com/milvus-io/milvus/internal/types"
"github.com/stretchr/testify/assert"
clientv3
"go.etcd.io/etcd/client/v3"
"github.com/milvus-io/milvus/internal/proto/commonpb"
"github.com/milvus-io/milvus/internal/proto/internalpb"
"github.com/milvus-io/milvus/internal/proto/milvuspb"
"github.com/milvus-io/milvus/internal/proto/querypb"
"github.com/stretchr/testify/assert"
clientv3
"go.etcd.io/etcd/client/v3"
)
///////////////////////////////////////////////////////////////////////////////////////////////////////////////////////
...
...
@@ -43,8 +44,8 @@ type MockQueryNode struct {
strResp
*
milvuspb
.
StringResponse
infoResp
*
querypb
.
GetSegmentInfoResponse
metricResp
*
milvuspb
.
GetMetricsResponse
searchResp
*
milvus
pb
.
SearchResults
queryResp
*
milvuspb
.
Query
Results
searchResp
*
internal
pb
.
SearchResults
queryResp
*
internalpb
.
Retrieve
Results
}
func
(
m
*
MockQueryNode
)
Init
()
error
{
...
...
@@ -115,11 +116,11 @@ func (m *MockQueryNode) GetMetrics(ctx context.Context, req *milvuspb.GetMetrics
return
m
.
metricResp
,
m
.
err
}
func
(
m
*
MockQueryNode
)
Search
(
ctx
context
.
Context
,
req
*
querypb
.
SearchRequest
)
(
*
milvus
pb
.
SearchResults
,
error
)
{
func
(
m
*
MockQueryNode
)
Search
(
ctx
context
.
Context
,
req
*
querypb
.
SearchRequest
)
(
*
internal
pb
.
SearchResults
,
error
)
{
return
m
.
searchResp
,
m
.
err
}
func
(
m
*
MockQueryNode
)
Query
(
ctx
context
.
Context
,
req
*
querypb
.
QueryRequest
)
(
*
milvuspb
.
Query
Results
,
error
)
{
func
(
m
*
MockQueryNode
)
Query
(
ctx
context
.
Context
,
req
*
querypb
.
QueryRequest
)
(
*
internalpb
.
Retrieve
Results
,
error
)
{
return
m
.
queryResp
,
m
.
err
}
...
...
internal/proto/query_coord.proto
浏览文件 @
028c5cb8
...
...
@@ -51,8 +51,8 @@ service QueryNode {
rpc
ReleaseSegments
(
ReleaseSegmentsRequest
)
returns
(
common.Status
)
{}
rpc
GetSegmentInfo
(
GetSegmentInfoRequest
)
returns
(
GetSegmentInfoResponse
)
{}
rpc
Search
(
SearchRequest
)
returns
(
milvus
.SearchResults
)
{}
rpc
Query
(
QueryRequest
)
returns
(
milvus.Query
Results
)
{}
rpc
Search
(
SearchRequest
)
returns
(
internal
.SearchResults
)
{}
rpc
Query
(
QueryRequest
)
returns
(
internal.Retrieve
Results
)
{}
// https://wiki.lfaidata.foundation/display/MIL/MEP+8+--+Add+metrics+for+proxy
rpc
GetMetrics
(
milvus.GetMetricsRequest
)
returns
(
milvus.GetMetricsResponse
)
{}
...
...
@@ -269,13 +269,13 @@ message ReleaseSegmentsRequest {
}
message
SearchRequest
{
milvus
.SearchRequest
req
=
1
;
internal
.SearchRequest
req
=
1
;
string
dml_channel
=
2
;
repeated
int64
segmentIDs
=
3
;
}
message
QueryRequest
{
milvus.Query
Request
req
=
1
;
internal.Retrieve
Request
req
=
1
;
string
dml_channel
=
2
;
repeated
int64
segmentIDs
=
3
;
}
...
...
internal/proto/querypb/query_coord.pb.go
浏览文件 @
028c5cb8
此差异已折叠。
点击以展开。
internal/querycoord/mock_querynode_client_test.go
浏览文件 @
028c5cb8
...
...
@@ -158,10 +158,10 @@ func (client *queryNodeClientMock) GetMetrics(ctx context.Context, req *milvuspb
return
client
.
grpcClient
.
GetMetrics
(
ctx
,
req
)
}
func
(
client
*
queryNodeClientMock
)
Search
(
ctx
context
.
Context
,
req
*
querypb
.
SearchRequest
)
(
*
milvus
pb
.
SearchResults
,
error
)
{
func
(
client
*
queryNodeClientMock
)
Search
(
ctx
context
.
Context
,
req
*
querypb
.
SearchRequest
)
(
*
internal
pb
.
SearchResults
,
error
)
{
return
client
.
grpcClient
.
Search
(
ctx
,
req
)
}
func
(
client
*
queryNodeClientMock
)
Query
(
ctx
context
.
Context
,
req
*
querypb
.
QueryRequest
)
(
*
milvuspb
.
Query
Results
,
error
)
{
func
(
client
*
queryNodeClientMock
)
Query
(
ctx
context
.
Context
,
req
*
querypb
.
QueryRequest
)
(
*
internalpb
.
Retrieve
Results
,
error
)
{
return
client
.
grpcClient
.
Query
(
ctx
,
req
)
}
internal/querynode/impl.go
浏览文件 @
028c5cb8
...
...
@@ -556,12 +556,12 @@ func (node *QueryNode) isHealthy() bool {
}
// Search performs replica search tasks.
func
(
node
*
QueryNode
)
Search
(
ctx
context
.
Context
,
req
*
queryPb
.
SearchRequest
)
(
*
milvus
pb
.
SearchResults
,
error
)
{
func
(
node
*
QueryNode
)
Search
(
ctx
context
.
Context
,
req
*
queryPb
.
SearchRequest
)
(
*
internal
pb
.
SearchResults
,
error
)
{
return
nil
,
errors
.
New
(
"not implemented"
)
}
// Query performs replica query tasks.
func
(
node
*
QueryNode
)
Query
(
ctx
context
.
Context
,
req
*
queryPb
.
QueryRequest
)
(
*
milvuspb
.
Query
Results
,
error
)
{
func
(
node
*
QueryNode
)
Query
(
ctx
context
.
Context
,
req
*
queryPb
.
QueryRequest
)
(
*
internalpb
.
Retrieve
Results
,
error
)
{
return
nil
,
errors
.
New
(
"not implemented"
)
}
...
...
internal/types/types.go
浏览文件 @
028c5cb8
...
...
@@ -19,6 +19,8 @@ package types
import
(
"context"
clientv3
"go.etcd.io/etcd/client/v3"
"github.com/milvus-io/milvus/internal/proto/commonpb"
"github.com/milvus-io/milvus/internal/proto/datapb"
"github.com/milvus-io/milvus/internal/proto/indexpb"
...
...
@@ -29,7 +31,6 @@ import (
"github.com/milvus-io/milvus/internal/proto/rootcoordpb"
"github.com/milvus-io/milvus/internal/util/sessionutil"
"github.com/milvus-io/milvus/internal/util/typeutil"
clientv3
"go.etcd.io/etcd/client/v3"
)
// TimeTickProvider is the interface all services implement
...
...
@@ -1132,8 +1133,8 @@ type QueryNode interface {
ReleaseSegments
(
ctx
context
.
Context
,
req
*
querypb
.
ReleaseSegmentsRequest
)
(
*
commonpb
.
Status
,
error
)
GetSegmentInfo
(
ctx
context
.
Context
,
req
*
querypb
.
GetSegmentInfoRequest
)
(
*
querypb
.
GetSegmentInfoResponse
,
error
)
Search
(
ctx
context
.
Context
,
req
*
querypb
.
SearchRequest
)
(
*
milvus
pb
.
SearchResults
,
error
)
Query
(
ctx
context
.
Context
,
req
*
querypb
.
QueryRequest
)
(
*
milvuspb
.
Query
Results
,
error
)
Search
(
ctx
context
.
Context
,
req
*
querypb
.
SearchRequest
)
(
*
internal
pb
.
SearchResults
,
error
)
Query
(
ctx
context
.
Context
,
req
*
querypb
.
QueryRequest
)
(
*
internalpb
.
Retrieve
Results
,
error
)
// GetMetrics gets the metrics about QueryNode.
GetMetrics
(
ctx
context
.
Context
,
req
*
milvuspb
.
GetMetricsRequest
)
(
*
milvuspb
.
GetMetricsResponse
,
error
)
...
...
internal/util/mock/querynode_client.go
浏览文件 @
028c5cb8
...
...
@@ -78,13 +78,14 @@ func (m *QueryNodeClient) GetSegmentInfo(ctx context.Context, in *querypb.GetSeg
return
&
querypb
.
GetSegmentInfoResponse
{},
m
.
Err
}
func
(
m
*
QueryNodeClient
)
Search
(
ctx
context
.
Context
,
in
*
querypb
.
SearchRequest
,
opts
...
grpc
.
CallOption
)
(
*
milvus
pb
.
SearchResults
,
error
)
{
return
&
milvus
pb
.
SearchResults
{},
m
.
Err
func
(
m
*
QueryNodeClient
)
Search
(
ctx
context
.
Context
,
in
*
querypb
.
SearchRequest
,
opts
...
grpc
.
CallOption
)
(
*
internal
pb
.
SearchResults
,
error
)
{
return
&
internal
pb
.
SearchResults
{},
m
.
Err
}
func
(
m
*
QueryNodeClient
)
Query
(
ctx
context
.
Context
,
in
*
querypb
.
QueryRequest
,
opts
...
grpc
.
CallOption
)
(
*
milvuspb
.
Query
Results
,
error
)
{
return
&
milvuspb
.
Query
Results
{},
m
.
Err
func
(
m
*
QueryNodeClient
)
Query
(
ctx
context
.
Context
,
in
*
querypb
.
QueryRequest
,
opts
...
grpc
.
CallOption
)
(
*
internalpb
.
Retrieve
Results
,
error
)
{
return
&
internalpb
.
Retrieve
Results
{},
m
.
Err
}
func
(
m
*
QueryNodeClient
)
GetMetrics
(
ctx
context
.
Context
,
in
*
milvuspb
.
GetMetricsRequest
,
opts
...
grpc
.
CallOption
)
(
*
milvuspb
.
GetMetricsResponse
,
error
)
{
return
&
milvuspb
.
GetMetricsResponse
{},
m
.
Err
}
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录