Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
milvus
milvus
提交
e68374b6
M
milvus
项目概览
milvus
/
milvus
大约 1 年 前同步成功
通知
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,发现更多精彩内容 >>
未验证
提交
e68374b6
编写于
1月 10, 2023
作者:
L
liliu-z
提交者:
GitHub
1月 10, 2023
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
Skip search GRPC call for standalon (#21630)
Signed-off-by:
N
Li Liu
<
li.liu@zilliz.com
>
上级
d78b17f4
变更
4
隐藏空白更改
内联
并排
Showing
4 changed file
with
39 addition
and
9 deletion
+39
-9
internal/proxy/task_search.go
internal/proxy/task_search.go
+11
-1
internal/querynode/impl.go
internal/querynode/impl.go
+2
-2
internal/querynode/query_node.go
internal/querynode/query_node.go
+15
-5
internal/querynode/shard_cluster.go
internal/querynode/shard_cluster.go
+11
-1
未找到文件。
internal/proxy/task_search.go
浏览文件 @
e68374b6
...
...
@@ -10,6 +10,7 @@ import (
"github.com/milvus-io/milvus/internal/common"
"github.com/milvus-io/milvus/internal/parser/planparserv2"
"github.com/milvus-io/milvus/internal/querynode"
"github.com/golang/protobuf/proto"
"go.uber.org/zap"
...
...
@@ -485,7 +486,16 @@ func (t *searchTask) searchShard(ctx context.Context, nodeID int64, qn types.Que
DmlChannels
:
channelIDs
,
Scope
:
querypb
.
DataScope_All
,
}
result
,
err
:=
qn
.
Search
(
ctx
,
req
)
queryNode
:=
querynode
.
GetQueryNode
()
var
result
*
internalpb
.
SearchResults
var
err
error
if
queryNode
!=
nil
&&
queryNode
.
IsStandAlone
{
result
,
err
=
queryNode
.
Search
(
ctx
,
req
)
}
else
{
result
,
err
=
qn
.
Search
(
ctx
,
req
)
}
if
err
!=
nil
{
log
.
Ctx
(
ctx
)
.
Warn
(
"QueryNode search return error"
,
zap
.
Int64
(
"msgID"
,
t
.
ID
()),
zap
.
Int64
(
"nodeID"
,
nodeID
),
zap
.
Strings
(
"channels"
,
channelIDs
),
zap
.
Error
(
err
))
...
...
internal/querynode/impl.go
浏览文件 @
e68374b6
...
...
@@ -729,7 +729,7 @@ func filterSegmentInfo(segmentInfos []*querypb.SegmentInfo, segmentIDs map[int64
// Search performs replica search tasks.
func
(
node
*
QueryNode
)
Search
(
ctx
context
.
Context
,
req
*
querypb
.
SearchRequest
)
(
*
internalpb
.
SearchResults
,
error
)
{
if
req
.
GetReq
()
.
GetBase
()
.
GetTargetID
()
!=
node
.
session
.
ServerID
{
if
!
node
.
IsStandAlone
&&
req
.
GetReq
()
.
GetBase
()
.
GetTargetID
()
!=
node
.
session
.
ServerID
{
return
&
internalpb
.
SearchResults
{
Status
:
&
commonpb
.
Status
{
ErrorCode
:
commonpb
.
ErrorCode_NodeIDNotMatch
,
...
...
@@ -1204,7 +1204,7 @@ func (node *QueryNode) SyncReplicaSegments(ctx context.Context, req *querypb.Syn
return
&
commonpb
.
Status
{
ErrorCode
:
commonpb
.
ErrorCode_Success
},
nil
}
//ShowConfigurations returns the configurations of queryNode matching req.Pattern
//
ShowConfigurations returns the configurations of queryNode matching req.Pattern
func
(
node
*
QueryNode
)
ShowConfigurations
(
ctx
context
.
Context
,
req
*
internalpb
.
ShowConfigurationsRequest
)
(
*
internalpb
.
ShowConfigurationsResponse
,
error
)
{
if
!
node
.
isHealthyOrStopping
()
{
log
.
Warn
(
"QueryNode.ShowConfigurations failed"
,
...
...
internal/querynode/query_node.go
浏览文件 @
e68374b6
...
...
@@ -123,22 +123,32 @@ type QueryNode struct {
// pool for load/release channel
taskPool
*
concurrency
.
Pool
IsStandAlone
bool
}
var
queryNode
*
QueryNode
=
nil
func
GetQueryNode
()
*
QueryNode
{
return
queryNode
}
// NewQueryNode will return a QueryNode with abnormal state.
func
NewQueryNode
(
ctx
context
.
Context
,
factory
dependency
.
Factory
)
*
QueryNode
{
ctx1
,
cancel
:=
context
.
WithCancel
(
ctx
)
node
:=
&
QueryNode
{
queryNode
=
&
QueryNode
{
queryNodeLoopCtx
:
ctx1
,
queryNodeLoopCancel
:
cancel
,
factory
:
factory
,
IsStandAlone
:
os
.
Getenv
(
metricsinfo
.
DeployModeEnvKey
)
==
metricsinfo
.
StandaloneDeployMode
,
}
n
ode
.
tSafeReplica
=
newTSafeReplica
()
node
.
scheduler
=
newTaskScheduler
(
ctx1
,
n
ode
.
tSafeReplica
)
n
ode
.
UpdateStateCode
(
commonpb
.
StateCode_Abnormal
)
queryN
ode
.
tSafeReplica
=
newTSafeReplica
()
queryNode
.
scheduler
=
newTaskScheduler
(
ctx1
,
queryN
ode
.
tSafeReplica
)
queryN
ode
.
UpdateStateCode
(
commonpb
.
StateCode_Abnormal
)
return
n
ode
return
queryN
ode
}
func
(
node
*
QueryNode
)
initSession
()
error
{
...
...
internal/querynode/shard_cluster.go
浏览文件 @
e68374b6
...
...
@@ -979,7 +979,17 @@ func (sc *ShardCluster) Search(ctx context.Context, req *querypb.SearchRequest,
wg
.
Add
(
1
)
go
func
()
{
defer
wg
.
Done
()
partialResult
,
nodeErr
:=
node
.
client
.
Search
(
reqCtx
,
nodeReq
)
queryNode
:=
GetQueryNode
()
var
partialResult
*
internalpb
.
SearchResults
var
nodeErr
error
if
queryNode
!=
nil
&&
queryNode
.
IsStandAlone
{
partialResult
,
nodeErr
=
queryNode
.
Search
(
reqCtx
,
nodeReq
)
}
else
{
partialResult
,
nodeErr
=
node
.
client
.
Search
(
reqCtx
,
nodeReq
)
}
resultMut
.
Lock
()
defer
resultMut
.
Unlock
()
if
nodeErr
!=
nil
||
partialResult
.
GetStatus
()
.
GetErrorCode
()
!=
commonpb
.
ErrorCode_Success
{
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录