Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
milvus
milvus
提交
a669440e
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 搜索 >>
未验证
提交
a669440e
编写于
7月 25, 2023
作者:
C
congqixia
提交者:
GitHub
7月 25, 2023
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
Allow querycoord executor load sealed segment with no index (#25902)
Signed-off-by:
N
Congqi Xia
<
congqi.xia@zilliz.com
>
上级
3735a097
变更
2
隐藏空白更改
内联
并排
Showing
2 changed file
with
101 addition
and
2 deletion
+101
-2
internal/querycoordv2/task/executor.go
internal/querycoordv2/task/executor.go
+6
-2
internal/querycoordv2/task/task_test.go
internal/querycoordv2/task/task_test.go
+95
-0
未找到文件。
internal/querycoordv2/task/executor.go
浏览文件 @
a669440e
...
...
@@ -21,6 +21,7 @@ import (
"sync"
"time"
"github.com/cockroachdb/errors"
"github.com/milvus-io/milvus/pkg/util/tsoutil"
"github.com/milvus-io/milvus/pkg/util/typeutil"
"go.uber.org/atomic"
...
...
@@ -269,8 +270,11 @@ func (ex *Executor) loadSegment(task *SegmentTask, step int) error {
segment
:=
resp
.
GetInfos
()[
0
]
indexes
,
err
:=
ex
.
broker
.
GetIndexInfo
(
ctx
,
task
.
CollectionID
(),
segment
.
GetID
())
if
err
!=
nil
{
log
.
Warn
(
"failed to get index of segment"
,
zap
.
Error
(
err
))
return
err
if
!
errors
.
Is
(
err
,
merr
.
ErrIndexNotFound
)
{
log
.
Warn
(
"failed to get index of segment"
,
zap
.
Error
(
err
))
return
err
}
indexes
=
nil
}
readableVersion
:=
int64
(
0
)
...
...
internal/querycoordv2/task/task_test.go
浏览文件 @
a669440e
...
...
@@ -158,6 +158,7 @@ func (suite *TaskSuite) BeforeTest(suiteName, testName string) {
switch
testName
{
case
"TestSubscribeChannelTask"
,
"TestLoadSegmentTask"
,
"TestLoadSegmentTaskNotIndex"
,
"TestLoadSegmentTaskFailed"
,
"TestSegmentTaskStale"
,
"TestTaskCanceled"
,
...
...
@@ -453,6 +454,100 @@ func (suite *TaskSuite) TestLoadSegmentTask() {
}
}
func
(
suite
*
TaskSuite
)
TestLoadSegmentTaskNotIndex
()
{
ctx
:=
context
.
Background
()
timeout
:=
10
*
time
.
Second
targetNode
:=
int64
(
3
)
partition
:=
int64
(
100
)
channel
:=
&
datapb
.
VchannelInfo
{
CollectionID
:
suite
.
collection
,
ChannelName
:
Params
.
CommonCfg
.
RootCoordDml
.
GetValue
()
+
"-test"
,
}
// Expect
suite
.
broker
.
EXPECT
()
.
GetCollectionSchema
(
mock
.
Anything
,
suite
.
collection
)
.
Return
(
&
schemapb
.
CollectionSchema
{
Name
:
"TestLoadSegmentTask"
,
Fields
:
[]
*
schemapb
.
FieldSchema
{
{
FieldID
:
100
,
Name
:
"vec"
,
DataType
:
schemapb
.
DataType_FloatVector
},
},
},
nil
)
suite
.
broker
.
EXPECT
()
.
DescribeIndex
(
mock
.
Anything
,
suite
.
collection
)
.
Return
([]
*
indexpb
.
IndexInfo
{
{
CollectionID
:
suite
.
collection
,
},
},
nil
)
for
_
,
segment
:=
range
suite
.
loadSegments
{
suite
.
broker
.
EXPECT
()
.
GetSegmentInfo
(
mock
.
Anything
,
segment
)
.
Return
(
&
datapb
.
GetSegmentInfoResponse
{
Infos
:
[]
*
datapb
.
SegmentInfo
{
{
ID
:
segment
,
CollectionID
:
suite
.
collection
,
PartitionID
:
partition
,
InsertChannel
:
channel
.
ChannelName
,
}},
},
nil
)
suite
.
broker
.
EXPECT
()
.
GetIndexInfo
(
mock
.
Anything
,
suite
.
collection
,
segment
)
.
Return
(
nil
,
merr
.
WrapErrIndexNotFound
())
}
suite
.
cluster
.
EXPECT
()
.
LoadSegments
(
mock
.
Anything
,
targetNode
,
mock
.
Anything
)
.
Return
(
merr
.
Status
(
nil
),
nil
)
// Test load segment task
suite
.
dist
.
ChannelDistManager
.
Update
(
targetNode
,
meta
.
DmChannelFromVChannel
(
&
datapb
.
VchannelInfo
{
CollectionID
:
suite
.
collection
,
ChannelName
:
channel
.
ChannelName
,
}))
tasks
:=
[]
Task
{}
segments
:=
make
([]
*
datapb
.
SegmentInfo
,
0
)
for
_
,
segment
:=
range
suite
.
loadSegments
{
segments
=
append
(
segments
,
&
datapb
.
SegmentInfo
{
ID
:
segment
,
InsertChannel
:
channel
.
ChannelName
,
PartitionID
:
1
,
})
task
,
err
:=
NewSegmentTask
(
ctx
,
timeout
,
0
,
suite
.
collection
,
suite
.
replica
,
NewSegmentAction
(
targetNode
,
ActionTypeGrow
,
channel
.
GetChannelName
(),
segment
),
)
suite
.
NoError
(
err
)
tasks
=
append
(
tasks
,
task
)
err
=
suite
.
scheduler
.
Add
(
task
)
suite
.
NoError
(
err
)
}
suite
.
broker
.
EXPECT
()
.
GetRecoveryInfoV2
(
mock
.
Anything
,
suite
.
collection
)
.
Return
(
nil
,
segments
,
nil
)
suite
.
target
.
UpdateCollectionNextTargetWithPartitions
(
suite
.
collection
,
int64
(
1
))
segmentsNum
:=
len
(
suite
.
loadSegments
)
suite
.
AssertTaskNum
(
0
,
segmentsNum
,
0
,
segmentsNum
)
// Process tasks
suite
.
dispatchAndWait
(
targetNode
)
suite
.
AssertTaskNum
(
segmentsNum
,
0
,
0
,
segmentsNum
)
// Process tasks done
// Dist contains channels
view
:=
&
meta
.
LeaderView
{
ID
:
targetNode
,
CollectionID
:
suite
.
collection
,
Segments
:
map
[
int64
]
*
querypb
.
SegmentDist
{},
}
for
_
,
segment
:=
range
suite
.
loadSegments
{
view
.
Segments
[
segment
]
=
&
querypb
.
SegmentDist
{
NodeID
:
targetNode
,
Version
:
0
}
}
distSegments
:=
lo
.
Map
(
segments
,
func
(
info
*
datapb
.
SegmentInfo
,
_
int
)
*
meta
.
Segment
{
return
meta
.
SegmentFromInfo
(
info
)
})
suite
.
dist
.
LeaderViewManager
.
Update
(
targetNode
,
view
)
suite
.
dist
.
SegmentDistManager
.
Update
(
targetNode
,
distSegments
...
)
suite
.
dispatchAndWait
(
targetNode
)
suite
.
AssertTaskNum
(
0
,
0
,
0
,
0
)
for
_
,
task
:=
range
tasks
{
suite
.
Equal
(
TaskStatusSucceeded
,
task
.
Status
())
suite
.
NoError
(
task
.
Err
())
}
}
func
(
suite
*
TaskSuite
)
TestLoadSegmentTaskFailed
()
{
ctx
:=
context
.
Background
()
timeout
:=
10
*
time
.
Second
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录