Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
milvus
milvus
提交
16798f4b
M
milvus
项目概览
milvus
/
milvus
大约 1 年 前同步成功
通知
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,发现更多精彩内容 >>
未验证
提交
16798f4b
编写于
1月 10, 2023
作者:
G
groot
提交者:
GitHub
1月 10, 2023
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
Update segment id for import task (#21583)
Signed-off-by:
N
groot
<
yihua.mo@zilliz.com
>
上级
7154b537
变更
2
隐藏空白更改
内联
并排
Showing
2 changed file
with
34 addition
and
1 deletion
+34
-1
internal/datanode/data_node.go
internal/datanode/data_node.go
+26
-0
internal/rootcoord/import_manager.go
internal/rootcoord/import_manager.go
+8
-1
未找到文件。
internal/datanode/data_node.go
浏览文件 @
16798f4b
...
...
@@ -1235,6 +1235,32 @@ func assignSegmentFunc(node *DataNode, req *datapb.ImportTaskRequest) importutil
zap
.
Int64
(
"segmentID"
,
segmentID
),
zap
.
Int
(
"shard ID"
,
shardID
),
zap
.
String
(
"target channel name"
,
targetChName
))
// call report to notify the rootcoord update the segment id list for this task
// ignore the returned error, since even report failed the segments still can be cleaned
retry
.
Do
(
context
.
Background
(),
func
()
error
{
importResult
:=
&
rootcoordpb
.
ImportResult
{
Status
:
&
commonpb
.
Status
{
ErrorCode
:
commonpb
.
ErrorCode_Success
,
},
TaskId
:
req
.
GetImportTask
()
.
TaskId
,
DatanodeId
:
Params
.
DataNodeCfg
.
GetNodeID
(),
State
:
commonpb
.
ImportState_ImportStarted
,
Segments
:
[]
int64
{
segmentID
},
AutoIds
:
make
([]
int64
,
0
),
RowCount
:
0
,
}
status
,
err
:=
node
.
rootCoord
.
ReportImport
(
context
.
Background
(),
importResult
)
if
err
!=
nil
{
log
.
Error
(
"fail to report import state to RootCoord"
,
zap
.
Error
(
err
))
return
err
}
if
status
!=
nil
&&
status
.
ErrorCode
!=
commonpb
.
ErrorCode_Success
{
return
errors
.
New
(
status
.
GetReason
())
}
return
nil
})
return
segmentID
,
targetChName
,
nil
}
}
...
...
internal/rootcoord/import_manager.go
浏览文件 @
16798f4b
...
...
@@ -595,7 +595,12 @@ func (m *importManager) updateTaskInfo(ir *rootcoordpb.ImportResult) (*datapb.Im
// Meta persist should be done before memory objs change.
toPersistImportTaskInfo
=
cloneImportTaskInfo
(
v
)
toPersistImportTaskInfo
.
State
.
StateCode
=
ir
.
GetState
()
toPersistImportTaskInfo
.
State
.
Segments
=
ir
.
GetSegments
()
// if is started state, append the new created segment id
if
v
.
GetState
()
.
GetStateCode
()
==
commonpb
.
ImportState_ImportStarted
{
toPersistImportTaskInfo
.
State
.
Segments
=
append
(
toPersistImportTaskInfo
.
State
.
Segments
,
ir
.
GetSegments
()
...
)
}
else
{
toPersistImportTaskInfo
.
State
.
Segments
=
ir
.
GetSegments
()
}
toPersistImportTaskInfo
.
State
.
RowCount
=
ir
.
GetRowCount
()
toPersistImportTaskInfo
.
State
.
RowIds
=
ir
.
GetAutoIds
()
for
_
,
kv
:=
range
ir
.
GetInfos
()
{
...
...
@@ -606,6 +611,8 @@ func (m *importManager) updateTaskInfo(ir *rootcoordpb.ImportResult) (*datapb.Im
toPersistImportTaskInfo
.
Infos
=
append
(
toPersistImportTaskInfo
.
Infos
,
kv
)
}
}
log
.
Info
(
"importManager update task info"
,
zap
.
Any
(
"toPersistImportTaskInfo"
,
toPersistImportTaskInfo
))
// Update task in task store.
if
err
:=
m
.
persistTaskInfo
(
toPersistImportTaskInfo
);
err
!=
nil
{
log
.
Error
(
"failed to update import task"
,
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录