Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
milvus
milvus
提交
33e1e170
M
milvus
项目概览
milvus
/
milvus
11 个月 前同步成功
通知
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,发现更多精彩内容 >>
提交
33e1e170
编写于
1月 04, 2021
作者:
B
bigsheeper
提交者:
yefu.chen
1月 04, 2021
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
Fix writer client DescribeSegment api
Signed-off-by:
N
bigsheeper
<
yihao.dai@zilliz.com
>
上级
0a383608
变更
8
隐藏空白更改
内联
并排
Showing
8 changed file
with
34 addition
and
39 deletion
+34
-39
internal/kv/etcd/etcd_kv.go
internal/kv/etcd/etcd_kv.go
+12
-0
internal/master/client.go
internal/master/client.go
+2
-2
internal/master/index_builder_scheduler.go
internal/master/index_builder_scheduler.go
+0
-1
internal/master/index_load_scheduler.go
internal/master/index_load_scheduler.go
+1
-8
internal/master/index_task.go
internal/master/index_task.go
+0
-1
internal/master/master.go
internal/master/master.go
+3
-15
internal/master/param_table.go
internal/master/param_table.go
+0
-12
internal/writenode/client/client.go
internal/writenode/client/client.go
+16
-0
未找到文件。
internal/kv/etcd/etcd_kv.go
浏览文件 @
33e1e170
...
...
@@ -68,6 +68,18 @@ func (kv *EtcdKV) Load(key string) (string, error) {
return
string
(
resp
.
Kvs
[
0
]
.
Value
),
nil
}
func
(
kv
*
EtcdKV
)
GetCount
(
key
string
)
(
int64
,
error
)
{
key
=
path
.
Join
(
kv
.
rootPath
,
key
)
ctx
,
cancel
:=
context
.
WithTimeout
(
context
.
TODO
(),
RequestTimeout
)
defer
cancel
()
resp
,
err
:=
kv
.
client
.
Get
(
ctx
,
key
)
if
err
!=
nil
{
return
-
1
,
err
}
return
resp
.
Count
,
nil
}
func
(
kv
*
EtcdKV
)
MultiLoad
(
keys
[]
string
)
([]
string
,
error
)
{
ops
:=
make
([]
clientv3
.
Op
,
0
,
len
(
keys
))
for
_
,
keyLoad
:=
range
keys
{
...
...
internal/master/client.go
浏览文件 @
33e1e170
...
...
@@ -90,12 +90,12 @@ func (m *MockBuildIndexClient) GetIndexFilePaths(indexID UniqueID) ([]string, er
}
type
LoadIndexClient
interface
{
LoadIndex
(
indexPaths
[]
string
,
segmentID
int64
,
fieldID
int64
,
fieldName
string
,
indexParams
map
[
string
]
string
)
error
LoadIndex
(
indexPaths
[]
string
,
segmentID
int64
,
fieldID
int64
,
fieldName
string
)
error
}
type
MockLoadIndexClient
struct
{
}
func
(
m
*
MockLoadIndexClient
)
LoadIndex
(
indexPaths
[]
string
,
segmentID
int64
,
fieldID
int64
,
fieldName
string
,
indexParams
map
[
string
]
string
)
error
{
func
(
m
*
MockLoadIndexClient
)
LoadIndex
(
indexPaths
[]
string
,
segmentID
int64
,
fieldID
int64
,
fieldName
string
)
error
{
return
nil
}
internal/master/index_builder_scheduler.go
浏览文件 @
33e1e170
...
...
@@ -133,7 +133,6 @@ func (scheduler *IndexBuildScheduler) describe() error {
fieldID
:
indexBuildInfo
.
fieldID
,
fieldName
:
fieldName
,
indexFilePaths
:
filePaths
,
indexParams
:
channelInfo
.
indexParams
,
}
// Save data to meta table
err
=
scheduler
.
metaTable
.
UpdateFieldIndexMeta
(
&
etcdpb
.
FieldIndexMeta
{
...
...
internal/master/index_load_scheduler.go
浏览文件 @
33e1e170
...
...
@@ -3,15 +3,12 @@ package master
import
(
"context"
"log"
"github.com/zilliztech/milvus-distributed/internal/proto/commonpb"
)
type
IndexLoadInfo
struct
{
segmentID
UniqueID
fieldID
UniqueID
fieldName
string
indexParams
[]
*
commonpb
.
KeyValuePair
indexFilePaths
[]
string
}
...
...
@@ -39,11 +36,7 @@ func NewIndexLoadScheduler(ctx context.Context, client LoadIndexClient, metaTabl
func
(
scheduler
*
IndexLoadScheduler
)
schedule
(
info
interface
{})
error
{
indexLoadInfo
:=
info
.
(
*
IndexLoadInfo
)
indexParams
:=
make
(
map
[
string
]
string
)
for
_
,
kv
:=
range
indexLoadInfo
.
indexParams
{
indexParams
[
kv
.
Key
]
=
kv
.
Value
}
err
:=
scheduler
.
client
.
LoadIndex
(
indexLoadInfo
.
indexFilePaths
,
indexLoadInfo
.
segmentID
,
indexLoadInfo
.
fieldID
,
indexLoadInfo
.
fieldName
,
indexParams
)
err
:=
scheduler
.
client
.
LoadIndex
(
indexLoadInfo
.
indexFilePaths
,
indexLoadInfo
.
segmentID
,
indexLoadInfo
.
fieldID
,
indexLoadInfo
.
fieldName
)
//TODO: Save data to meta table
if
err
!=
nil
{
return
err
...
...
internal/master/index_task.go
浏览文件 @
33e1e170
...
...
@@ -68,7 +68,6 @@ func (task *createIndexTask) Execute() error {
fieldID
:
fieldID
,
fieldName
:
task
.
req
.
FieldName
,
indexFilePaths
:
indexMeta
.
IndexFilePaths
,
indexParams
:
indexMeta
.
IndexParams
,
})
if
err
!=
nil
{
return
err
...
...
internal/master/master.go
浏览文件 @
33e1e170
...
...
@@ -10,12 +10,6 @@ import (
"sync/atomic"
"time"
"github.com/zilliztech/milvus-distributed/internal/querynode/client"
indexbuilderclient
"github.com/zilliztech/milvus-distributed/internal/indexbuilder/client"
writerclient
"github.com/zilliztech/milvus-distributed/internal/writenode/client"
etcdkv
"github.com/zilliztech/milvus-distributed/internal/kv/etcd"
ms
"github.com/zilliztech/milvus-distributed/internal/msgstream"
"github.com/zilliztech/milvus-distributed/internal/proto/masterpb"
...
...
@@ -181,15 +175,9 @@ func CreateServer(ctx context.Context) (*Master, error) {
m
.
scheduler
.
SetDDMsgStream
(
pulsarDDStream
)
m
.
scheduler
.
SetIDAllocator
(
func
()
(
UniqueID
,
error
)
{
return
m
.
idAllocator
.
AllocOne
()
})
flushClient
,
err
:=
writerclient
.
NewWriterClient
(
Params
.
EtcdAddress
,
kvRootPath
,
Params
.
WriteNodeSegKvSubPath
,
pulsarDDStream
)
if
err
!=
nil
{
return
nil
,
err
}
buildIndexClient
,
err
:=
indexbuilderclient
.
NewBuildIndexClient
(
ctx
,
Params
.
IndexBuilderAddress
)
if
err
!=
nil
{
return
nil
,
err
}
loadIndexClient
:=
client
.
NewLoadIndexClient
(
ctx
,
Params
.
PulsarAddress
,
Params
.
LoadIndexChannelNames
)
flushClient
:=
&
MockWriteNodeClient
{}
buildIndexClient
:=
&
MockBuildIndexClient
{}
loadIndexClient
:=
&
MockLoadIndexClient
{}
m
.
indexLoadSch
=
NewIndexLoadScheduler
(
ctx
,
loadIndexClient
,
m
.
metaTable
)
m
.
indexBuildSch
=
NewIndexBuildScheduler
(
ctx
,
buildIndexClient
,
m
.
metaTable
,
m
.
indexLoadSch
)
...
...
internal/master/param_table.go
浏览文件 @
33e1e170
...
...
@@ -50,8 +50,6 @@ type ParamTable struct {
MaxPartitionNum
int64
DefaultPartitionTag
string
LoadIndexChannelNames
[]
string
}
var
Params
ParamTable
...
...
@@ -99,8 +97,6 @@ func (p *ParamTable) Init() {
p
.
initMsgChannelSubName
()
p
.
initMaxPartitionNum
()
p
.
initDefaultPartitionTag
()
p
.
initLoadIndexChannelNames
()
}
func
(
p
*
ParamTable
)
initAddress
()
{
...
...
@@ -360,11 +356,3 @@ func (p *ParamTable) initDefaultPartitionTag() {
p
.
DefaultPartitionTag
=
defaultTag
}
func
(
p
*
ParamTable
)
initLoadIndexChannelNames
()
{
loadIndexChannelName
,
err
:=
p
.
Load
(
"msgChannel.chanNamePrefix.cmd"
)
if
err
!=
nil
{
panic
(
err
)
}
p
.
LoadIndexChannelNames
=
[]
string
{
loadIndexChannelName
}
}
internal/writenode/client/client.go
浏览文件 @
33e1e170
...
...
@@ -6,6 +6,7 @@ import (
"github.com/golang/protobuf/proto"
"go.etcd.io/etcd/clientv3"
"github.com/zilliztech/milvus-distributed/internal/errors"
"github.com/zilliztech/milvus-distributed/internal/kv"
etcdkv
"github.com/zilliztech/milvus-distributed/internal/kv/etcd"
"github.com/zilliztech/milvus-distributed/internal/msgstream"
...
...
@@ -79,6 +80,21 @@ func (c *Client) DescribeSegment(segmentID UniqueID) (*SegmentDescription, error
}
key
:=
c
.
kvPrefix
+
strconv
.
FormatInt
(
segmentID
,
10
)
etcdKV
,
ok
:=
c
.
kvClient
.
(
*
etcdkv
.
EtcdKV
)
if
!
ok
{
return
nil
,
errors
.
New
(
"type assertion failed for etcd kv"
)
}
count
,
err
:=
etcdKV
.
GetCount
(
key
)
if
err
!=
nil
{
return
nil
,
err
}
if
count
<=
0
{
ret
.
IsClosed
=
false
return
ret
,
nil
}
value
,
err
:=
c
.
kvClient
.
Load
(
key
)
if
err
!=
nil
{
return
ret
,
err
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录