Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
Crayon鑫
Paddle
提交
e5b0d9e1
P
Paddle
项目概览
Crayon鑫
/
Paddle
与 Fork 源项目一致
Fork自
PaddlePaddle / Paddle
通知
1
Star
1
Fork
0
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
1
列表
看板
标记
里程碑
合并请求
0
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
P
Paddle
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
1
Issue
1
列表
看板
标记
里程碑
合并请求
0
合并请求
0
Pages
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
提交
Issue看板
未验证
提交
e5b0d9e1
编写于
1月 20, 2021
作者:
L
liuyuhui
提交者:
GitHub
1月 20, 2021
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
[Kunlun] Add condition_variable and notify() in BindThreadedSSAGraphExecutor (#30586)
上级
ca338214
变更
2
隐藏空白更改
内联
并排
Showing
2 changed file
with
22 addition
and
12 deletion
+22
-12
paddle/fluid/framework/details/bind_threaded_ssa_graph_executor.cc
...uid/framework/details/bind_threaded_ssa_graph_executor.cc
+15
-12
paddle/fluid/framework/details/bind_threaded_ssa_graph_executor.h
...luid/framework/details/bind_threaded_ssa_graph_executor.h
+7
-0
未找到文件。
paddle/fluid/framework/details/bind_threaded_ssa_graph_executor.cc
浏览文件 @
e5b0d9e1
...
...
@@ -30,9 +30,6 @@ namespace paddle {
namespace
framework
{
namespace
details
{
static
std
::
atomic
<
unsigned
int
>
exec_op_count_
;
static
std
::
atomic
<
int
>
error_state
;
BindThreadedSSAGraphExecutor
::
BindThreadedSSAGraphExecutor
(
const
ExecutionStrategy
&
strategy
,
const
std
::
vector
<
Scope
*>
&
local_scopes
,
const
std
::
vector
<
Scope
*>
&
local_exec_scopes
,
...
...
@@ -125,7 +122,7 @@ FetchResultType BindThreadedSSAGraphExecutor::RunMainStream(
for
(
auto
cur_op
:
ready_fetch_ops
)
{
ready_ops
->
Push
(
cur_op
);
}
// Atomic variable, no need to lock
exec_op_count_
=
0
;
platform
::
XPUPlace
cur_place
;
...
...
@@ -134,9 +131,8 @@ FetchResultType BindThreadedSSAGraphExecutor::RunMainStream(
while
(
cur_count
<
op_deps_
.
size
())
{
cur_count
++
;
auto
cur_op
=
ready_ops
->
Pop
();
// when execption, get cur_op == nullptr
if
(
cur_op
==
nullptr
)
{
// sleep a while to make sure worker thread quit
sleep
(
10
);
exec_op_count_
=
op_deps_
.
size
();
break
;
}
...
...
@@ -151,14 +147,16 @@ FetchResultType BindThreadedSSAGraphExecutor::RunMainStream(
RunOpAsyncMainStream
(
cur_op
,
op_deps
.
get
(),
ready_ops
,
cur_index
);
}
}
while
(
exec_op_count_
<
op_deps_
.
size
())
{
{
std
::
unique_lock
<
std
::
mutex
>
lock
(
mutex_
);
cv_
.
wait
(
lock
,
[
&
]
{
return
exec_op_count_
>=
op_deps_
.
size
();
});
}
// Wait FetchOps.
ClearFetchOp
(
graph_
,
&
fetch_ops
);
if
(
exception_
.
IsCaught
())
{
ExecutionFinal
(
&
fetch_ops
);
}
// Wait FetchOps.
ClearFetchOp
(
graph_
,
&
fetch_ops
);
return
fetches
;
}
...
...
@@ -222,7 +220,8 @@ void BindThreadedSSAGraphExecutor::InsertFetchOps(
}
}
}
// RunMultiDeviceOpAsync function is used for Communicated OPs
// like all_reduce\broadcast among multicards.
void
BindThreadedSSAGraphExecutor
::
RunMultiDeviceOpAsync
(
OpHandleBase
*
op
,
std
::
unordered_map
<
OpHandleBase
*
,
struct
RunningItem
>
*
op_deps
,
...
...
@@ -256,10 +255,12 @@ void BindThreadedSSAGraphExecutor::RunMultiDeviceOpAsync(
ready_ops
->
Push
(
nullptr
);
exception_
.
Catch
(
std
::
current_exception
());
}
// Atomic variable, no need to lock
exec_op_count_
++
;
cv_
.
notify_all
();
});
}
// RunOpAsyncMainStream function is used for computed OPs
void
BindThreadedSSAGraphExecutor
::
RunOpAsyncMainStream
(
OpHandleBase
*
op
,
std
::
unordered_map
<
OpHandleBase
*
,
struct
RunningItem
>
*
op_deps
,
...
...
@@ -285,7 +286,9 @@ void BindThreadedSSAGraphExecutor::RunOpAsyncMainStream(
ready_ops
->
Push
(
nullptr
);
exception_
.
Catch
(
std
::
current_exception
());
}
// Atomic variable, no need to lock
exec_op_count_
++
;
cv_
.
notify_all
();
});
}
...
...
paddle/fluid/framework/details/bind_threaded_ssa_graph_executor.h
浏览文件 @
e5b0d9e1
...
...
@@ -14,7 +14,9 @@
#pragma once
#include <ThreadPool.h>
#include <condition_variable> // NOLINT
#include <memory>
#include <mutex> // NOLINT
#include <string>
#include <unordered_map>
#include <vector>
...
...
@@ -76,6 +78,11 @@ class BindThreadedSSAGraphExecutor : public SSAGraphExecutor {
::
ThreadPool
prepare_pool_
;
::
ThreadPool
multi_device_op_pool_
;
std
::
mutex
mutex_
;
std
::
condition_variable
cv_
;
std
::
atomic
<
unsigned
int
>
exec_op_count_
;
std
::
atomic
<
int
>
error_state
;
void
RunOpAsyncMainStream
(
OpHandleBase
*
op
,
std
::
unordered_map
<
OpHandleBase
*
,
struct
RunningItem
>
*
op_deps
,
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录