Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
milvus
milvus
提交
71a7fef5
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 搜索 >>
未验证
提交
71a7fef5
编写于
5月 26, 2023
作者:
Y
yah01
提交者:
GitHub
5月 26, 2023
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
Consume threads of the number of nq (#24410)
Signed-off-by:
N
yah01
<
yang.cen@zilliz.com
>
上级
6209d5d7
变更
2
隐藏空白更改
内联
并排
Showing
2 changed file
with
26 addition
and
1 deletion
+26
-1
internal/querynodev2/tasks/scheduler.go
internal/querynodev2/tasks/scheduler.go
+21
-1
internal/querynodev2/tasks/task.go
internal/querynodev2/tasks/task.go
+5
-0
未找到文件。
internal/querynodev2/tasks/scheduler.go
浏览文件 @
71a7fef5
...
...
@@ -3,6 +3,7 @@ package tasks
import
(
"context"
"fmt"
"sync"
"go.uber.org/atomic"
...
...
@@ -25,7 +26,9 @@ type Scheduler struct {
queryProcessQueue
chan
*
QueryTask
queryWaitQueue
chan
*
QueryTask
pool
*
conc
.
Pool
[
any
]
pool
*
conc
.
Pool
[
any
]
runningThreadNum
int
cond
*
sync
.
Cond
}
func
NewScheduler
()
*
Scheduler
{
...
...
@@ -39,6 +42,7 @@ func NewScheduler() *Scheduler {
// queryProcessQueue: make(chan),
pool
:
conc
.
NewPool
[
any
](
maxReadConcurrency
,
ants
.
WithPreAlloc
(
true
)),
cond
:
sync
.
NewCond
(
&
sync
.
Mutex
{}),
}
}
...
...
@@ -151,7 +155,23 @@ func (s *Scheduler) processAll(ctx context.Context) {
}
func
(
s
*
Scheduler
)
process
(
t
Task
)
{
s
.
cond
.
L
.
Lock
()
for
s
.
runningThreadNum
>=
s
.
pool
.
Cap
()
{
s
.
cond
.
Wait
()
}
s
.
runningThreadNum
+=
t
.
Weight
()
s
.
cond
.
L
.
Unlock
()
s
.
pool
.
Submit
(
func
()
(
any
,
error
)
{
defer
func
()
{
s
.
cond
.
L
.
Lock
()
defer
s
.
cond
.
L
.
Unlock
()
s
.
runningThreadNum
-=
t
.
Weight
()
if
s
.
runningThreadNum
<
s
.
pool
.
Cap
()
{
s
.
cond
.
Broadcast
()
}
}()
metrics
.
QueryNodeReadTaskConcurrency
.
WithLabelValues
(
fmt
.
Sprint
(
paramtable
.
GetNodeID
()))
.
Inc
()
err
:=
t
.
Execute
()
...
...
internal/querynodev2/tasks/task.go
浏览文件 @
71a7fef5
...
...
@@ -25,6 +25,7 @@ type Task interface {
Done
(
err
error
)
Canceled
()
error
Wait
()
error
Weight
()
int
}
type
SearchTask
struct
{
...
...
@@ -235,6 +236,10 @@ func (t *SearchTask) Wait() error {
return
<-
t
.
notifier
}
func
(
t
*
SearchTask
)
Weight
()
int
{
return
int
(
t
.
nq
)
}
func
(
t
*
SearchTask
)
Result
()
*
internalpb
.
SearchResults
{
return
t
.
result
}
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录