Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
milvus
milvus
提交
387ad98b
M
milvus
项目概览
milvus
/
milvus
9 个月 前同步成功
通知
260
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,体验更适合开发者的 AI 搜索 >>
未验证
提交
387ad98b
编写于
8月 20, 2023
作者:
W
wayblink
提交者:
GitHub
8月 20, 2023
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
Check and reset if grpc client serverID mismatch with session (#26473)
Signed-off-by:
N
wayblink
<
anyang.wang@zilliz.com
>
上级
2c32d0db
变更
6
隐藏空白更改
内联
并排
Showing
6 changed file
with
40 addition
and
5 deletion
+40
-5
internal/distributed/datacoord/client/client.go
internal/distributed/datacoord/client/client.go
+1
-0
internal/distributed/indexcoord/client/client.go
internal/distributed/indexcoord/client/client.go
+1
-0
internal/distributed/querycoord/client/client.go
internal/distributed/querycoord/client/client.go
+1
-0
internal/distributed/rootcoord/client/client.go
internal/distributed/rootcoord/client/client.go
+1
-0
internal/util/grpcclient/client.go
internal/util/grpcclient/client.go
+30
-5
internal/util/mock/grpcclient.go
internal/util/mock/grpcclient.go
+6
-0
未找到文件。
internal/distributed/datacoord/client/client.go
浏览文件 @
387ad98b
...
...
@@ -78,6 +78,7 @@ func NewClient(ctx context.Context, metaRoot string, etcdCli *clientv3.Client) (
client
.
grpcClient
.
SetRole
(
typeutil
.
DataCoordRole
)
client
.
grpcClient
.
SetGetAddrFunc
(
client
.
getDataCoordAddr
)
client
.
grpcClient
.
SetNewGrpcClientFunc
(
client
.
newGrpcClient
)
client
.
grpcClient
.
SetSession
(
sess
)
return
client
,
nil
}
...
...
internal/distributed/indexcoord/client/client.go
浏览文件 @
387ad98b
...
...
@@ -74,6 +74,7 @@ func NewClient(ctx context.Context, metaRoot string, etcdCli *clientv3.Client) (
client
.
grpcClient
.
SetRole
(
typeutil
.
IndexCoordRole
)
client
.
grpcClient
.
SetGetAddrFunc
(
client
.
getIndexCoordAddr
)
client
.
grpcClient
.
SetNewGrpcClientFunc
(
client
.
newGrpcClient
)
client
.
grpcClient
.
SetSession
(
sess
)
return
client
,
nil
}
...
...
internal/distributed/querycoord/client/client.go
浏览文件 @
387ad98b
...
...
@@ -73,6 +73,7 @@ func NewClient(ctx context.Context, metaRoot string, etcdCli *clientv3.Client) (
client
.
grpcClient
.
SetRole
(
typeutil
.
QueryCoordRole
)
client
.
grpcClient
.
SetGetAddrFunc
(
client
.
getQueryCoordAddr
)
client
.
grpcClient
.
SetNewGrpcClientFunc
(
client
.
newGrpcClient
)
client
.
grpcClient
.
SetSession
(
sess
)
return
client
,
nil
}
...
...
internal/distributed/rootcoord/client/client.go
浏览文件 @
387ad98b
...
...
@@ -81,6 +81,7 @@ func NewClient(ctx context.Context, metaRoot string, etcdCli *clientv3.Client) (
client
.
grpcClient
.
SetRole
(
typeutil
.
RootCoordRole
)
client
.
grpcClient
.
SetGetAddrFunc
(
client
.
getRootCoordAddr
)
client
.
grpcClient
.
SetNewGrpcClientFunc
(
client
.
newGrpcClient
)
client
.
grpcClient
.
SetSession
(
sess
)
return
client
,
nil
}
...
...
internal/util/grpcclient/client.go
浏览文件 @
387ad98b
...
...
@@ -42,6 +42,7 @@ import (
"github.com/milvus-io/milvus/internal/util/generic"
"github.com/milvus-io/milvus/internal/util/interceptor"
"github.com/milvus-io/milvus/internal/util/retry"
"github.com/milvus-io/milvus/internal/util/sessionutil"
"github.com/milvus-io/milvus/internal/util/trace"
)
...
...
@@ -60,6 +61,7 @@ type GrpcClient[T interface {
Close
()
error
SetNodeID
(
int64
)
GetNodeID
()
int64
SetSession
(
sess
*
sessionutil
.
Session
)
}
// ClientBase is a base of grpc client
...
...
@@ -87,6 +89,7 @@ type ClientBase[T interface {
MaxBackoff
float32
BackoffMultiplier
float32
NodeID
int64
sess
*
sessionutil
.
Session
sf
singleflight
.
Group
}
...
...
@@ -291,11 +294,6 @@ func (c *ClientBase[T]) callOnce(ctx context.Context, caller func(client T) (any
return
ret
,
nil
}
if
!
funcutil
.
CheckCtxValid
(
ctx
)
{
// start bg check in case of https://github.com/milvus-io/milvus/issues/22435
go
c
.
bgHealthCheck
(
client
)
return
generic
.
Zero
[
T
](),
err
}
if
IsCrossClusterRoutingErr
(
err
)
{
log
.
Warn
(
"CrossClusterRoutingErr, start to reset connection"
,
zap
.
Error
(
err
))
c
.
resetConnection
(
client
)
...
...
@@ -306,6 +304,28 @@ func (c *ClientBase[T]) callOnce(ctx context.Context, caller func(client T) (any
c
.
resetConnection
(
client
)
return
ret
,
err
}
if
!
funcutil
.
CheckCtxValid
(
ctx
)
{
// check if server ID matches coord session, if not, reset connection
if
c
.
sess
!=
nil
{
sessions
,
_
,
getSessionErr
:=
c
.
sess
.
GetSessions
(
c
.
GetRole
())
if
getSessionErr
!=
nil
{
// Only log but not handle this error as it is an auxiliary logic
log
.
Warn
(
"Fail to GetSessions"
,
zap
.
Error
(
getSessionErr
))
}
if
coordSess
,
exist
:=
sessions
[
c
.
GetRole
()];
exist
{
if
c
.
GetNodeID
()
!=
coordSess
.
ServerID
{
log
.
Warn
(
"Server ID mismatch, may connected to a old server, start to reset connection"
,
zap
.
Error
(
err
))
c
.
resetConnection
(
client
)
return
ret
,
err
}
}
}
// start bg check in case of https://github.com/milvus-io/milvus/issues/22435
go
c
.
bgHealthCheck
(
client
)
return
generic
.
Zero
[
T
](),
err
}
if
!
funcutil
.
IsGrpcErr
(
err
)
{
log
.
Warn
(
"ClientBase:isNotGrpcErr"
,
zap
.
Error
(
err
))
return
generic
.
Zero
[
T
](),
err
...
...
@@ -398,6 +418,11 @@ func (c *ClientBase[T]) GetNodeID() int64 {
return
c
.
NodeID
}
// SetSession set session role of client
func
(
c
*
ClientBase
[
T
])
SetSession
(
sess
*
sessionutil
.
Session
)
{
c
.
sess
=
sess
}
func
IsCrossClusterRoutingErr
(
err
error
)
bool
{
// GRPC utilizes `status.Status` to encapsulate errors,
// hence it is not viable to employ the `errors.Is` for assessment.
...
...
internal/util/mock/grpcclient.go
浏览文件 @
387ad98b
...
...
@@ -30,6 +30,7 @@ import (
"github.com/milvus-io/milvus/internal/util/funcutil"
"github.com/milvus-io/milvus/internal/util/generic"
"github.com/milvus-io/milvus/internal/util/retry"
"github.com/milvus-io/milvus/internal/util/sessionutil"
"github.com/milvus-io/milvus/internal/util/trace"
)
...
...
@@ -43,6 +44,7 @@ type GRPCClientBase[T any] struct {
GetGrpcClientErr
error
role
string
nodeID
int64
sess
*
sessionutil
.
Session
}
func
(
c
*
GRPCClientBase
[
T
])
SetGetAddrFunc
(
f
func
()
(
string
,
error
))
{
...
...
@@ -169,3 +171,7 @@ func (c *GRPCClientBase[T]) SetNodeID(nodeID int64) {
func
SuccessStatus
()
*
commonpb
.
Status
{
return
&
commonpb
.
Status
{
ErrorCode
:
commonpb
.
ErrorCode_Success
}
}
func
(
c
*
GRPCClientBase
[
T
])
SetSession
(
sess
*
sessionutil
.
Session
)
{
c
.
sess
=
sess
}
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录