Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
milvus
milvus
提交
be407742
M
milvus
项目概览
milvus
/
milvus
10 个月 前同步成功
通知
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 搜索 >>
未验证
提交
be407742
编写于
6月 16, 2022
作者:
L
Letian Jiang
提交者:
GitHub
6月 16, 2022
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
Filter duplicated sealed segments which are both online and offline (#17593)
Signed-off-by:
N
Letian Jiang
<
letian.jiang@zilliz.com
>
上级
1f6fbf91
变更
2
隐藏空白更改
内联
并排
Showing
2 changed file
with
110 addition
and
1 deletion
+110
-1
internal/querynode/query_node.go
internal/querynode/query_node.go
+42
-1
internal/querynode/query_node_test.go
internal/querynode/query_node_test.go
+68
-0
未找到文件。
internal/querynode/query_node.go
浏览文件 @
be407742
...
...
@@ -50,6 +50,8 @@ import (
clientv3
"go.etcd.io/etcd/client/v3"
"go.uber.org/zap"
v3rpc
"go.etcd.io/etcd/api/v3/v3rpc/rpctypes"
etcdkv
"github.com/milvus-io/milvus/internal/kv/etcd"
"github.com/milvus-io/milvus/internal/log"
"github.com/milvus-io/milvus/internal/proto/internalpb"
...
...
@@ -61,7 +63,6 @@ import (
"github.com/milvus-io/milvus/internal/util/paramtable"
"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"
)
// make sure QueryNode implements types.QueryNode
...
...
@@ -376,6 +377,8 @@ func (node *QueryNode) handleSealedSegmentsChangeInfo(info *querypb.SealedSegmen
log
.
Warn
(
"failed to validate vchannel for SegmentChangeInfo"
,
zap
.
Error
(
err
))
continue
}
// ignore segments that are online and offline in the same QueryNode
filterDuplicateChangeInfo
(
line
)
node
.
ShardClusterService
.
HandoffVChannelSegments
(
vchannel
,
line
)
}
...
...
@@ -407,3 +410,41 @@ func validateChangeChannel(info *querypb.SegmentChangeInfo) (string, error) {
return
channelName
,
nil
}
// filterDuplicateChangeInfo filters out duplicated sealed segments which are both online and offline (Fix issue#17347)
func
filterDuplicateChangeInfo
(
line
*
querypb
.
SegmentChangeInfo
)
{
if
line
.
OnlineNodeID
==
line
.
OfflineNodeID
{
dupSegmentIDs
:=
make
(
map
[
UniqueID
]
struct
{})
for
_
,
onlineSegment
:=
range
line
.
OnlineSegments
{
for
_
,
offlineSegment
:=
range
line
.
OfflineSegments
{
if
onlineSegment
.
SegmentID
==
offlineSegment
.
SegmentID
&&
onlineSegment
.
SegmentState
==
segmentTypeSealed
&&
offlineSegment
.
SegmentState
==
segmentTypeSealed
{
dupSegmentIDs
[
onlineSegment
.
SegmentID
]
=
struct
{}{}
}
}
}
if
len
(
dupSegmentIDs
)
==
0
{
return
}
var
dupSegmentIDsList
[]
UniqueID
for
sid
:=
range
dupSegmentIDs
{
dupSegmentIDsList
=
append
(
dupSegmentIDsList
,
sid
)
}
log
.
Warn
(
"Found sealed segments are that are online and offline."
,
zap
.
Int64s
(
"SegmentIDs"
,
dupSegmentIDsList
))
var
filteredOnlineSegments
[]
*
querypb
.
SegmentInfo
for
_
,
onlineSegment
:=
range
line
.
OnlineSegments
{
if
_
,
ok
:=
dupSegmentIDs
[
onlineSegment
.
SegmentID
];
!
ok
{
filteredOnlineSegments
=
append
(
filteredOnlineSegments
,
onlineSegment
)
}
}
line
.
OnlineSegments
=
filteredOnlineSegments
var
filteredOfflineSegments
[]
*
querypb
.
SegmentInfo
for
_
,
offlineSegment
:=
range
line
.
OfflineSegments
{
if
_
,
ok
:=
dupSegmentIDs
[
offlineSegment
.
SegmentID
];
!
ok
{
filteredOfflineSegments
=
append
(
filteredOfflineSegments
,
offlineSegment
)
}
}
line
.
OfflineSegments
=
filteredOfflineSegments
}
}
internal/querynode/query_node_test.go
浏览文件 @
be407742
...
...
@@ -413,6 +413,74 @@ func TestQueryNode_validateChangeChannel(t *testing.T) {
}
}
func
TestQueryNode_filterDuplicateChangeInfo
(
t
*
testing
.
T
)
{
t
.
Run
(
"dup change info"
,
func
(
t
*
testing
.
T
)
{
info
:=
&
querypb
.
SegmentChangeInfo
{
OnlineNodeID
:
233
,
OnlineSegments
:
[]
*
querypb
.
SegmentInfo
{
{
SegmentID
:
23333
,
SegmentState
:
segmentTypeSealed
,
},
},
OfflineNodeID
:
233
,
OfflineSegments
:
[]
*
querypb
.
SegmentInfo
{
{
SegmentID
:
23333
,
SegmentState
:
segmentTypeSealed
,
},
},
}
filterDuplicateChangeInfo
(
info
)
assert
.
Equal
(
t
,
0
,
len
(
info
.
OnlineSegments
))
assert
.
Equal
(
t
,
0
,
len
(
info
.
OfflineSegments
))
})
t
.
Run
(
"normal change info1"
,
func
(
t
*
testing
.
T
)
{
info
:=
&
querypb
.
SegmentChangeInfo
{
OnlineNodeID
:
233
,
OnlineSegments
:
[]
*
querypb
.
SegmentInfo
{
{
SegmentID
:
23333
,
SegmentState
:
segmentTypeSealed
,
},
},
OfflineNodeID
:
234
,
OfflineSegments
:
[]
*
querypb
.
SegmentInfo
{
{
SegmentID
:
23333
,
SegmentState
:
segmentTypeSealed
,
},
},
}
filterDuplicateChangeInfo
(
info
)
assert
.
Equal
(
t
,
1
,
len
(
info
.
OnlineSegments
))
assert
.
Equal
(
t
,
1
,
len
(
info
.
OfflineSegments
))
})
t
.
Run
(
"normal change info2"
,
func
(
t
*
testing
.
T
)
{
info
:=
&
querypb
.
SegmentChangeInfo
{
OnlineNodeID
:
233
,
OnlineSegments
:
[]
*
querypb
.
SegmentInfo
{
{
SegmentID
:
23333
,
SegmentState
:
segmentTypeSealed
,
},
},
OfflineNodeID
:
234
,
OfflineSegments
:
[]
*
querypb
.
SegmentInfo
{
{
SegmentID
:
23333
,
SegmentState
:
segmentTypeSealed
,
},
},
}
filterDuplicateChangeInfo
(
info
)
assert
.
Equal
(
t
,
1
,
len
(
info
.
OnlineSegments
))
assert
.
Equal
(
t
,
1
,
len
(
info
.
OfflineSegments
))
})
}
func
TestQueryNode_handleSealedSegmentsChangeInfo
(
t
*
testing
.
T
)
{
ctx
,
cancel
:=
context
.
WithCancel
(
context
.
Background
())
defer
cancel
()
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录