Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
大炮V587
oceanbase
提交
1abd119e
O
oceanbase
项目概览
大炮V587
/
oceanbase
与 Fork 源项目一致
Fork自
oceanbase / oceanbase
通知
1
Star
0
Fork
0
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
0
列表
看板
标记
里程碑
合并请求
0
DevOps
流水线
流水线任务
计划
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
O
oceanbase
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
0
Issue
0
列表
看板
标记
里程碑
合并请求
0
合并请求
0
Pages
DevOps
DevOps
流水线
流水线任务
计划
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
流水线任务
提交
Issue看板
前往新版Gitcode,体验更适合开发者的 AI 搜索 >>
提交
1abd119e
编写于
3月 17, 2023
作者:
O
obdev
提交者:
ob-robot
3月 17, 2023
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
revert 'Reset PX DataHub whole msg in rescan scenario'
上级
843888c9
变更
3
隐藏空白更改
内联
并排
Showing
3 changed file
with
4 addition
and
24 deletion
+4
-24
src/sql/engine/px/datahub/components/ob_dh_range_dist_wf.cpp
src/sql/engine/px/datahub/components/ob_dh_range_dist_wf.cpp
+0
-2
src/sql/engine/px/datahub/ob_dh_msg_provider.h
src/sql/engine/px/datahub/ob_dh_msg_provider.h
+4
-9
src/sql/engine/px/ob_px_sqc_proxy.h
src/sql/engine/px/ob_px_sqc_proxy.h
+0
-13
未找到文件。
src/sql/engine/px/datahub/components/ob_dh_range_dist_wf.cpp
浏览文件 @
1abd119e
...
...
@@ -298,8 +298,6 @@ int ObRDWFPieceMsgCtx::send_whole_msg(common::ObIArray<ObPxSqcMeta *> &sqcs)
void
ObRDWFPieceMsgCtx
::
reset_resource
()
{
received_
=
0
;
infos_
.
reset
();
arena_alloc_
.
reset
();
}
int
ObRDWFWholeMsg
::
assign
(
const
ObRDWFWholeMsg
&
msg
)
...
...
src/sql/engine/px/datahub/ob_dh_msg_provider.h
浏览文件 @
1abd119e
...
...
@@ -25,26 +25,20 @@ namespace sql
class
ObPxDatahubDataProvider
{
public:
ObPxDatahubDataProvider
()
:
op_id_
(
-
1
),
msg_type_
(
dtl
::
TESTING
),
send_msg_cnt_
(
0
),
msg_set_
(
false
)
{
}
virtual
int
get_msg_nonblock
(
const
dtl
::
ObDtlMsg
*&
msg
,
int64_t
timeout_ts
)
=
0
;
virtual
void
reset
()
{
msg_set_
=
false
;
}
virtual
void
reset
()
{}
TO_STRING_KV
(
K_
(
op_id
),
K_
(
msg_type
));
uint64_t
op_id_
;
// 注册本 provider 的算子 id,用于 provder 数组里寻址对应 provider
dtl
::
ObDtlMsgType
msg_type_
;
volatile
int64_t
send_msg_cnt_
;
bool
msg_set_
;
};
template
<
typename
T
>
class
ObWholeMsgProvider
:
public
ObPxDatahubDataProvider
{
public:
ObWholeMsgProvider
()
{}
ObWholeMsgProvider
()
:
msg_set_
(
false
)
{}
virtual
~
ObWholeMsgProvider
()
=
default
;
virtual
void
reset
()
override
{
msg_
.
reset
();
ObPxDatahubDataProvider
::
reset
()
;
}
virtual
void
reset
()
override
{
msg_
.
reset
();
msg_set_
=
false
;
}
int
get_msg_nonblock
(
const
dtl
::
ObDtlMsg
*&
msg
,
int64_t
timeout_ts
)
{
int
ret
=
OB_SUCCESS
;
...
...
@@ -89,6 +83,7 @@ private:
return
ret
;
}
private:
bool
msg_set_
;
T
msg_
;
common
::
ObThreadCond
msg_ready_cond_
;
};
...
...
src/sql/engine/px/ob_px_sqc_proxy.h
浏览文件 @
1abd119e
...
...
@@ -225,19 +225,6 @@ int ObPxSQCProxy::get_dh_msg(
SQL_LOG
(
WARN
,
"fail push data to channel"
,
K
(
ret
));
}
else
if
(
OB_FAIL
(
ch
->
flush
()))
{
SQL_LOG
(
WARN
,
"fail flush dtl data"
,
K
(
ret
));
}
else
{
// The whole message should be reset in next rescan, we reset it after last piece msg
// send in sending piece msg and receiving whole msg scenario (need send && wait whole msg).
if
(
need_wait_whole_msg
)
{
const
int64_t
task_cnt
=
get_task_count
();
if
(
provider
->
send_msg_cnt_
%
task_cnt
==
0
)
{
provider
->
msg_set_
=
false
;
}
provider
->
send_msg_cnt_
+=
1
;
if
(
provider
->
send_msg_cnt_
%
task_cnt
==
0
)
{
provider
->
reset
();
// reset whole message
}
}
}
}
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录