Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
milvus
milvus
提交
6cd73d4e
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,发现更多精彩内容 >>
未验证
提交
6cd73d4e
编写于
4月 11, 2023
作者:
C
congqixia
提交者:
GitHub
4月 11, 2023
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
Wait safe signal before release growing (#23358)
Signed-off-by:
N
Congqi Xia
<
congqi.xia@zilliz.com
>
上级
49eb4b8a
变更
3
隐藏空白更改
内联
并排
Showing
3 changed file
with
65 addition
and
15 deletion
+65
-15
internal/querynodev2/delegator/delegator_data.go
internal/querynodev2/delegator/delegator_data.go
+2
-1
internal/querynodev2/delegator/distribution.go
internal/querynodev2/delegator/distribution.go
+21
-6
internal/querynodev2/delegator/distribution_test.go
internal/querynodev2/delegator/distribution_test.go
+42
-8
未找到文件。
internal/querynodev2/delegator/delegator_data.go
浏览文件 @
6cd73d4e
...
...
@@ -360,11 +360,12 @@ func (sd *shardDelegator) LoadSegments(ctx context.Context, req *querypb.LoadSeg
Version
:
req
.
GetVersion
(),
}
})
removed
:=
sd
.
distribution
.
AddDistributions
(
entries
...
)
removed
,
signal
:=
sd
.
distribution
.
AddDistributions
(
entries
...
)
// release possible matched growing segments async
if
len
(
removed
)
>
0
{
go
func
()
{
<-
signal
worker
,
err
:=
sd
.
workerManager
.
GetWorker
(
paramtable
.
GetNodeID
())
if
err
!=
nil
{
log
.
Warn
(
"failed to get local worker when try to release related growing"
,
zap
.
Error
(
err
))
...
...
internal/querynodev2/delegator/distribution.go
浏览文件 @
6cd73d4e
...
...
@@ -29,6 +29,19 @@ const (
wildcardNodeID
=
int64
(
-
1
)
)
var
(
closedCh
chan
struct
{}
closeOnce
sync
.
Once
)
func
getClosedCh
()
chan
struct
{}
{
closeOnce
.
Do
(
func
()
{
closedCh
=
make
(
chan
struct
{})
close
(
closedCh
)
})
return
closedCh
}
// distribution is the struct to store segment distribution.
// it contains both growing and sealed segments.
type
distribution
struct
{
...
...
@@ -107,7 +120,7 @@ func (d *distribution) Serviceable() bool {
}
// AddDistributions add multiple segment entries.
func
(
d
*
distribution
)
AddDistributions
(
entries
...
SegmentEntry
)
[]
int64
{
func
(
d
*
distribution
)
AddDistributions
(
entries
...
SegmentEntry
)
([]
int64
,
chan
struct
{})
{
d
.
mut
.
Lock
()
defer
d
.
mut
.
Unlock
()
...
...
@@ -123,8 +136,12 @@ func (d *distribution) AddDistributions(entries ...SegmentEntry) []int64 {
}
}
d
.
genSnapshot
()
return
removed
ch
:=
d
.
genSnapshot
()
// no offline growing, return closed ch to skip wait
if
len
(
removed
)
==
0
{
return
removed
,
getClosedCh
()
}
return
removed
,
ch
}
// AddGrowing adds growing segment distribution.
...
...
@@ -192,9 +209,7 @@ func (d *distribution) RemoveDistributions(sealedSegments []SegmentEntry, growin
if
!
changed
{
// no change made, return closed signal channel
ch
:=
make
(
chan
struct
{})
close
(
ch
)
return
ch
return
getClosedCh
()
}
return
d
.
genSnapshot
()
...
...
internal/querynodev2/delegator/distribution_test.go
浏览文件 @
6cd73d4e
...
...
@@ -38,14 +38,16 @@ func (s *DistributionSuite) TearDownTest() {
func
(
s
*
DistributionSuite
)
TestAddDistribution
()
{
type
testCase
struct
{
tag
string
input
[]
SegmentEntry
expected
[]
SnapshotItem
tag
string
input
[]
SegmentEntry
growing
[]
SegmentEntry
expected
[]
SnapshotItem
expectedSignalClosed
bool
}
cases
:=
[]
testCase
{
{
tag
:
"one
node"
,
tag
:
"one
_
node"
,
input
:
[]
SegmentEntry
{
{
NodeID
:
1
,
...
...
@@ -71,9 +73,10 @@ func (s *DistributionSuite) TestAddDistribution() {
},
},
},
expectedSignalClosed
:
true
,
},
{
tag
:
"multiple
nodes"
,
tag
:
"multiple
_
nodes"
,
input
:
[]
SegmentEntry
{
{
NodeID
:
1
,
...
...
@@ -113,6 +116,27 @@ func (s *DistributionSuite) TestAddDistribution() {
},
},
},
expectedSignalClosed
:
true
,
},
{
tag
:
"remove_growing"
,
growing
:
[]
SegmentEntry
{
{
NodeID
:
1
,
SegmentID
:
1
},
},
input
:
[]
SegmentEntry
{
{
NodeID
:
1
,
SegmentID
:
1
},
{
NodeID
:
1
,
SegmentID
:
2
},
},
expected
:
[]
SnapshotItem
{
{
NodeID
:
1
,
Segments
:
[]
SegmentEntry
{
{
NodeID
:
1
,
SegmentID
:
1
},
{
NodeID
:
1
,
SegmentID
:
2
},
},
},
},
expectedSignalClosed
:
false
,
},
}
...
...
@@ -120,14 +144,24 @@ func (s *DistributionSuite) TestAddDistribution() {
s
.
Run
(
tc
.
tag
,
func
()
{
s
.
SetupTest
()
defer
s
.
TearDownTest
()
s
.
dist
.
Add
Distributions
(
tc
.
input
...
)
sealed
,
_
,
version
:=
s
.
dist
.
GetCurrent
(
)
defer
s
.
dist
.
FinishUsage
(
version
)
s
.
dist
.
Add
Growing
(
tc
.
growing
...
)
_
,
signal
:=
s
.
dist
.
AddDistributions
(
tc
.
input
...
)
sealed
,
_
:=
s
.
dist
.
Peek
(
)
s
.
compareSnapshotItems
(
tc
.
expected
,
sealed
)
s
.
Equal
(
tc
.
expectedSignalClosed
,
s
.
isClosedCh
(
signal
))
})
}
}
func
(
s
*
DistributionSuite
)
isClosedCh
(
ch
chan
struct
{})
bool
{
select
{
case
<-
ch
:
return
true
default
:
return
false
}
}
func
(
s
*
DistributionSuite
)
compareSnapshotItems
(
target
,
value
[]
SnapshotItem
)
{
if
!
s
.
Equal
(
len
(
target
),
len
(
value
))
{
return
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录