forked from huawei/mindspore2022
!21140 [MS][LITE]fix bug about binding cpu core
Merge pull request !21140 from 张学同/to_merge
This commit is contained in:
commit
f7e3f31736
|
|
@ -13,7 +13,8 @@
|
|||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
#include <sched.h>
|
||||
#include <unistd.h>
|
||||
#include "thread/actor_threadpool.h"
|
||||
#include "thread/core_affinity.h"
|
||||
|
||||
|
|
@ -26,6 +27,7 @@ void ActorWorker::CreateThread(ActorThreadPool *pool) {
|
|||
}
|
||||
|
||||
void ActorWorker::RunWithSpin() {
|
||||
SetAffinity();
|
||||
#if !defined(__APPLE__) && !defined(SUPPORT_MSVC)
|
||||
static std::atomic_int index = {0};
|
||||
pthread_setname_np(pthread_self(), ("ActorThread_" + std::to_string(index++)).c_str());
|
||||
|
|
@ -116,7 +118,7 @@ void ActorThreadPool::PushActorToQueue(ActorBase *actor) {
|
|||
actor_queue_.push(actor);
|
||||
#endif
|
||||
}
|
||||
THREAD_INFO("actor[%s] enqueue success", actor->GetAID().Name().c_str());
|
||||
THREAD_DEBUG("actor[%s] enqueue success", actor->GetAID().Name().c_str());
|
||||
// active one idle actor thread if exist
|
||||
for (size_t i = 0; i < actor_thread_num_; ++i) {
|
||||
auto worker = reinterpret_cast<ActorWorker *>(workers_[i]);
|
||||
|
|
@ -126,11 +128,13 @@ void ActorThreadPool::PushActorToQueue(ActorBase *actor) {
|
|||
}
|
||||
}
|
||||
|
||||
int ActorThreadPool::CreateThreads(size_t actor_thread_num, size_t all_thread_num) {
|
||||
int ActorThreadPool::CreateThreads(size_t actor_thread_num, size_t all_thread_num, const std::vector<int> &core_list) {
|
||||
#ifdef USE_HQUEUE
|
||||
actor_queue_.Init(MAX_READY_ACTOR_NR);
|
||||
#endif
|
||||
|
||||
#ifdef BIND_CORE
|
||||
affinity_->SetCoreId(core_list);
|
||||
#endif
|
||||
size_t core_num = std::thread::hardware_concurrency();
|
||||
THREAD_INFO("ThreadInfo, Actor: [%zu], All: [%zu], CoreNum: [%zu]", actor_thread_num, all_thread_num, core_num);
|
||||
actor_thread_num_ = actor_thread_num < core_num ? actor_thread_num : core_num;
|
||||
|
|
@ -142,27 +146,56 @@ int ActorThreadPool::CreateThreads(size_t actor_thread_num, size_t all_thread_nu
|
|||
std::lock_guard<std::mutex> _l(pool_mutex_);
|
||||
auto worker = new (std::nothrow) ActorWorker();
|
||||
THREAD_ERROR_IF_NULL(worker);
|
||||
#ifdef BIND_CORE
|
||||
cpu_set_t mask;
|
||||
CPU_ZERO(&mask);
|
||||
if (core_list.size() > 0) {
|
||||
CPU_SET(core_list[workers_.size() % core_list.size()], &mask);
|
||||
}
|
||||
worker->set_mask(mask);
|
||||
#endif
|
||||
worker->CreateThread(this);
|
||||
workers_.push_back(worker);
|
||||
THREAD_INFO("create actor thread[%zu]", i);
|
||||
}
|
||||
size_t kernel_thread_num = all_thread_num - actor_thread_num_;
|
||||
if (kernel_thread_num > 0) {
|
||||
return ThreadPool::CreateThreads(kernel_thread_num);
|
||||
return ThreadPool::CreateThreads(kernel_thread_num, core_list);
|
||||
}
|
||||
return THREAD_OK;
|
||||
}
|
||||
|
||||
ActorThreadPool *ActorThreadPool::CreateThreadPool(size_t actor_thread_num, size_t all_thread_num) {
|
||||
ActorThreadPool *ActorThreadPool::CreateThreadPool(size_t actor_thread_num, size_t all_thread_num, BindMode bind_mode) {
|
||||
ActorThreadPool *pool = new (std::nothrow) ActorThreadPool();
|
||||
if (pool == nullptr) {
|
||||
return nullptr;
|
||||
}
|
||||
int ret = pool->CreateThreads(actor_thread_num, all_thread_num);
|
||||
int ret;
|
||||
std::vector<int> core_list;
|
||||
#ifdef BIND_CORE
|
||||
ret = pool->InitAffinityInfo();
|
||||
if (ret != THREAD_OK) {
|
||||
delete pool;
|
||||
return nullptr;
|
||||
}
|
||||
core_list = pool->affinity_->GetCoreId(all_thread_num, bind_mode);
|
||||
#endif // BIND_CORE
|
||||
ret = pool->CreateThreads(actor_thread_num, all_thread_num, core_list);
|
||||
if (ret != THREAD_OK) {
|
||||
delete pool;
|
||||
return nullptr;
|
||||
}
|
||||
|
||||
return pool;
|
||||
}
|
||||
|
||||
ActorThreadPool *ActorThreadPool::CreateThreadPool(size_t actor_thread_num, size_t all_thread_num,
|
||||
const std::vector<int> &core_list) {
|
||||
ActorThreadPool *pool = new (std::nothrow) ActorThreadPool();
|
||||
if (pool == nullptr) {
|
||||
return nullptr;
|
||||
}
|
||||
int ret;
|
||||
#ifdef BIND_CORE
|
||||
ret = pool->InitAffinityInfo();
|
||||
if (ret != THREAD_OK) {
|
||||
|
|
@ -170,6 +203,12 @@ ActorThreadPool *ActorThreadPool::CreateThreadPool(size_t actor_thread_num, size
|
|||
return nullptr;
|
||||
}
|
||||
#endif // BIND_CORE
|
||||
ret = pool->CreateThreads(actor_thread_num, all_thread_num, core_list);
|
||||
if (ret != THREAD_OK) {
|
||||
delete pool;
|
||||
return nullptr;
|
||||
}
|
||||
|
||||
return pool;
|
||||
}
|
||||
|
||||
|
|
@ -178,7 +217,7 @@ ActorThreadPool *ActorThreadPool::CreateThreadPool(size_t thread_num) {
|
|||
if (pool == nullptr) {
|
||||
return nullptr;
|
||||
}
|
||||
int ret = pool->CreateThreads(thread_num, thread_num);
|
||||
int ret = pool->CreateThreads(thread_num, thread_num, {});
|
||||
if (ret != THREAD_OK) {
|
||||
delete pool;
|
||||
return nullptr;
|
||||
|
|
|
|||
|
|
@ -18,6 +18,7 @@
|
|||
#define MINDSPORE_CORE_MINDRT_RUNTIME_ACTOR_THREADPOOL_H_
|
||||
|
||||
#include <queue>
|
||||
#include <vector>
|
||||
#include <mutex>
|
||||
#include <atomic>
|
||||
#include <condition_variable>
|
||||
|
|
@ -43,7 +44,10 @@ class ActorWorker : public Worker {
|
|||
class ActorThreadPool : public ThreadPool {
|
||||
public:
|
||||
// create ThreadPool that contains actor thread and kernel thread
|
||||
static ActorThreadPool *CreateThreadPool(size_t actor_thread_num, size_t all_thread_num);
|
||||
static ActorThreadPool *CreateThreadPool(size_t actor_thread_num, size_t all_thread_num, BindMode bind_mode);
|
||||
|
||||
static ActorThreadPool *CreateThreadPool(size_t actor_thread_num, size_t all_thread_num,
|
||||
const std::vector<int> &core_list);
|
||||
// create ThreadPool that contains only actor thread
|
||||
static ActorThreadPool *CreateThreadPool(size_t thread_num);
|
||||
~ActorThreadPool() override;
|
||||
|
|
@ -53,7 +57,7 @@ class ActorThreadPool : public ThreadPool {
|
|||
|
||||
private:
|
||||
ActorThreadPool() {}
|
||||
int CreateThreads(size_t actor_thread_num, size_t all_thread_num);
|
||||
int CreateThreads(size_t actor_thread_num, size_t all_thread_num, const std::vector<int> &core_list);
|
||||
size_t actor_thread_num_{0};
|
||||
|
||||
std::mutex actor_mutex_;
|
||||
|
|
|
|||
|
|
@ -248,21 +248,31 @@ int CoreAffinity::InitHardwareCoreInfo() {
|
|||
return THREAD_OK;
|
||||
}
|
||||
|
||||
int CoreAffinity::InitBindCoreId(size_t thread_num, BindMode bind_mode) {
|
||||
std::vector<int> CoreAffinity::GetCoreId(size_t thread_num, BindMode bind_mode) {
|
||||
std::vector<int> bind_id;
|
||||
if (core_num_ != sorted_id_.size()) {
|
||||
THREAD_ERROR("init sorted core id failed");
|
||||
return THREAD_ERROR;
|
||||
return bind_id;
|
||||
}
|
||||
bind_id_.clear();
|
||||
if (bind_mode == Power_Higher || bind_mode == Power_NoBind) {
|
||||
for (size_t i = 0; i < thread_num; ++i) {
|
||||
bind_id_.push_back(sorted_id_[i % core_num_]);
|
||||
bind_id.push_back(sorted_id_[i % core_num_]);
|
||||
}
|
||||
} else if (bind_mode == Power_Middle) {
|
||||
for (size_t i = 0; i < thread_num; ++i) {
|
||||
bind_id_.push_back(sorted_id_[(i + higher_num_) % core_num_]);
|
||||
bind_id.push_back(sorted_id_[(i + higher_num_) % core_num_]);
|
||||
}
|
||||
} else {
|
||||
return bind_id;
|
||||
}
|
||||
return bind_id;
|
||||
}
|
||||
void CoreAffinity::SetCoreId(const std::vector<int> &core_list) { bind_id_ = core_list; }
|
||||
|
||||
int CoreAffinity::InitBindCoreId(size_t thread_num, BindMode bind_mode) {
|
||||
bind_id_.clear();
|
||||
bind_id_ = GetCoreId(thread_num, bind_mode);
|
||||
if (bind_id_.empty()) {
|
||||
return THREAD_ERROR;
|
||||
}
|
||||
return THREAD_OK;
|
||||
|
|
|
|||
|
|
@ -43,6 +43,8 @@ class CoreAffinity {
|
|||
int BindThreads(const std::vector<Worker *> &workers, const std::vector<int> &core_list);
|
||||
int BindThreads(const std::vector<Worker *> &workers, BindMode bind_mode);
|
||||
int BindProcess(BindMode bind_mode) const;
|
||||
std::vector<int> GetCoreId(size_t thread_num, BindMode bind_mode);
|
||||
void SetCoreId(const std::vector<int> &core_list);
|
||||
|
||||
private:
|
||||
#ifdef BIND_CORE
|
||||
|
|
|
|||
|
|
@ -20,14 +20,23 @@
|
|||
namespace mindspore {
|
||||
#ifdef THREAD_POOL_DEBUG
|
||||
#include <stdio.h>
|
||||
#define THREAD_DEBUG(content, args...) \
|
||||
{ printf("[DEBUG] %s|%d: " #content "\r\n", __func__, __LINE__, ##args); }
|
||||
#define THREAD_INFO(content, args...) \
|
||||
{ printf("[INFO] %s|%d: " #content "\r\n", __func__, __LINE__, ##args); }
|
||||
#define THREAD_ERROR(content, args...) \
|
||||
{ printf("[ERROR] %s|%d: " #content "\r\n", __func__, __LINE__, ##args); }
|
||||
#else
|
||||
#define THREAD_DEBUG(content, ...)
|
||||
#define THREAD_INFO(content, ...)
|
||||
#if defined(__ANDROID__)
|
||||
#include <android/log.h>
|
||||
#define THREAD_ERROR(content, args...) \
|
||||
{ __android_log_print(ANDROID_LOG_ERROR, "MS_LITE", "%s|%d: " #content "\r\n", __func__, __LINE__, ##args); }
|
||||
#else
|
||||
#define THREAD_ERROR(content, ...)
|
||||
#endif
|
||||
#endif
|
||||
|
||||
#define THREAD_ERROR_IF_NULL(ptr) \
|
||||
do { \
|
||||
|
|
|
|||
|
|
@ -13,7 +13,8 @@
|
|||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
|
||||
#include <sched.h>
|
||||
#include <unistd.h>
|
||||
#include "thread/threadpool.h"
|
||||
#include "thread/core_affinity.h"
|
||||
|
||||
|
|
@ -31,7 +32,28 @@ Worker::~Worker() {
|
|||
|
||||
void Worker::CreateThread() { thread_ = std::thread(&Worker::Run, this); }
|
||||
|
||||
void Worker::SetAffinity() {
|
||||
#ifdef BIND_CORE
|
||||
#ifdef __ANDROID__
|
||||
int ret = sched_setaffinity(gettid(), sizeof(cpu_set_t), &mask_);
|
||||
if (ret != THREAD_OK) {
|
||||
THREAD_ERROR("bind thread %d to cpu failed. ERROR %d", gettid(), errno);
|
||||
}
|
||||
return;
|
||||
#else
|
||||
#if !defined(__APPLE__) && !defined(SUPPORT_MSVC)
|
||||
int ret = pthread_setaffinity_np(pthread_self(), sizeof(cpu_set_t), &mask_);
|
||||
if (ret != THREAD_OK) {
|
||||
THREAD_ERROR("bind thread %lu to cpu failed. ERROR %d", pthread_self(), errno);
|
||||
}
|
||||
return;
|
||||
#endif
|
||||
#endif
|
||||
#endif
|
||||
}
|
||||
|
||||
void Worker::Run() {
|
||||
SetAffinity();
|
||||
#if !defined(__APPLE__) && !defined(SUPPORT_MSVC)
|
||||
static std::atomic_int index = {0};
|
||||
pthread_setname_np(pthread_self(), ("KernelThread_" + std::to_string(index++)).c_str());
|
||||
|
|
@ -105,7 +127,7 @@ ThreadPool::~ThreadPool() {
|
|||
THREAD_INFO("destruct success");
|
||||
}
|
||||
|
||||
int ThreadPool::CreateThreads(size_t thread_num) {
|
||||
int ThreadPool::CreateThreads(size_t thread_num, const std::vector<int> &core_list) {
|
||||
size_t core_num = std::thread::hardware_concurrency();
|
||||
thread_num = thread_num < core_num ? thread_num : core_num;
|
||||
THREAD_INFO("ThreadInfo, Num: [%zu], CoreNum: [%zu]", thread_num, core_num);
|
||||
|
|
@ -117,6 +139,14 @@ int ThreadPool::CreateThreads(size_t thread_num) {
|
|||
for (size_t i = 0; i < thread_num; ++i) {
|
||||
auto worker = new (std::nothrow) Worker();
|
||||
THREAD_ERROR_IF_NULL(worker);
|
||||
#ifdef BIND_CORE
|
||||
cpu_set_t mask;
|
||||
CPU_ZERO(&mask);
|
||||
if (core_list.size() > 0) {
|
||||
CPU_SET(core_list[workers_.size() % core_list.size()], &mask);
|
||||
}
|
||||
worker->set_mask(mask);
|
||||
#endif
|
||||
worker->CreateThread();
|
||||
workers_.push_back(worker);
|
||||
THREAD_INFO("create kernel thread[%zu]", i);
|
||||
|
|
@ -127,7 +157,7 @@ int ThreadPool::CreateThreads(size_t thread_num) {
|
|||
int ThreadPool::ParallelLaunch(const Func &func, Content content, int task_num) const {
|
||||
// distribute task to the KernelThread and the idle ActorThread,
|
||||
// if the task num is greater than the KernelThread num
|
||||
THREAD_INFO("launch: %d", task_num);
|
||||
THREAD_DEBUG("launch: %d", task_num);
|
||||
Task task = {func, content};
|
||||
|
||||
DistributeTask(&task, task_num);
|
||||
|
|
@ -266,12 +296,12 @@ int ThreadPool::SetProcessAffinity(BindMode bind_mode) const {
|
|||
#endif // BIND_CORE
|
||||
}
|
||||
|
||||
ThreadPool *ThreadPool::CreateThreadPool(size_t thread_num) {
|
||||
ThreadPool *ThreadPool::CreateThreadPool(size_t thread_num, const std::vector<int> &core_list) {
|
||||
ThreadPool *pool = new (std::nothrow) ThreadPool();
|
||||
if (pool == nullptr) {
|
||||
return nullptr;
|
||||
}
|
||||
int ret = pool->CreateThreads(thread_num);
|
||||
int ret = pool->CreateThreads(thread_num, core_list);
|
||||
if (ret != THREAD_OK) {
|
||||
delete pool;
|
||||
return nullptr;
|
||||
|
|
|
|||
|
|
@ -73,16 +73,21 @@ class Worker {
|
|||
|
||||
std::thread::id thread_id() const { return thread_.get_id(); }
|
||||
#ifdef BIND_CORE
|
||||
void set_mask(const cpu_set_t &mask) { mask_ = mask; }
|
||||
pthread_t handle() { return thread_.native_handle(); }
|
||||
#endif
|
||||
|
||||
protected:
|
||||
void SetAffinity();
|
||||
void Run();
|
||||
void YieldAndDeactive();
|
||||
void WaitUntilActive();
|
||||
|
||||
bool alive_{true};
|
||||
std::thread thread_;
|
||||
#ifdef BIND_CORE
|
||||
cpu_set_t mask_;
|
||||
#endif
|
||||
std::atomic_int status_{kThreadBusy};
|
||||
|
||||
std::mutex mutex_;
|
||||
|
|
@ -98,7 +103,7 @@ class Worker {
|
|||
|
||||
class ThreadPool {
|
||||
public:
|
||||
static ThreadPool *CreateThreadPool(size_t thread_num);
|
||||
static ThreadPool *CreateThreadPool(size_t thread_num, const std::vector<int> &core_list = {});
|
||||
virtual ~ThreadPool();
|
||||
|
||||
size_t thread_num() const { return workers_.size(); }
|
||||
|
|
@ -112,7 +117,7 @@ class ThreadPool {
|
|||
protected:
|
||||
ThreadPool() = default;
|
||||
|
||||
int CreateThreads(size_t thread_num);
|
||||
int CreateThreads(size_t thread_num, const std::vector<int> &core_list);
|
||||
|
||||
int InitAffinityInfo();
|
||||
|
||||
|
|
|
|||
|
|
@ -72,16 +72,21 @@ int InnerContext::Init() {
|
|||
}
|
||||
if (this->thread_pool_ == nullptr && this->IsCpuEnabled()) {
|
||||
int actor_parallel_thread = this->enable_parallel_ ? kDefaultParallelNum : 1;
|
||||
thread_pool_ = ActorThreadPool::CreateThreadPool(actor_parallel_thread, this->thread_num_);
|
||||
if (thread_pool_ == nullptr) {
|
||||
MS_LOG(ERROR) << "Create ThreadPool failed";
|
||||
return RET_NULL_PTR;
|
||||
}
|
||||
|
||||
if (this->affinity_core_list_.empty()) {
|
||||
thread_pool_->SetCpuAffinity(
|
||||
static_cast<BindMode>(this->device_list_.front().device_info_.cpu_device_info_.cpu_bind_mode_));
|
||||
auto bind_mode = static_cast<BindMode>(this->device_list_.front().device_info_.cpu_device_info_.cpu_bind_mode_);
|
||||
thread_pool_ = ActorThreadPool::CreateThreadPool(actor_parallel_thread, this->thread_num_, bind_mode);
|
||||
if (thread_pool_ == nullptr) {
|
||||
MS_LOG(ERROR) << "Create ThreadPool failed";
|
||||
return RET_NULL_PTR;
|
||||
}
|
||||
} else {
|
||||
thread_pool_->SetCpuAffinity(this->affinity_core_list_);
|
||||
thread_pool_ =
|
||||
ActorThreadPool::CreateThreadPool(actor_parallel_thread, this->thread_num_, this->affinity_core_list_);
|
||||
if (thread_pool_ == nullptr) {
|
||||
MS_LOG(ERROR) << "Create ThreadPool failed";
|
||||
return RET_NULL_PTR;
|
||||
}
|
||||
}
|
||||
}
|
||||
if (this->allocator == nullptr) {
|
||||
|
|
@ -115,7 +120,6 @@ int InnerContext::Init() {
|
|||
|
||||
InnerContext::~InnerContext() {
|
||||
if (this->thread_pool_ != nullptr) {
|
||||
thread_pool_->SetCpuAffinity(static_cast<BindMode>(NO_BIND));
|
||||
delete thread_pool_;
|
||||
this->thread_pool_ = nullptr;
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in New Issue