From 74579393dfd568c9a6a58923c9caa07a196866ad Mon Sep 17 00:00:00 2001 From: sunxiaoxia2022 Date: Sun, 5 May 2024 11:50:12 +0800 Subject: [PATCH 01/18] add cpu tp flow --- .../cpp/benchmark_app/infer_request_wrap.hpp | 1 + .../openvino/runtime/iasync_infer_request.hpp | 2 + .../openvino/runtime/icompiled_model.hpp | 3 + .../threading/cpu_streams_executor.hpp | 8 ++ .../runtime/threading/istreams_executor.hpp | 16 ++- .../src/dev/iasync_infer_request.cpp | 4 +- .../dev/threading/cpu_streams_executor.cpp | 123 +++++++++++++++++- .../cpu_streams_executor_internal.cpp | 17 ++- .../src/dev/threading/executor_manager.cpp | 2 +- .../intel_cpu/src/async_infer_request.cpp | 5 + .../intel_cpu/src/async_infer_request.h | 8 ++ src/plugins/intel_cpu/src/compiled_model.cpp | 52 ++++++-- src/plugins/intel_cpu/src/compiled_model.h | 7 + src/plugins/intel_cpu/src/config.h | 5 +- .../intel_cpu/src/cpu_streams_calculation.cpp | 55 +++++++- .../intel_cpu/src/cpu_streams_calculation.hpp | 12 ++ src/plugins/intel_cpu/src/graph.cpp | 3 + src/plugins/intel_cpu/src/infer_request.cpp | 64 +++++++++ src/plugins/intel_cpu/src/infer_request.h | 2 + 19 files changed, 361 insertions(+), 28 deletions(-) diff --git a/samples/cpp/benchmark_app/infer_request_wrap.hpp b/samples/cpp/benchmark_app/infer_request_wrap.hpp index 0c09a6d6ada..fa45fde4a61 100644 --- a/samples/cpp/benchmark_app/infer_request_wrap.hpp +++ b/samples/cpp/benchmark_app/infer_request_wrap.hpp @@ -151,6 +151,7 @@ public: const double latency, const std::exception_ptr& ptr = nullptr) { std::unique_lock lock(_mutex); + // std::cout << "put_idle_request: id: " << id << " lat_group_id: " << lat_group_id << " latency: " << latency << "\n"; if (ptr) { inferenceException = ptr; } else { diff --git a/src/inference/dev_api/openvino/runtime/iasync_infer_request.hpp b/src/inference/dev_api/openvino/runtime/iasync_infer_request.hpp index 14c6fa2d657..bfca199464b 100644 --- a/src/inference/dev_api/openvino/runtime/iasync_infer_request.hpp +++ b/src/inference/dev_api/openvino/runtime/iasync_infer_request.hpp @@ -277,6 +277,8 @@ private: m_sync_callback_executor; //!< Used to run post inference callback in synchronous pipline mutable std::mutex m_mutex; std::function m_callback; + + std::vector> m_sub_infer_requests; }; } // namespace ov diff --git a/src/inference/dev_api/openvino/runtime/icompiled_model.hpp b/src/inference/dev_api/openvino/runtime/icompiled_model.hpp index eca22b3b003..7efdee598b4 100644 --- a/src/inference/dev_api/openvino/runtime/icompiled_model.hpp +++ b/src/inference/dev_api/openvino/runtime/icompiled_model.hpp @@ -145,6 +145,9 @@ private: std::shared_ptr m_task_executor = nullptr; //!< Holds a task executor std::shared_ptr m_callback_executor = nullptr; //!< Holds a callback executor + bool m_subCompileModel = true; + std::vector> m_sub_compilemodels; + friend ov::CoreImpl; protected: diff --git a/src/inference/dev_api/openvino/runtime/threading/cpu_streams_executor.hpp b/src/inference/dev_api/openvino/runtime/threading/cpu_streams_executor.hpp index 1638ada4282..f0b68be6ee6 100644 --- a/src/inference/dev_api/openvino/runtime/threading/cpu_streams_executor.hpp +++ b/src/inference/dev_api/openvino/runtime/threading/cpu_streams_executor.hpp @@ -55,6 +55,14 @@ public: void run_sub_stream(Task task, int id) override; + void send_message(MessageInfo msg_info); + + void wait_message(); + + void infer_wait(); + + void server_wait(int streams_num); + private: struct Impl; std::unique_ptr _impl; diff --git a/src/inference/dev_api/openvino/runtime/threading/istreams_executor.hpp b/src/inference/dev_api/openvino/runtime/threading/istreams_executor.hpp index 0ad4dd181a9..e79e9311588 100644 --- a/src/inference/dev_api/openvino/runtime/threading/istreams_executor.hpp +++ b/src/inference/dev_api/openvino/runtime/threading/istreams_executor.hpp @@ -38,6 +38,20 @@ public: */ using Ptr = std::shared_ptr; + enum MsgType{ + TP, + START_INFER, + CALL_BACK + }; + + struct MessageInfo{ + MsgType msg_type; + int rank; + int data; + void* buf; + Task task; + }; + /** * @brief Defines inference thread binding type */ @@ -203,7 +217,7 @@ public: return _threadBindingOffset; } int get_sub_streams() const { - return _sub_streams; + return _sub_streams > 0 ? _sub_streams : _streams; } StreamsMode get_sub_stream_mode() const { const auto proc_type_table = get_proc_type_table(); diff --git a/src/inference/src/dev/iasync_infer_request.cpp b/src/inference/src/dev/iasync_infer_request.cpp index 9e914a7c38b..c914772c23f 100644 --- a/src/inference/src/dev/iasync_infer_request.cpp +++ b/src/inference/src/dev/iasync_infer_request.cpp @@ -208,12 +208,12 @@ std::vector ov::IAsyncInferRequest::get_profiling_info() cons } ov::SoPtr ov::IAsyncInferRequest::get_tensor(const ov::Output& port) const { - check_state(); + // check_state(); return m_sync_request->get_tensor(port); } void ov::IAsyncInferRequest::set_tensor(const ov::Output& port, const ov::SoPtr& tensor) { - check_state(); + // check_state(); return m_sync_request->set_tensor(port, tensor); } diff --git a/src/inference/src/dev/threading/cpu_streams_executor.cpp b/src/inference/src/dev/threading/cpu_streams_executor.cpp index 02849fedd33..6ce4e7a548e 100644 --- a/src/inference/src/dev/threading/cpu_streams_executor.cpp +++ b/src/inference/src/dev/threading/cpu_streams_executor.cpp @@ -66,11 +66,12 @@ struct CPUStreamsExecutor::Impl { } } _numaNodeId = - _impl->_config.get_streams() - ? _impl->_usedNumaNodes.at((_streamId % _impl->_config.get_streams()) / - ((_impl->_config.get_streams() + _impl->_usedNumaNodes.size() - 1) / + _impl->_config.get_sub_streams() + ? _impl->_usedNumaNodes.at((_streamId % _impl->_config.get_sub_streams()) / + ((_impl->_config.get_sub_streams() + _impl->_usedNumaNodes.size() - 1) / _impl->_usedNumaNodes.size())) : _impl->_usedNumaNodes.at(_streamId % _impl->_usedNumaNodes.size()); + // std::cout << "[ Stream ] " << _impl->_config.get_name() << " : " << _streamId << ", " << _impl->_config.get_sub_streams() << "\n"; #if OV_THREAD == OV_THREAD_TBB || OV_THREAD == OV_THREAD_TBB_AUTO if (is_cpu_map_available() && _impl->_config.get_streams_info_table().size() > 0) { init_stream(); @@ -152,6 +153,8 @@ struct CPUStreamsExecutor::Impl { _taskArena.reset(new custom::task_arena{concurrency}); _cpu_ids = stream_id < static_cast(stream_processors.size()) ? stream_processors[stream_id] : _cpu_ids; + std::cout << "stream_id: " << stream_id << " size: " << stream_processors.size() + << " cpu size: " << _cpu_ids.size() << " addr: " << _impl << " , " << _cpu_ids[0] << "\n"; if (_cpu_ids.size() > 0) { CpuSet processMask; int ncpus = 0; @@ -170,7 +173,7 @@ struct CPUStreamsExecutor::Impl { int max_threads_per_core; StreamCreateType stream_type; const auto org_proc_type_table = get_org_proc_type_table(); - int streams_num = _impl->_config.get_streams(); + int streams_num = _impl->_config.get_sub_streams(); const auto stream_id = streams_num == 0 ? 0 : (_sub_stream_id >= 0 ? streams_num + _sub_stream_id : _streamId % streams_num); get_cur_stream_info(stream_id, @@ -320,8 +323,9 @@ struct CPUStreamsExecutor::Impl { this) { _exectorMgr = executor_manager(); auto numaNodes = get_available_numa_nodes(); - int streams_num = _config.get_streams(); - int sub_streams_num = _config.get_sub_streams(); + int streams_num = _config.get_sub_streams(); + int sub_streams_num = 0;//_config.get_sub_streams(); + // std::cout << "[ Impl ] " << _config.get_name() << " : " << streams_num << "\n"; if (streams_num != 0) { std::copy_n(std::begin(numaNodes), std::min(streams_num, numaNodes.size()), @@ -340,6 +344,7 @@ struct CPUStreamsExecutor::Impl { { std::unique_lock lock(_mutex); _queueCondVar.wait(lock, [&] { + std::cout << _config.get_name() << " addr: " << this << " in : " << streamId << "\n"; return !_taskQueue.empty() || (stopped = _isStopped); }); if (!_taskQueue.empty()) { @@ -422,6 +427,78 @@ struct CPUStreamsExecutor::Impl { } } + void send_message(MessageInfo msg_info) { + { + std::lock_guard lock(_msgMutex); + _messageQueue.push_back(msg_info); + // std::cout << "send_" << _streamId << " : " << msg_info.msg_type << "\n"; + } + _msgCondVar.notify_all(); + } + + void wait_message() { + std::unique_lock lock(_readMutex); + _readCondVar.wait(lock, [&] { + std::cout << "wait_" << _streamId << " : " << _readQueue[_streamId].size() << " / " + << _config.get_sub_streams() - 1 << "\n"; + return _readQueue[_streamId].size() >= _config.get_sub_streams() - 1; + }); + std::cout << "wait_" << _streamId << " end\n"; + } + + void infer_wait() { + std::unique_lock lock(_inferMutex); + // std::cout << "infer_wait ......\n"; + _inferCondVar.wait(lock); + } + + void server_wait(int streams_num) { + if (!_serverThread.joinable()) { + _messageQueue.clear(); + _readQueue.assign(streams_num, std::vector()); + MsgType msg_type; + _serverThread = std::thread([&, streams_num]() { + int count = 0; + while (!_isServerStopped) { + std::vector msgQueue; + { + // std::cout << "server_wait ........\n"; + std::unique_lock lock(_msgMutex); + while (_messageQueue.empty()) { + _msgCondVar.wait(lock); + } + std::swap(_messageQueue, msgQueue); + // std::cout << "server_wait receive: " << msgQueue[0].msg_type << " data:" << msgQueue[0].data + // << " / " << msgQueue.size() << "\n"; + } + + for (auto rec_info : msgQueue) { + msg_type = rec_info.msg_type; + if (msg_type == START_INFER) { + Task task = std::move(rec_info.task); + task(); + } else if (msg_type == TP) { + for (int i = 0; i < streams_num; i++) { + if (rec_info.data != i) { + std::lock_guard lock(_readMutex); + _readQueue[i].push_back(rec_info); + } + } + _readCondVar.notify_all(); + } else if (msg_type == CALL_BACK) { // CALL_BACK + count++; + std::cout << "server_wait CALL_BACK: " << count << "/" << streams_num << "\n"; + if (count == streams_num) { + _inferCondVar.notify_one(); + count = 0; + } + } + } + } + }); + } + } + struct SubQueue { std::mutex _subMutex; std::condition_variable _subQueueCondVar; @@ -460,14 +537,25 @@ struct CPUStreamsExecutor::Impl { int _subStreamsNum = 0; std::vector _threads; std::vector _subThreads; + std::thread _serverThread; std::mutex _mutex; std::condition_variable _queueCondVar; std::queue _taskQueue; bool _isStopped = false; + bool _isServerStopped = false; std::vector> _subTaskThread; std::vector _usedNumaNodes; CustomThreadLocal _streams; std::shared_ptr _exectorMgr; + std::vector _messageQueue; + std::vector> _readQueue; + std::mutex _msgMutex; + std::mutex _readMutex; + std::mutex _inferMutex; + std::condition_variable _msgCondVar; + std::condition_variable _readCondVar; + std::condition_variable _inferCondVar; + bool _isExit = false; }; int CPUStreamsExecutor::get_stream_id() { @@ -485,6 +573,22 @@ int CPUStreamsExecutor::get_socket_id() { return stream->_socketId; } +void CPUStreamsExecutor::send_message(MessageInfo msg_info) { + _impl->send_message(msg_info); +} + +void CPUStreamsExecutor::wait_message() { + _impl->wait_message(); +} + +void CPUStreamsExecutor::infer_wait() { + _impl->infer_wait(); +} + +void CPUStreamsExecutor::server_wait(int streams_num) { + _impl->server_wait(streams_num); +} + CPUStreamsExecutor::CPUStreamsExecutor(const IStreamsExecutor::Config& config) : _impl{new Impl{config}} {} CPUStreamsExecutor::~CPUStreamsExecutor() { @@ -510,6 +614,11 @@ CPUStreamsExecutor::~CPUStreamsExecutor() { thread.join(); } } + _impl->_isServerStopped = true; + _impl->_msgCondVar.notify_one(); + if (_impl->_serverThread.joinable()) { + _impl->_serverThread.join(); + } } void CPUStreamsExecutor::execute(Task task) { @@ -517,7 +626,7 @@ void CPUStreamsExecutor::execute(Task task) { } void CPUStreamsExecutor::run(Task task) { - if (0 == _impl->_config.get_streams()) { + if (0 == _impl->_config.get_sub_streams()) { _impl->Defer(std::move(task)); } else { _impl->Enqueue(std::move(task)); diff --git a/src/inference/src/dev/threading/cpu_streams_executor_internal.cpp b/src/inference/src/dev/threading/cpu_streams_executor_internal.cpp index db47593d895..c3a76741119 100644 --- a/src/inference/src/dev/threading/cpu_streams_executor_internal.cpp +++ b/src/inference/src/dev/threading/cpu_streams_executor_internal.cpp @@ -91,8 +91,15 @@ void reserve_cpu_by_streams_info(const std::vector> _streams_in int num_conditions = 0; int condition_idx = 0; bool last_all_proc = false; + bool sub_stream_enable = false; for (size_t i = 0; i < _streams_info_table.size(); i++) { + if (i > 0 && _streams_info_table[i][NUMBER_OF_STREAMS] < 0 && + _streams_info_table[i - 1][NUMBER_OF_STREAMS] > 0) { + stream_pos.clear(); + num_streams = 0; + sub_stream_enable = true; + } if (_streams_info_table[i][NUMBER_OF_STREAMS] != 0) { stream_pos.push_back(num_streams); } @@ -107,11 +114,15 @@ void reserve_cpu_by_streams_info(const std::vector> _streams_in std::vector proc_types; std::vector numa_nodes; std::vector sockets; - if (_streams_info_table[i][NUMBER_OF_STREAMS] != 0) { + if ((_streams_info_table[i][NUMBER_OF_STREAMS] > 0 && !sub_stream_enable) || + (_streams_info_table[i][NUMBER_OF_STREAMS] < 0 && sub_stream_enable)) { streams_table.push_back(_streams_info_table[i]); if (_streams_info_table[i][NUMBER_OF_STREAMS] < 0) { - streams_table[streams_table.size() - 1][NUMBER_OF_STREAMS] = 1; + streams_table[streams_table.size() - 1][NUMBER_OF_STREAMS] = + std::abs(_streams_info_table[i][NUMBER_OF_STREAMS]); } + } else { + continue; } if (last_all_proc && _streams_info_table[i][NUMBER_OF_STREAMS] != 0) { last_all_proc = false; @@ -135,7 +146,7 @@ void reserve_cpu_by_streams_info(const std::vector> _streams_in } } } - if (_streams_info_table[i][PROC_TYPE] > ALL_PROC && _streams_info_table[i][NUMBER_OF_STREAMS] > 0) { + if (_streams_info_table[i][PROC_TYPE] > ALL_PROC && _streams_info_table[i][NUMBER_OF_STREAMS] != 0) { condition_idx++; } } diff --git a/src/inference/src/dev/threading/executor_manager.cpp b/src/inference/src/dev/threading/executor_manager.cpp index 3f358234a2e..c1aad790698 100644 --- a/src/inference/src/dev/threading/executor_manager.cpp +++ b/src/inference/src/dev/threading/executor_manager.cpp @@ -130,7 +130,7 @@ std::shared_ptr ExecutorManagerImpl::get_idle_c std::lock_guard guard(streamExecutorMutex); for (auto& it : cpuStreamsExecutors) { const auto& executor = it.second; - if (executor.use_count() != 1) + if (executor.use_count() != 1 && config.get_name() != "CPUStreamsExecutor") continue; auto& executorConfig = it.first; diff --git a/src/plugins/intel_cpu/src/async_infer_request.cpp b/src/plugins/intel_cpu/src/async_infer_request.cpp index b77d0059db7..b77410579e9 100644 --- a/src/plugins/intel_cpu/src/async_infer_request.cpp +++ b/src/plugins/intel_cpu/src/async_infer_request.cpp @@ -19,3 +19,8 @@ ov::intel_cpu::AsyncInferRequest::~AsyncInferRequest() { void ov::intel_cpu::AsyncInferRequest::throw_if_canceled() const { check_cancelled_state(); } + +void ov::intel_cpu::AsyncInferRequest::setSubInferRequest( + const std::vector>& requests) { + m_sub_infer_requests = requests; +} diff --git a/src/plugins/intel_cpu/src/async_infer_request.h b/src/plugins/intel_cpu/src/async_infer_request.h index ebf1a210f30..0b69fa11e3a 100644 --- a/src/plugins/intel_cpu/src/async_infer_request.h +++ b/src/plugins/intel_cpu/src/async_infer_request.h @@ -17,7 +17,15 @@ public: const std::shared_ptr& callback_executor); ~AsyncInferRequest(); + void setSubInferRequest(const std::vector>& requests); + + std::vector> getSubInferRequest() const { + return m_sub_infer_requests; + } + void throw_if_canceled() const; +private: + std::vector> m_sub_infer_requests; }; } // namespace intel_cpu diff --git a/src/plugins/intel_cpu/src/compiled_model.cpp b/src/plugins/intel_cpu/src/compiled_model.cpp index 96bfe2f4672..1bcdef7b716 100644 --- a/src/plugins/intel_cpu/src/compiled_model.cpp +++ b/src/plugins/intel_cpu/src/compiled_model.cpp @@ -60,7 +60,22 @@ CompiledModel::CompiledModel(const std::shared_ptr& model, // special case when all InferRequests are muxed into a single queue m_task_executor = m_plugin->get_executor_manager()->get_executor("CPU"); } else { - stream_executor = m_plugin->get_executor_manager()->get_idle_cpu_streams_executor(m_cfg.streamExecutorConfig); + IStreamsExecutor::Config executor_confg; + if (m_cfg.enableSubStreams) { + executor_confg = IStreamsExecutor::Config{"CPUMainStreamExecutor", + 1, + 1, + IStreamsExecutor::ThreadBindingType::NONE, + 1, + 0, + 1, + IStreamsExecutor::Config::PreferredCoreType::ANY, + {}, + true}; + } else { + executor_confg = std::move(m_cfg.streamExecutorConfig); + } + stream_executor = m_plugin->get_executor_manager()->get_idle_cpu_streams_executor(executor_confg); m_task_executor = stream_executor; } if (0 != m_cfg.streamExecutorConfig.get_streams()) { @@ -75,11 +90,12 @@ CompiledModel::CompiledModel(const std::shared_ptr& model, if (m_callback_executor) set_callback_executor(m_callback_executor); - int streams = std::max(1, m_cfg.streamExecutorConfig.get_streams()); + int streams = m_cfg.enableSubStreams ? 1 : std::max(1, m_cfg.streamExecutorConfig.get_sub_streams()); + std::cout << "streams: " << streams << "\n"; std::vector tasks; tasks.resize(streams); m_graphs.resize(streams); - if (m_cfg.streamExecutorConfig.get_streams() != 0) { + if (m_cfg.streamExecutorConfig.get_sub_streams() != 0) { auto all_graphs_ready = [&] { return std::all_of(m_graphs.begin(), m_graphs.end(), [&](Graph& graph) { return graph.IsReady(); @@ -103,16 +119,23 @@ CompiledModel::CompiledModel(const std::shared_ptr& model, } else { CompiledModel::get_graph(); } - // init sub stream threads of executor - int sub_streams = m_cfg.streamExecutorConfig.get_sub_streams(); - if (sub_streams > 0 && stream_executor != nullptr) { - std::vector tasks; - tasks.resize(sub_streams); - for (auto&& task : tasks) { - task = [] {}; + std::cout << "xxxxxxxxx: m_subCompileModel: " << m_cfg.enableSubStreams << ", " << m_cfg.streamExecutorConfig.get_sub_streams() << "\n"; + if (m_cfg.enableSubStreams) { + m_cfg.enableSubStreams = false; + for (int i = 0; i < m_cfg.streamExecutorConfig.get_sub_streams(); i++) { + m_sub_compilemodels.push_back(std::make_shared(model, plugin, m_cfg, loaded_from_cache)); } - stream_executor->run_sub_stream_and_wait(tasks); } + // init sub stream threads of executor + // int sub_streams = m_cfg.streamExecutorConfig.get_sub_streams(); + // if (sub_streams > 0 && stream_executor != nullptr) { + // std::vector tasks; + // tasks.resize(sub_streams); + // for (auto&& task : tasks) { + // task = [] {}; + // } + // stream_executor->run_sub_stream_and_wait(tasks); + // } } CompiledModel::GraphGuard::Lock CompiledModel::get_graph() const { @@ -168,6 +191,13 @@ std::shared_ptr CompiledModel::create_infer_request() co std::make_shared(std::static_pointer_cast(internal_request), get_task_executor(), get_callback_executor()); + if (m_sub_compilemodels.size() > 0) { + std::vector> requests; + for (int i = 0; i < m_sub_compilemodels.size(); i++) { + requests.push_back(m_sub_compilemodels[i]->create_infer_request()); + } + async_infer_request->setSubInferRequest(requests); + } return async_infer_request; } diff --git a/src/plugins/intel_cpu/src/compiled_model.h b/src/plugins/intel_cpu/src/compiled_model.h index 64caee60aeb..e3a25f2bd87 100644 --- a/src/plugins/intel_cpu/src/compiled_model.h +++ b/src/plugins/intel_cpu/src/compiled_model.h @@ -73,6 +73,13 @@ private: * even from main thread */ GraphGuard::Lock get_graph() const; + + std::vector> get_sub_compilemodles() const { + return m_sub_compilemodels; + } + + // bool m_subCompileModel = true; + std::vector> m_sub_compilemodels; }; } // namespace intel_cpu diff --git a/src/plugins/intel_cpu/src/config.h b/src/plugins/intel_cpu/src/config.h index ef968ab88a7..7e2ccfd2e24 100644 --- a/src/plugins/intel_cpu/src/config.h +++ b/src/plugins/intel_cpu/src/config.h @@ -70,7 +70,10 @@ struct Config { bool enableCpuPinning = true; bool changedCpuPinning = false; ov::hint::SchedulingCoreType schedulingCoreType = ov::hint::SchedulingCoreType::ANY_CORE; - std::set modelDistributionPolicy = {}; + std::set modelDistributionPolicy = {ov::hint::ModelDistributionPolicy::TENSOR_PARALLEL}; + int streamsRankLevel = 1; + int numSubStreams = 0; + bool enableSubStreams = false; bool enableHyperThreading = true; bool changedHyperThreading = false; #if defined(OPENVINO_ARCH_X86) || defined(OPENVINO_ARCH_X86_64) diff --git a/src/plugins/intel_cpu/src/cpu_streams_calculation.cpp b/src/plugins/intel_cpu/src/cpu_streams_calculation.cpp index fb110778152..789d14c9203 100644 --- a/src/plugins/intel_cpu/src/cpu_streams_calculation.cpp +++ b/src/plugins/intel_cpu/src/cpu_streams_calculation.cpp @@ -54,9 +54,9 @@ std::vector> get_streams_info_table(const int input_streams, const IStreamsExecutor::Config::StreamsMode sub_streams_model, const int& target_proc) { stream_info[PROC_TYPE] = ALL_PROC; + stream_info[NUMBER_OF_STREAMS] = 1; stream_info[NUMBER_OF_STREAMS] = sub_streams_model == IStreamsExecutor::Config::StreamsMode::SUB_STREAMS_NULL ? 1 : -1; - stream_info[THREADS_PER_STREAM] = num_threads; update_ids_method(one_proc_info); streams_info_table.push_back(stream_info); stream_info[NUMBER_OF_STREAMS] = 0; @@ -370,7 +370,7 @@ std::vector> get_streams_info_table(const int input_streams, create_one_stream(proc_socket_table[current_socket_id], proc_type_table, proc_socket_table[current_socket_id][ALL_PROC], - IStreamsExecutor::Config::StreamsMode::SUB_STREAMS_NULL); + IStreamsExecutor::Config::StreamsMode::SUB_STREAMS_FOR_SOCKET); for (size_t n_node = 0; n_node < proc_socket_table.size(); n_node++) { if (n_node != size_t(current_socket_id)) { create_one_stream(proc_socket_table[n_node], @@ -379,6 +379,23 @@ std::vector> get_streams_info_table(const int input_streams, IStreamsExecutor::Config::StreamsMode::SUB_STREAMS_FOR_SOCKET); } } + stream_info = streams_info_table[0]; + stream_info[NUMBER_OF_STREAMS] = 1; + for (size_t n = 1; n < streams_info_table.size(); n++) { + if (streams_info_table[n][NUMBER_OF_STREAMS] == -1) { + if (stream_info[PROC_TYPE] != streams_info_table[n][PROC_TYPE]) { + stream_info[PROC_TYPE] = ALL_PROC; + } + stream_info[THREADS_PER_STREAM] += streams_info_table[n][THREADS_PER_STREAM]; + if (stream_info[STREAM_NUMA_NODE_ID] != streams_info_table[n][STREAM_NUMA_NODE_ID]) { + stream_info[STREAM_NUMA_NODE_ID] = -1; + } + if (stream_info[STREAM_SOCKET_ID] != streams_info_table[n][STREAM_SOCKET_ID]) { + stream_info[STREAM_SOCKET_ID] = -1; + } + } + } + streams_info_table.insert(streams_info_table.begin(), stream_info); n_streams--; } @@ -461,6 +478,33 @@ std::vector> get_streams_info_table(const int input_streams, return streams_info_table; } +std::vector> get_streams_rank_table(const std::vector>& streams_info_table, + const int input_rank_level, + int& num_sub_streams) { + std::vector> rank_table = {}; + num_sub_streams = 0; + std::vector init_rank = {}; + int rank_level = input_rank_level == 0 ? 1 : input_rank_level; + init_rank.resize(rank_level, 0); + + for (auto& row : streams_info_table) { + if (row[NUMBER_OF_STREAMS] < 0) { + for (int i = 0; i < abs(row[NUMBER_OF_STREAMS]); i++) { + init_rank[rank_level - 1] = num_sub_streams + i; + rank_table.push_back(init_rank); + } + num_sub_streams -= row[NUMBER_OF_STREAMS]; + } + } + if (rank_level == 2) { + for (int i = num_sub_streams / 2; i < num_sub_streams; i++) { + rank_table[i][0] = 1; + rank_table[i][1] -= num_sub_streams / 2; + } + } + return rank_table; +} + int get_model_prefer_threads(const int num_streams, const std::vector> proc_type_table, const std::shared_ptr& model, @@ -607,6 +651,13 @@ std::vector> generate_stream_info(const int streams, ov::util::to_string(config.hintPerfMode), config.modelDistributionPolicy, proc_type_table); + // streams_info_table = {{1, 1, 32, -1, -1}, {-1, 1, 16, 0, 0}, {-1, 1, 16, 1, 1}}; + if (config.modelDistributionPolicy.find(ov::hint::ModelDistributionPolicy::TENSOR_PARALLEL) != + config.modelDistributionPolicy.end()) { + auto streams_rank_table = + get_streams_rank_table(streams_info_table, config.streamsRankLevel, config.numSubStreams); + config.enableSubStreams = true; + } auto cpu_reservation = get_cpu_pinning(config.enableCpuPinning, config.changedCpuPinning, proc_type_table, streams_info_table); diff --git a/src/plugins/intel_cpu/src/cpu_streams_calculation.hpp b/src/plugins/intel_cpu/src/cpu_streams_calculation.hpp index 584f9a2d9a0..e362c0373d8 100644 --- a/src/plugins/intel_cpu/src/cpu_streams_calculation.hpp +++ b/src/plugins/intel_cpu/src/cpu_streams_calculation.hpp @@ -53,6 +53,18 @@ std::vector> get_streams_info_table(const int input_streams, const std::string input_perf_hint, const std::set hint_llm_distribution_policy, const std::vector>& proc_type_table); + +/** + * @brief Generate streams rank table for tensor parallel according to streams info table. + * @param[in] streams_info_table is streams information table for tensor parallel. + * @param[in] input_rank_level is depth of rank nesting. + * @param[out] num_sub_streams is number of sub streams for tensor parallel. + * @return streams rank table which will be used by StreamsExecutor. + */ +std::vector> get_streams_rank_table(const std::vector>& streams_info_table, + const int input_rank_level, + int& num_sub_streams); + /** * @brief Get model_prefer_threads * @param[in] num_streams is target streams set by user via NUM_STREAMS or hints. diff --git a/src/plugins/intel_cpu/src/graph.cpp b/src/plugins/intel_cpu/src/graph.cpp index f23994a6e2a..4ba3c49719c 100644 --- a/src/plugins/intel_cpu/src/graph.cpp +++ b/src/plugins/intel_cpu/src/graph.cpp @@ -1094,6 +1094,9 @@ void Graph::PullOutputData(std::unordered_map>& void Graph::InferStatic(SyncInferRequest* request) { dnnl::stream stream(getEngine()); + const auto& cpuExecutor = context->getCPUStreamExecutor(); + // std::cout << "execute stream: " << cpuExecutor->get_stream_id() << "\n"; + for (const auto& node : executableGraphNodes) { VERBOSE(node, getConfig().debugCaps.verbose); PERF(node, getConfig().collectPerfCounters); diff --git a/src/plugins/intel_cpu/src/infer_request.cpp b/src/plugins/intel_cpu/src/infer_request.cpp index 93bfb0117f7..fa7d72ac8f5 100644 --- a/src/plugins/intel_cpu/src/infer_request.cpp +++ b/src/plugins/intel_cpu/src/infer_request.cpp @@ -108,6 +108,22 @@ void SyncInferRequest::infer() { OV_ITT_SCOPED_TASK(itt::domains::intel_cpu, m_profiling_task); auto graphLock = m_compiled_model->get_graph(); m_graph = &(graphLock._graph); + auto streams_executor = m_graph->context->getCPUStreamExecutor(); + + auto requests = m_asyncRequest->getSubInferRequest(); + // std::cout << "[ infer ] " << requests.size() << "\n"; + if (requests.size() > 0) { + streams_executor->server_wait(requests.size()); + ov::threading::IStreamsExecutor::MessageInfo msg_info; + msg_info.msg_type = ov::threading::IStreamsExecutor::MsgType::START_INFER; + ov::threading::Task task = [&] { + SyncInferRequest::sub_streams_infer(); + }; + msg_info.task = std::move(task); + streams_executor->send_message(msg_info); + streams_executor->infer_wait(); + return; + } throw_if_canceled(); convert_batched_tensors(); @@ -608,6 +624,54 @@ SyncInferRequest::OutputControlBlock::OutputControlBlock(const ov::element::Type m_tensor = std::make_shared(memory); } +void SyncInferRequest::sub_streams_infer() { + std::map, ov::SoPtr> input_tensors; + auto requests = m_asyncRequest->getSubInferRequest(); + auto inputs = m_asyncRequest->get_inputs(); + auto outputs = m_asyncRequest->get_outputs(); + // auto sub_models = m_compiled_model->get_sub_compilemodles(); + // auto graphLock = sub_models[0]->get_graph(); + // m_graph = &(graphLock._graph); + auto streams_executor = m_graph->context->getCPUStreamExecutor(); + size_t requests_num = requests.size(); + size_t requests_count = 0; + // std::cout << "[ sub_streams_infer ] inputs: " << inputs.size() << " requests: " << requests_num << "\n"; + + if (requests.size() > 0) { + for (const auto& input : inputs) { + auto tensor = m_asyncRequest->get_tensor(input); + input_tensors.insert({input, tensor}); + } + for (const auto& output : outputs) { + auto tensor = requests[0]->get_tensor(output); + m_asyncRequest->set_tensor(output, tensor); + } + for (size_t i = 0; i < requests_num; i++) { + for (auto& input : input_tensors) { + requests[i]->set_tensor(input.first, input.second); + } + + requests[i]->set_callback([i, requests, streams_executor](const std::exception_ptr& ptr) { + // std::cout << "set_callback------ " << i << "\n"; + ov::threading::IStreamsExecutor::MessageInfo msg_info; + msg_info.msg_type = ov::threading::IStreamsExecutor::MsgType::CALL_BACK; + streams_executor->send_message(msg_info); + }); + } + { + auto sub_models = m_compiled_model->get_sub_compilemodles(); + auto graphLock = sub_models[0]->get_graph(); + m_graph = &(graphLock._graph); + auto streams_executor = m_graph->context->getCPUStreamExecutor(); + + for (size_t i = 0; i < requests_num; i++) { + std::cout << "start_async : " << i << "\n"; + requests[i]->start_async(); + } + } + } +} + } // namespace intel_cpu } // namespace ov diff --git a/src/plugins/intel_cpu/src/infer_request.h b/src/plugins/intel_cpu/src/infer_request.h index 95eca7e4d96..5ead6f1c226 100644 --- a/src/plugins/intel_cpu/src/infer_request.h +++ b/src/plugins/intel_cpu/src/infer_request.h @@ -105,6 +105,8 @@ private: const ov::Output& get_internal_port(const ov::Output& port) const; + void sub_streams_infer(); + private: std::unordered_map m_outputControlBlocks; From bd9d37b07e02366a7f3e0bc8288514c7b7dffab5 Mon Sep 17 00:00:00 2001 From: sunxiaoxia2022 Date: Mon, 6 May 2024 14:01:36 +0800 Subject: [PATCH 02/18] change to muti streamexecutors --- .../runtime/threading/istreams_executor.hpp | 2 +- .../dev/threading/cpu_streams_executor.cpp | 18 +++++------ .../src/dev/threading/executor_manager.cpp | 2 +- src/plugins/intel_cpu/src/compiled_model.cpp | 32 ++++++++++++------- src/plugins/intel_cpu/src/config.h | 1 + src/plugins/intel_cpu/src/infer_request.cpp | 2 +- 6 files changed, 34 insertions(+), 23 deletions(-) diff --git a/src/inference/dev_api/openvino/runtime/threading/istreams_executor.hpp b/src/inference/dev_api/openvino/runtime/threading/istreams_executor.hpp index f501296f186..2af9c496e8c 100644 --- a/src/inference/dev_api/openvino/runtime/threading/istreams_executor.hpp +++ b/src/inference/dev_api/openvino/runtime/threading/istreams_executor.hpp @@ -209,7 +209,7 @@ public: return _threadBindingOffset; } int get_sub_streams() const { - return _sub_streams > 0 ? _sub_streams : _streams; + return _sub_streams; } StreamsMode get_sub_stream_mode() const { const auto proc_type_table = get_proc_type_table(); diff --git a/src/inference/src/dev/threading/cpu_streams_executor.cpp b/src/inference/src/dev/threading/cpu_streams_executor.cpp index bc5b3c6c75c..3e84a4179f3 100644 --- a/src/inference/src/dev/threading/cpu_streams_executor.cpp +++ b/src/inference/src/dev/threading/cpu_streams_executor.cpp @@ -66,12 +66,12 @@ struct CPUStreamsExecutor::Impl { } } _numaNodeId = - _impl->_config.get_sub_streams() - ? _impl->_usedNumaNodes.at((_streamId % _impl->_config.get_sub_streams()) / - ((_impl->_config.get_sub_streams() + _impl->_usedNumaNodes.size() - 1) / + _impl->_config.get_streams() + ? _impl->_usedNumaNodes.at((_streamId % _impl->_config.get_streams()) / + ((_impl->_config.get_streams() + _impl->_usedNumaNodes.size() - 1) / _impl->_usedNumaNodes.size())) : _impl->_usedNumaNodes.at(_streamId % _impl->_usedNumaNodes.size()); - // std::cout << "[ Stream ] " << _impl->_config.get_name() << " : " << _streamId << ", " << _impl->_config.get_sub_streams() << "\n"; + // std::cout << "[ Stream ] " << _impl->_config.get_name() << " : " << _streamId << ", " << _impl->_config.get_streams() << "\n"; #if OV_THREAD == OV_THREAD_TBB || OV_THREAD == OV_THREAD_TBB_AUTO if (is_cpu_map_available() && _impl->_config.get_streams_info_table().size() > 0) { init_stream(); @@ -174,7 +174,7 @@ struct CPUStreamsExecutor::Impl { int max_threads_per_core; StreamCreateType stream_type; const auto org_proc_type_table = get_org_proc_type_table(); - int streams_num = _impl->_config.get_sub_streams(); + int streams_num = _impl->_config.get_streams(); const auto stream_id = streams_num == 0 ? 0 : (_sub_stream_id >= 0 ? streams_num + _sub_stream_id : _streamId % streams_num); get_cur_stream_info(stream_id, @@ -324,7 +324,7 @@ struct CPUStreamsExecutor::Impl { this) { _exectorMgr = executor_manager(); auto numaNodes = get_available_numa_nodes(); - int streams_num = _config.get_sub_streams(); + int streams_num = _config.get_streams(); int sub_streams_num = 0;//_config.get_sub_streams(); // std::cout << "[ Impl ] " << _config.get_name() << " : " << streams_num << "\n"; if (streams_num != 0) { @@ -345,7 +345,7 @@ struct CPUStreamsExecutor::Impl { { std::unique_lock lock(_mutex); _queueCondVar.wait(lock, [&] { - std::cout << _config.get_name() << " addr: " << this << " in : " << streamId << "\n"; + // std::cout << _config.get_name() << " addr: " << this << " in : " << streamId << "\n"; return !_taskQueue.empty() || (stopped = _isStopped); }); if (!_taskQueue.empty()) { @@ -488,7 +488,7 @@ struct CPUStreamsExecutor::Impl { _readCondVar.notify_all(); } else if (msg_type == CALL_BACK) { // CALL_BACK count++; - std::cout << "server_wait CALL_BACK: " << count << "/" << streams_num << "\n"; + // std::cout << "server_wait CALL_BACK: " << count << "/" << streams_num << "\n"; if (count == streams_num) { _inferCondVar.notify_one(); count = 0; @@ -627,7 +627,7 @@ void CPUStreamsExecutor::execute(Task task) { } void CPUStreamsExecutor::run(Task task) { - if (0 == _impl->_config.get_sub_streams()) { + if (0 == _impl->_config.get_streams()) { _impl->Defer(std::move(task)); } else { _impl->Enqueue(std::move(task)); diff --git a/src/inference/src/dev/threading/executor_manager.cpp b/src/inference/src/dev/threading/executor_manager.cpp index 96e496a1382..efd9e9835ee 100644 --- a/src/inference/src/dev/threading/executor_manager.cpp +++ b/src/inference/src/dev/threading/executor_manager.cpp @@ -129,7 +129,7 @@ std::shared_ptr ExecutorManagerImpl::get_idle_c std::lock_guard guard(streamExecutorMutex); for (auto& it : cpuStreamsExecutors) { const auto& executor = it.second; - if (executor.use_count() != 1 && config.get_name() != "CPUStreamsExecutor") + if (executor.use_count() != 1) continue; auto& executorConfig = it.first; diff --git a/src/plugins/intel_cpu/src/compiled_model.cpp b/src/plugins/intel_cpu/src/compiled_model.cpp index 73f639b6f60..a595fab607b 100644 --- a/src/plugins/intel_cpu/src/compiled_model.cpp +++ b/src/plugins/intel_cpu/src/compiled_model.cpp @@ -18,6 +18,7 @@ #include "openvino/util/common_util.hpp" #include "openvino/runtime/threading/cpu_streams_executor.hpp" #include "transformations/utils/utils.hpp" +#include "openvino/runtime/threading/cpu_streams_info.hpp" #include "cpu/x64/cpu_isa_traits.hpp" #include @@ -56,24 +57,22 @@ CompiledModel::CompiledModel(const std::shared_ptr& model, OPENVINO_THROW("Unable to get API version. Core is unavailable"); ov::threading::IStreamsExecutor::Ptr stream_executor = nullptr; + IStreamsExecutor::Config executor_confg; if (cfg.exclusiveAsyncRequests) { // special case when all InferRequests are muxed into a single queue m_task_executor = m_plugin->get_executor_manager()->get_executor("CPU"); } else { - IStreamsExecutor::Config executor_confg; if (m_cfg.enableSubStreams) { executor_confg = IStreamsExecutor::Config{"CPUMainStreamExecutor", 1, 1, - IStreamsExecutor::ThreadBindingType::NONE, - 1, - 0, - 1, - IStreamsExecutor::Config::PreferredCoreType::ANY, - {}, + ov::hint::SchedulingCoreType::ANY_CORE, + false, true}; } else { - executor_confg = std::move(m_cfg.streamExecutorConfig); + executor_confg = m_cfg.subStreamExecConfig.get_name() != "StreamExecutor" + ? std::move(m_cfg.subStreamExecConfig) + : std::move(m_cfg.streamExecutorConfig); } stream_executor = m_plugin->get_executor_manager()->get_idle_cpu_streams_executor(executor_confg); m_task_executor = stream_executor; @@ -90,12 +89,11 @@ CompiledModel::CompiledModel(const std::shared_ptr& model, if (m_callback_executor) set_callback_executor(m_callback_executor); - int streams = m_cfg.enableSubStreams ? 1 : std::max(1, m_cfg.streamExecutorConfig.get_sub_streams()); - std::cout << "streams: " << streams << "\n"; + int streams = std::max(1, executor_confg.get_streams()); std::vector tasks; tasks.resize(streams); m_graphs.resize(streams); - if (m_cfg.streamExecutorConfig.get_sub_streams() != 0) { + if (executor_confg.get_streams() != 0) { auto all_graphs_ready = [&] { return std::all_of(m_graphs.begin(), m_graphs.end(), [&](Graph& graph) { return graph.IsReady(); @@ -123,6 +121,18 @@ CompiledModel::CompiledModel(const std::shared_ptr& model, if (m_cfg.enableSubStreams) { m_cfg.enableSubStreams = false; for (int i = 0; i < m_cfg.streamExecutorConfig.get_sub_streams(); i++) { + auto streams_info_table = m_cfg.streamExecutorConfig.get_streams_info_table(); + std::vector> info_table; + info_table.push_back(streams_info_table[i + 1]); + info_table[0][NUMBER_OF_STREAMS] = 1; + auto streamExecutorConfig = IStreamsExecutor::Config{"CPUStreamsExecutor", + 1, + 1, + ov::hint::SchedulingCoreType::ANY_CORE, + false, + true, + info_table}; + m_cfg.subStreamExecConfig = std::move(streamExecutorConfig); m_sub_compilemodels.push_back(std::make_shared(model, plugin, m_cfg, loaded_from_cache)); } } diff --git a/src/plugins/intel_cpu/src/config.h b/src/plugins/intel_cpu/src/config.h index 7e2ccfd2e24..1ea21ad1635 100644 --- a/src/plugins/intel_cpu/src/config.h +++ b/src/plugins/intel_cpu/src/config.h @@ -58,6 +58,7 @@ struct Config { size_t rtCacheCapacity = 0ul; #endif ov::threading::IStreamsExecutor::Config streamExecutorConfig; + ov::threading::IStreamsExecutor::Config subStreamExecConfig; int streams = 1; bool streamsChanged = false; int threads = 0; diff --git a/src/plugins/intel_cpu/src/infer_request.cpp b/src/plugins/intel_cpu/src/infer_request.cpp index ee22d9cbb75..5d78e13d396 100644 --- a/src/plugins/intel_cpu/src/infer_request.cpp +++ b/src/plugins/intel_cpu/src/infer_request.cpp @@ -665,7 +665,7 @@ void SyncInferRequest::sub_streams_infer() { auto streams_executor = m_graph->context->getCPUStreamExecutor(); for (size_t i = 0; i < requests_num; i++) { - std::cout << "start_async : " << i << "\n"; + // std::cout << "start_async : " << i << "\n"; requests[i]->start_async(); } } From 62bfd4662939087a4fb238fdc0fe334c7257114c Mon Sep 17 00:00:00 2001 From: sunxiaoxia2022 Date: Tue, 7 May 2024 11:31:31 +0800 Subject: [PATCH 03/18] add api of getting rank --- .../threading/cpu_streams_executor.hpp | 2 ++ .../runtime/threading/istreams_executor.hpp | 21 +++++++++++++++---- .../dev/threading/cpu_streams_executor.cpp | 13 ++++++++---- src/plugins/intel_cpu/src/compiled_model.cpp | 3 ++- src/plugins/intel_cpu/src/config.h | 1 + .../intel_cpu/src/cpu_streams_calculation.cpp | 2 +- src/plugins/intel_cpu/src/infer_request.cpp | 3 --- 7 files changed, 32 insertions(+), 13 deletions(-) diff --git a/src/inference/dev_api/openvino/runtime/threading/cpu_streams_executor.hpp b/src/inference/dev_api/openvino/runtime/threading/cpu_streams_executor.hpp index f0b68be6ee6..a7372073bd4 100644 --- a/src/inference/dev_api/openvino/runtime/threading/cpu_streams_executor.hpp +++ b/src/inference/dev_api/openvino/runtime/threading/cpu_streams_executor.hpp @@ -53,6 +53,8 @@ public: int get_socket_id() override; + std::vector get_rank() override; + void run_sub_stream(Task task, int id) override; void send_message(MessageInfo msg_info); diff --git a/src/inference/dev_api/openvino/runtime/threading/istreams_executor.hpp b/src/inference/dev_api/openvino/runtime/threading/istreams_executor.hpp index 2af9c496e8c..e96acbe8964 100644 --- a/src/inference/dev_api/openvino/runtime/threading/istreams_executor.hpp +++ b/src/inference/dev_api/openvino/runtime/threading/istreams_executor.hpp @@ -46,8 +46,7 @@ public: struct MessageInfo{ MsgType msg_type; - int rank; - int data; + std::vector rank; void* buf; Task task; }; @@ -109,6 +108,7 @@ public: std::vector> _streams_info_table = {}; std::vector> _stream_processor_ids; int _sub_streams = 0; + std::vector _rank = {}; /** * @brief Get and reserve cpu ids based on configuration and hardware information, @@ -138,6 +138,7 @@ public: * @param[in] cpu_reservation @copybrief Config::_cpu_reservation * @param[in] cpu_pinning @copybrief Config::_cpu_pinning * @param[in] streams_info_table @copybrief Config::_streams_info_table + * @param[in] rank @copybrief Config::_rank */ Config(std::string name = "StreamsExecutor", int streams = 1, @@ -145,14 +146,16 @@ public: ov::hint::SchedulingCoreType thread_preferred_core_type = ov::hint::SchedulingCoreType::ANY_CORE, bool cpu_reservation = false, bool cpu_pinning = false, - std::vector> streams_info_table = {}) + std::vector> streams_info_table = {}, + std::vector rank = {}) : _name{name}, _streams{streams}, _threads_per_stream{threads_per_stream}, _thread_preferred_core_type(thread_preferred_core_type), _cpu_reservation{cpu_reservation}, _cpu_pinning{cpu_pinning}, - _streams_info_table{streams_info_table} { + _streams_info_table{streams_info_table}, + _rank{rank} { update_executor_config(); } @@ -211,6 +214,9 @@ public: int get_sub_streams() const { return _sub_streams; } + std::vector get_rank() const { + return _rank; + } StreamsMode get_sub_stream_mode() const { const auto proc_type_table = get_proc_type_table(); int sockets = proc_type_table.size() > 1 ? static_cast(proc_type_table.size()) - 1 : 1; @@ -263,6 +269,13 @@ public: */ virtual int get_socket_id() = 0; + /** + * @brief Return the rank of current stream + * Return {} when current stream has no rank + * @return Rank array, or throws exceptions if called not from stream thread + */ + virtual std::vector get_rank() = 0; + /** * @brief Execute the task in the current thread using streams executor configuration and constraints * @param task A task to start diff --git a/src/inference/src/dev/threading/cpu_streams_executor.cpp b/src/inference/src/dev/threading/cpu_streams_executor.cpp index 3e84a4179f3..e195cf699bc 100644 --- a/src/inference/src/dev/threading/cpu_streams_executor.cpp +++ b/src/inference/src/dev/threading/cpu_streams_executor.cpp @@ -154,8 +154,6 @@ struct CPUStreamsExecutor::Impl { _taskArena.reset(new custom::task_arena{concurrency}); _cpu_ids = stream_id < static_cast(stream_processors.size()) ? stream_processors[stream_id] : _cpu_ids; - std::cout << "stream_id: " << stream_id << " size: " << stream_processors.size() - << " cpu size: " << _cpu_ids.size() << " addr: " << _impl << " , " << _cpu_ids[0] << "\n"; if (_cpu_ids.size() > 0) { CpuSet processMask; int ncpus = 0; @@ -177,6 +175,7 @@ struct CPUStreamsExecutor::Impl { int streams_num = _impl->_config.get_streams(); const auto stream_id = streams_num == 0 ? 0 : (_sub_stream_id >= 0 ? streams_num + _sub_stream_id : _streamId % streams_num); + _rank = _impl->_config.get_rank(); get_cur_stream_info(stream_id, _impl->_config.get_cpu_pinning(), org_proc_type_table, @@ -204,6 +203,7 @@ struct CPUStreamsExecutor::Impl { int _socketId = 0; bool _execute = false; int _sub_stream_id = -1; + std::vector _rank; std::queue _taskQueue; #if OV_THREAD == OV_THREAD_TBB || OV_THREAD == OV_THREAD_TBB_AUTO std::unique_ptr _taskArena; @@ -469,7 +469,7 @@ struct CPUStreamsExecutor::Impl { _msgCondVar.wait(lock); } std::swap(_messageQueue, msgQueue); - // std::cout << "server_wait receive: " << msgQueue[0].msg_type << " data:" << msgQueue[0].data + // std::cout << "server_wait receive: " << msgQueue[0].msg_type << " rank:" << msgQueue[0].rank[0] // << " / " << msgQueue.size() << "\n"; } @@ -480,7 +480,7 @@ struct CPUStreamsExecutor::Impl { task(); } else if (msg_type == TP) { for (int i = 0; i < streams_num; i++) { - if (rec_info.data != i) { + if (rec_info.rank[0] != i) { std::lock_guard lock(_readMutex); _readQueue[i].push_back(rec_info); } @@ -574,6 +574,11 @@ int CPUStreamsExecutor::get_socket_id() { return stream->_socketId; } +std::vector CPUStreamsExecutor::get_rank() { + auto stream = _impl->_streams.local(); + return stream->_rank; +} + void CPUStreamsExecutor::send_message(MessageInfo msg_info) { _impl->send_message(msg_info); } diff --git a/src/plugins/intel_cpu/src/compiled_model.cpp b/src/plugins/intel_cpu/src/compiled_model.cpp index a595fab607b..9f97a0b750a 100644 --- a/src/plugins/intel_cpu/src/compiled_model.cpp +++ b/src/plugins/intel_cpu/src/compiled_model.cpp @@ -131,7 +131,8 @@ CompiledModel::CompiledModel(const std::shared_ptr& model, ov::hint::SchedulingCoreType::ANY_CORE, false, true, - info_table}; + info_table, + m_cfg.streamsRankTable[i]}; m_cfg.subStreamExecConfig = std::move(streamExecutorConfig); m_sub_compilemodels.push_back(std::make_shared(model, plugin, m_cfg, loaded_from_cache)); } diff --git a/src/plugins/intel_cpu/src/config.h b/src/plugins/intel_cpu/src/config.h index 1ea21ad1635..9728ddbb2ee 100644 --- a/src/plugins/intel_cpu/src/config.h +++ b/src/plugins/intel_cpu/src/config.h @@ -65,6 +65,7 @@ struct Config { int threadsPerStream = 0; ov::threading::IStreamsExecutor::ThreadBindingType threadBindingType = ov::threading::IStreamsExecutor::ThreadBindingType::NONE; ov::hint::PerformanceMode hintPerfMode = ov::hint::PerformanceMode::LATENCY; + std::vector> streamsRankTable; bool changedHintPerfMode = false; ov::log::Level logLevel = ov::log::Level::NO; uint32_t hintNumRequests = 0; diff --git a/src/plugins/intel_cpu/src/cpu_streams_calculation.cpp b/src/plugins/intel_cpu/src/cpu_streams_calculation.cpp index c8e8fd611fe..d8ec38bbed1 100644 --- a/src/plugins/intel_cpu/src/cpu_streams_calculation.cpp +++ b/src/plugins/intel_cpu/src/cpu_streams_calculation.cpp @@ -655,7 +655,7 @@ std::vector> generate_stream_info(const int streams, // streams_info_table = {{1, 1, 32, -1, -1}, {-1, 1, 16, 0, 0}, {-1, 1, 16, 1, 1}}; if (config.modelDistributionPolicy.find(ov::hint::ModelDistributionPolicy::TENSOR_PARALLEL) != config.modelDistributionPolicy.end()) { - auto streams_rank_table = + config.streamsRankTable = get_streams_rank_table(streams_info_table, config.streamsRankLevel, config.numSubStreams); config.enableSubStreams = true; } diff --git a/src/plugins/intel_cpu/src/infer_request.cpp b/src/plugins/intel_cpu/src/infer_request.cpp index 5d78e13d396..bdbe830a9af 100644 --- a/src/plugins/intel_cpu/src/infer_request.cpp +++ b/src/plugins/intel_cpu/src/infer_request.cpp @@ -629,9 +629,6 @@ void SyncInferRequest::sub_streams_infer() { auto requests = m_asyncRequest->getSubInferRequest(); auto inputs = m_asyncRequest->get_inputs(); auto outputs = m_asyncRequest->get_outputs(); - // auto sub_models = m_compiled_model->get_sub_compilemodles(); - // auto graphLock = sub_models[0]->get_graph(); - // m_graph = &(graphLock._graph); auto streams_executor = m_graph->context->getCPUStreamExecutor(); size_t requests_num = requests.size(); size_t requests_count = 0; From f664f2a14b034403afa8a7a90d04855853e2163a Mon Sep 17 00:00:00 2001 From: sunxiaoxia2022 Date: Wed, 8 May 2024 09:10:46 +0800 Subject: [PATCH 04/18] move apis of send_message into new file --- .../runtime/threading/cpu_message.hpp | 63 +++++++++ .../threading/cpu_streams_executor.hpp | 8 -- .../src/dev/threading/cpu_message.cpp | 126 ++++++++++++++++++ .../dev/threading/cpu_streams_executor.cpp | 103 -------------- src/plugins/intel_cpu/src/compiled_model.cpp | 2 +- src/plugins/intel_cpu/src/infer_request.cpp | 38 +++--- src/plugins/intel_cpu/src/plugin.cpp | 3 + 7 files changed, 210 insertions(+), 133 deletions(-) create mode 100644 src/inference/dev_api/openvino/runtime/threading/cpu_message.hpp create mode 100644 src/inference/src/dev/threading/cpu_message.cpp diff --git a/src/inference/dev_api/openvino/runtime/threading/cpu_message.hpp b/src/inference/dev_api/openvino/runtime/threading/cpu_message.hpp new file mode 100644 index 00000000000..78d6606722d --- /dev/null +++ b/src/inference/dev_api/openvino/runtime/threading/cpu_message.hpp @@ -0,0 +1,63 @@ +// Copyright (C) 2018-2024 Intel Corporation +// SPDX-License-Identifier: Apache-2.0 +// + +/** + * @brief OpenVINO Runtime Executor Manager + * @file openvino/runtime/threading/executor_manager.hpp + */ + +#pragma once + +#include +#include +#include +#include +#include +#include +#include +#include + +#include "openvino/runtime/common.hpp" +#include "openvino/runtime/threading/istreams_executor.hpp" +#include "openvino/runtime/threading/itask_executor.hpp" + +namespace ov { + +namespace threading { +enum MsgType { TP, START_INFER, CALL_BACK }; + +struct MessageInfo { + MsgType msg_type; + std::vector rank; + void* buf; + Task task; +}; +class OPENVINO_RUNTIME_API MessageManage { +public: + MessageManage(); + void send_message(MessageInfo msg_info); + + std::vector wait_message(int cur_rank, int streams_num); + + void infer_wait(); + + void server_wait(int streams_num); + + ~MessageManage(); +private: + std::thread _serverThread; + bool _isServerStopped = false; + std::vector _messageQueue; + std::vector> _readQueue; + std::mutex _msgMutex; + std::mutex _readMutex; + std::mutex _inferMutex; + std::condition_variable _msgCondVar; + std::condition_variable _readCondVar; + std::condition_variable _inferCondVar; +}; + +OPENVINO_RUNTIME_API std::shared_ptr message_manager(); +} // namespace threading +} // namespace ov \ No newline at end of file diff --git a/src/inference/dev_api/openvino/runtime/threading/cpu_streams_executor.hpp b/src/inference/dev_api/openvino/runtime/threading/cpu_streams_executor.hpp index a7372073bd4..1f694fe9b70 100644 --- a/src/inference/dev_api/openvino/runtime/threading/cpu_streams_executor.hpp +++ b/src/inference/dev_api/openvino/runtime/threading/cpu_streams_executor.hpp @@ -57,14 +57,6 @@ public: void run_sub_stream(Task task, int id) override; - void send_message(MessageInfo msg_info); - - void wait_message(); - - void infer_wait(); - - void server_wait(int streams_num); - private: struct Impl; std::unique_ptr _impl; diff --git a/src/inference/src/dev/threading/cpu_message.cpp b/src/inference/src/dev/threading/cpu_message.cpp new file mode 100644 index 00000000000..256d422a1e3 --- /dev/null +++ b/src/inference/src/dev/threading/cpu_message.cpp @@ -0,0 +1,126 @@ +// Copyright (C) 2018-2024 Intel Corporation +// SPDX-License-Identifier: Apache-2.0 +// +#include "openvino/runtime/threading/cpu_message.hpp" +#include +#include +#include +#include +#include +#include + +namespace ov { +namespace threading { + +MessageManage::MessageManage() {} + +void MessageManage::send_message(MessageInfo msg_info) { + { + std::lock_guard lock(_msgMutex); + _messageQueue.push_back(msg_info); + // std::cout << "send : " << msg_info.msg_type << ", " << msg_info.rank.size() << "\n"; + } + _msgCondVar.notify_all(); +} + +std::vector MessageManage::wait_message(int cur_rank, int streams_num) { + std::vector messages_total; + std::unique_lock lock(_readMutex); + _readCondVar.wait(lock, [&] { + // std::cout << "wait_" << cur_rank << " : " << _readQueue[cur_rank].size() << " / " << streams_num << "\n"; + return _readQueue[cur_rank].size() >= streams_num; + }); + std::swap(_readQueue[cur_rank], messages_total); + // std::cout << "wait_" << cur_rank << " " << _readQueue[cur_rank].size() << " end\n"; + return messages_total; +} + +void MessageManage::infer_wait() { + std::unique_lock lock(_inferMutex); + _inferCondVar.wait(lock); +} + +void MessageManage::server_wait(int streams_num) { + if (!_serverThread.joinable()) { + _readQueue.assign(streams_num, std::vector()); + MsgType msg_type; + _serverThread = std::thread([&, streams_num]() { + int count = 0; + while (!_isServerStopped) { + std::vector msgQueue; + { + // std::cout << "server_wait ........\n"; + std::unique_lock lock(_msgMutex); + while (_messageQueue.empty()) { + _msgCondVar.wait(lock); + } + std::swap(_messageQueue, msgQueue); + // std::cout << "server_wait receive: " << msgQueue[0].msg_type << " rank:" << msgQueue[0].rank.size() + // << " / " << msgQueue.size() << "\n"; + } + + for (auto rec_info : msgQueue) { + msg_type = rec_info.msg_type; + if (msg_type == START_INFER) { + Task task = std::move(rec_info.task); + task(); + } else if (msg_type == TP) { + for (int i = 0; i < streams_num; i++) { + std::lock_guard lock(_readMutex); + _readQueue[i].push_back(rec_info); + } + _readCondVar.notify_all(); + } else if (msg_type == CALL_BACK) { // CALL_BACK + count++; + // std::cout << "server_wait CALL_BACK: " << count << "/" << streams_num << "\n"; + if (count == streams_num) { + _inferCondVar.notify_one(); + count = 0; + } + } + } + } + std::cout << "-------- server_wait end ---------\n"; + }); + } +} + +MessageManage::~MessageManage() { + _isServerStopped = true; + _msgCondVar.notify_one(); + if (_serverThread.joinable()) { + _serverThread.join(); + } +} + +namespace { + +class MessageManageHolder { + std::mutex _mutex; + std::weak_ptr _manager; + +public: + MessageManageHolder(const MessageManageHolder&) = delete; + MessageManageHolder& operator=(const MessageManageHolder&) = delete; + + MessageManageHolder() = default; + + std::shared_ptr get() { + std::lock_guard lock(_mutex); + auto manager = _manager.lock(); + if (!manager) { + _manager = manager = std::make_shared(); + } + return manager; + } +}; + +} // namespace + +std::shared_ptr message_manager(){ + static MessageManageHolder message_manage; + return message_manage.get(); +} + +} // namespace threading +} // namespace ov \ No newline at end of file diff --git a/src/inference/src/dev/threading/cpu_streams_executor.cpp b/src/inference/src/dev/threading/cpu_streams_executor.cpp index e195cf699bc..13964f1b60e 100644 --- a/src/inference/src/dev/threading/cpu_streams_executor.cpp +++ b/src/inference/src/dev/threading/cpu_streams_executor.cpp @@ -428,78 +428,6 @@ struct CPUStreamsExecutor::Impl { } } - void send_message(MessageInfo msg_info) { - { - std::lock_guard lock(_msgMutex); - _messageQueue.push_back(msg_info); - // std::cout << "send_" << _streamId << " : " << msg_info.msg_type << "\n"; - } - _msgCondVar.notify_all(); - } - - void wait_message() { - std::unique_lock lock(_readMutex); - _readCondVar.wait(lock, [&] { - std::cout << "wait_" << _streamId << " : " << _readQueue[_streamId].size() << " / " - << _config.get_sub_streams() - 1 << "\n"; - return _readQueue[_streamId].size() >= _config.get_sub_streams() - 1; - }); - std::cout << "wait_" << _streamId << " end\n"; - } - - void infer_wait() { - std::unique_lock lock(_inferMutex); - // std::cout << "infer_wait ......\n"; - _inferCondVar.wait(lock); - } - - void server_wait(int streams_num) { - if (!_serverThread.joinable()) { - _messageQueue.clear(); - _readQueue.assign(streams_num, std::vector()); - MsgType msg_type; - _serverThread = std::thread([&, streams_num]() { - int count = 0; - while (!_isServerStopped) { - std::vector msgQueue; - { - // std::cout << "server_wait ........\n"; - std::unique_lock lock(_msgMutex); - while (_messageQueue.empty()) { - _msgCondVar.wait(lock); - } - std::swap(_messageQueue, msgQueue); - // std::cout << "server_wait receive: " << msgQueue[0].msg_type << " rank:" << msgQueue[0].rank[0] - // << " / " << msgQueue.size() << "\n"; - } - - for (auto rec_info : msgQueue) { - msg_type = rec_info.msg_type; - if (msg_type == START_INFER) { - Task task = std::move(rec_info.task); - task(); - } else if (msg_type == TP) { - for (int i = 0; i < streams_num; i++) { - if (rec_info.rank[0] != i) { - std::lock_guard lock(_readMutex); - _readQueue[i].push_back(rec_info); - } - } - _readCondVar.notify_all(); - } else if (msg_type == CALL_BACK) { // CALL_BACK - count++; - // std::cout << "server_wait CALL_BACK: " << count << "/" << streams_num << "\n"; - if (count == streams_num) { - _inferCondVar.notify_one(); - count = 0; - } - } - } - } - }); - } - } - struct SubQueue { std::mutex _subMutex; std::condition_variable _subQueueCondVar; @@ -538,24 +466,14 @@ struct CPUStreamsExecutor::Impl { int _subStreamsNum = 0; std::vector _threads; std::vector _subThreads; - std::thread _serverThread; std::mutex _mutex; std::condition_variable _queueCondVar; std::queue _taskQueue; bool _isStopped = false; - bool _isServerStopped = false; std::vector> _subTaskThread; std::vector _usedNumaNodes; CustomThreadLocal _streams; std::shared_ptr _exectorMgr; - std::vector _messageQueue; - std::vector> _readQueue; - std::mutex _msgMutex; - std::mutex _readMutex; - std::mutex _inferMutex; - std::condition_variable _msgCondVar; - std::condition_variable _readCondVar; - std::condition_variable _inferCondVar; bool _isExit = false; }; @@ -579,22 +497,6 @@ std::vector CPUStreamsExecutor::get_rank() { return stream->_rank; } -void CPUStreamsExecutor::send_message(MessageInfo msg_info) { - _impl->send_message(msg_info); -} - -void CPUStreamsExecutor::wait_message() { - _impl->wait_message(); -} - -void CPUStreamsExecutor::infer_wait() { - _impl->infer_wait(); -} - -void CPUStreamsExecutor::server_wait(int streams_num) { - _impl->server_wait(streams_num); -} - CPUStreamsExecutor::CPUStreamsExecutor(const IStreamsExecutor::Config& config) : _impl{new Impl{config}} {} CPUStreamsExecutor::~CPUStreamsExecutor() { @@ -620,11 +522,6 @@ CPUStreamsExecutor::~CPUStreamsExecutor() { thread.join(); } } - _impl->_isServerStopped = true; - _impl->_msgCondVar.notify_one(); - if (_impl->_serverThread.joinable()) { - _impl->_serverThread.join(); - } } void CPUStreamsExecutor::execute(Task task) { diff --git a/src/plugins/intel_cpu/src/compiled_model.cpp b/src/plugins/intel_cpu/src/compiled_model.cpp index 9f97a0b750a..5f8f727762f 100644 --- a/src/plugins/intel_cpu/src/compiled_model.cpp +++ b/src/plugins/intel_cpu/src/compiled_model.cpp @@ -125,7 +125,7 @@ CompiledModel::CompiledModel(const std::shared_ptr& model, std::vector> info_table; info_table.push_back(streams_info_table[i + 1]); info_table[0][NUMBER_OF_STREAMS] = 1; - auto streamExecutorConfig = IStreamsExecutor::Config{"CPUStreamsExecutor", + auto streamExecutorConfig = IStreamsExecutor::Config{"CPUSubStreamsExecutor" + std::to_string(i), 1, 1, ov::hint::SchedulingCoreType::ANY_CORE, diff --git a/src/plugins/intel_cpu/src/infer_request.cpp b/src/plugins/intel_cpu/src/infer_request.cpp index bdbe830a9af..ba890d965d5 100644 --- a/src/plugins/intel_cpu/src/infer_request.cpp +++ b/src/plugins/intel_cpu/src/infer_request.cpp @@ -18,6 +18,7 @@ #include "proxy_mem_mgr.h" #include "utils/general_utils.h" #include "utils/ngraph_utils.hpp" +#include "openvino/runtime/threading/cpu_message.hpp" using OvString = ov::element_type_traits::value_type; @@ -110,22 +111,24 @@ void SyncInferRequest::infer() { m_graph = &(graphLock._graph); auto streams_executor = m_graph->context->getCPUStreamExecutor(); + throw_if_canceled(); auto requests = m_asyncRequest->getSubInferRequest(); // std::cout << "[ infer ] " << requests.size() << "\n"; if (requests.size() > 0) { - streams_executor->server_wait(requests.size()); - ov::threading::IStreamsExecutor::MessageInfo msg_info; - msg_info.msg_type = ov::threading::IStreamsExecutor::MsgType::START_INFER; + auto message = ov::threading::message_manager(); + message->server_wait(requests.size()); + ov::threading::MessageInfo msg_info; + msg_info.msg_type = ov::threading::MsgType::START_INFER; ov::threading::Task task = [&] { SyncInferRequest::sub_streams_infer(); }; msg_info.task = std::move(task); - streams_executor->send_message(msg_info); - streams_executor->infer_wait(); + message->send_message(msg_info); + message->infer_wait(); + // std::cout << "------ infer end -----\n"; return; } - throw_if_canceled(); convert_batched_tensors(); if (m_batched_tensors.size() > 0) { // batched_tensors will be updated for each infer, external_ptr should be update together @@ -629,7 +632,7 @@ void SyncInferRequest::sub_streams_infer() { auto requests = m_asyncRequest->getSubInferRequest(); auto inputs = m_asyncRequest->get_inputs(); auto outputs = m_asyncRequest->get_outputs(); - auto streams_executor = m_graph->context->getCPUStreamExecutor(); + auto message = ov::threading::message_manager(); size_t requests_num = requests.size(); size_t requests_count = 0; // std::cout << "[ sub_streams_infer ] inputs: " << inputs.size() << " requests: " << requests_num << "\n"; @@ -648,23 +651,16 @@ void SyncInferRequest::sub_streams_infer() { requests[i]->set_tensor(input.first, input.second); } - requests[i]->set_callback([i, requests, streams_executor](const std::exception_ptr& ptr) { + requests[i]->set_callback([i, requests, message](const std::exception_ptr& ptr) { // std::cout << "set_callback------ " << i << "\n"; - ov::threading::IStreamsExecutor::MessageInfo msg_info; - msg_info.msg_type = ov::threading::IStreamsExecutor::MsgType::CALL_BACK; - streams_executor->send_message(msg_info); + ov::threading::MessageInfo msg_info; + msg_info.msg_type = ov::threading::MsgType::CALL_BACK; + message->send_message(msg_info); }); } - { - auto sub_models = m_compiled_model->get_sub_compilemodles(); - auto graphLock = sub_models[0]->get_graph(); - m_graph = &(graphLock._graph); - auto streams_executor = m_graph->context->getCPUStreamExecutor(); - - for (size_t i = 0; i < requests_num; i++) { - // std::cout << "start_async : " << i << "\n"; - requests[i]->start_async(); - } + for (size_t i = 0; i < requests_num; i++) { + // std::cout << "start_async : " << i << "\n"; + requests[i]->start_async(); } } } diff --git a/src/plugins/intel_cpu/src/plugin.cpp b/src/plugins/intel_cpu/src/plugin.cpp index e4540ea4344..0aa82a94d2e 100644 --- a/src/plugins/intel_cpu/src/plugin.cpp +++ b/src/plugins/intel_cpu/src/plugin.cpp @@ -136,6 +136,9 @@ Plugin::Plugin() : deviceFullName(getDeviceFullName()), specialSetup(new CPUSpec Plugin::~Plugin() { executor_manager()->clear("CPU"); executor_manager()->clear("CPUStreamsExecutor"); + executor_manager()->clear("CPUSubStreamsExecutor0"); + executor_manager()->clear("CPUSubStreamsExecutor1"); + executor_manager()->clear("CPUMainStreamExecutor"); executor_manager()->clear("CPUCallbackExecutor"); } From 670a59fe63f6f84782275bbeb3287b80d04dba32 Mon Sep 17 00:00:00 2001 From: sunxiaoxia2022 Date: Thu, 9 May 2024 18:12:19 +0800 Subject: [PATCH 05/18] fix probabilistic no response of wait_message --- src/inference/src/dev/threading/cpu_message.cpp | 15 ++++++++++++++- 1 file changed, 14 insertions(+), 1 deletion(-) diff --git a/src/inference/src/dev/threading/cpu_message.cpp b/src/inference/src/dev/threading/cpu_message.cpp index 256d422a1e3..9d595e188b9 100644 --- a/src/inference/src/dev/threading/cpu_message.cpp +++ b/src/inference/src/dev/threading/cpu_message.cpp @@ -18,7 +18,7 @@ void MessageManage::send_message(MessageInfo msg_info) { { std::lock_guard lock(_msgMutex); _messageQueue.push_back(msg_info); - // std::cout << "send : " << msg_info.msg_type << ", " << msg_info.rank.size() << "\n"; + // std::cout << "send : " << msg_info.msg_type << ", " << (msg_info.rank.size() > 0 ? msg_info.rank[0] : -1) << "\n"; } _msgCondVar.notify_all(); } @@ -65,6 +65,19 @@ void MessageManage::server_wait(int streams_num) { Task task = std::move(rec_info.task); task(); } else if (msg_type == TP) { + // Resend _readQueue that failed last time + bool stop = false; + while (!stop) { + stop = true; + for (int i = 0; i < streams_num; i++) { + if (_readQueue[i].size() == streams_num) { + stop = false; + } + } + if (!stop) { + _readCondVar.notify_all(); + } + } for (int i = 0; i < streams_num; i++) { std::lock_guard lock(_readMutex); _readQueue[i].push_back(rec_info); From 155030320adafc3a6b9da537c4e790191ee494e1 Mon Sep 17 00:00:00 2001 From: sunxiaoxia2022 Date: Mon, 13 May 2024 22:04:44 +0800 Subject: [PATCH 06/18] fix llm model infer failed in second test with genai --- .../intel_cpu/src/cpu_streams_calculation.cpp | 2 +- src/plugins/intel_cpu/src/infer_request.cpp | 25 +++++++++++++------ src/plugins/intel_cpu/src/plugin.cpp | 10 ++++---- 3 files changed, 23 insertions(+), 14 deletions(-) diff --git a/src/plugins/intel_cpu/src/cpu_streams_calculation.cpp b/src/plugins/intel_cpu/src/cpu_streams_calculation.cpp index d8ec38bbed1..20af0f2fa84 100644 --- a/src/plugins/intel_cpu/src/cpu_streams_calculation.cpp +++ b/src/plugins/intel_cpu/src/cpu_streams_calculation.cpp @@ -652,7 +652,7 @@ std::vector> generate_stream_info(const int streams, ov::util::to_string(config.hintPerfMode), config.modelDistributionPolicy, proc_type_table); - // streams_info_table = {{1, 1, 32, -1, -1}, {-1, 1, 16, 0, 0}, {-1, 1, 16, 1, 1}}; + // streams_info_table = {{1, 1, 56, 1, 1}, {-1, 1, 28, 1, 1}, {-1, 1, 28, 0, 0}}; if (config.modelDistributionPolicy.find(ov::hint::ModelDistributionPolicy::TENSOR_PARALLEL) != config.modelDistributionPolicy.end()) { config.streamsRankTable = diff --git a/src/plugins/intel_cpu/src/infer_request.cpp b/src/plugins/intel_cpu/src/infer_request.cpp index ba890d965d5..29862756cfb 100644 --- a/src/plugins/intel_cpu/src/infer_request.cpp +++ b/src/plugins/intel_cpu/src/infer_request.cpp @@ -330,6 +330,15 @@ void SyncInferRequest::change_default_ptr() { } std::vector> SyncInferRequest::query_state() const { + auto requests = m_asyncRequest->getSubInferRequest(); + if (requests.size() > 0) { + std::vector> states; + for (auto request : requests) { + auto cur = request->query_state(); + states.insert(states.end(), cur.begin(), cur.end()); + } + return states; + } return {m_memory_states.begin(), m_memory_states.end()}; } @@ -638,17 +647,17 @@ void SyncInferRequest::sub_streams_infer() { // std::cout << "[ sub_streams_infer ] inputs: " << inputs.size() << " requests: " << requests_num << "\n"; if (requests.size() > 0) { - for (const auto& input : inputs) { - auto tensor = m_asyncRequest->get_tensor(input); - input_tensors.insert({input, tensor}); - } for (const auto& output : outputs) { - auto tensor = requests[0]->get_tensor(output); - m_asyncRequest->set_tensor(output, tensor); + auto main_tensor = get_tensor(output); + if (main_tensor->get_size() == 0) { + auto tensor = requests[0]->get_tensor(output); + set_tensor(output, tensor); + } } for (size_t i = 0; i < requests_num; i++) { - for (auto& input : input_tensors) { - requests[i]->set_tensor(input.first, input.second); + for (auto& input : inputs) { + auto tensor = m_asyncRequest->get_tensor(input); + requests[i]->set_tensor(input, tensor); } requests[i]->set_callback([i, requests, message](const std::exception_ptr& ptr) { diff --git a/src/plugins/intel_cpu/src/plugin.cpp b/src/plugins/intel_cpu/src/plugin.cpp index 0aa82a94d2e..e2aa4bbcbcf 100644 --- a/src/plugins/intel_cpu/src/plugin.cpp +++ b/src/plugins/intel_cpu/src/plugin.cpp @@ -284,11 +284,11 @@ std::shared_ptr Plugin::compile_model(const std::shared_ptr< conf.readProperties(config, modelType); calculate_streams(conf, cloned_model); - if (conf.streamExecutorConfig.get_sub_stream_mode() == - IStreamsExecutor::Config::StreamsMode::SUB_STREAMS_FOR_SOCKET) { - int num_sub_streams = conf.streamExecutorConfig.get_sub_streams(); - transformations.SetSubStreamNum(num_sub_streams); - } + // if (conf.streamExecutorConfig.get_sub_stream_mode() == + // IStreamsExecutor::Config::StreamsMode::SUB_STREAMS_FOR_SOCKET) { + // int num_sub_streams = conf.streamExecutorConfig.get_sub_streams(); + // transformations.SetSubStreamNum(num_sub_streams); + // } transformations.PostLpt(); transformations.Snippets(); From feb5d0f855641750a413e28e09475693bab1ffbc Mon Sep 17 00:00:00 2001 From: sunxiaoxia2022 Date: Tue, 14 May 2024 15:44:33 +0800 Subject: [PATCH 07/18] fix output is unstable --- src/plugins/intel_cpu/src/infer_request.cpp | 7 ++----- 1 file changed, 2 insertions(+), 5 deletions(-) diff --git a/src/plugins/intel_cpu/src/infer_request.cpp b/src/plugins/intel_cpu/src/infer_request.cpp index 29862756cfb..76ae6756a43 100644 --- a/src/plugins/intel_cpu/src/infer_request.cpp +++ b/src/plugins/intel_cpu/src/infer_request.cpp @@ -648,11 +648,8 @@ void SyncInferRequest::sub_streams_infer() { if (requests.size() > 0) { for (const auto& output : outputs) { - auto main_tensor = get_tensor(output); - if (main_tensor->get_size() == 0) { - auto tensor = requests[0]->get_tensor(output); - set_tensor(output, tensor); - } + auto tensor = requests[0]->get_tensor(output); + set_tensor(output, tensor); } for (size_t i = 0; i < requests_num; i++) { for (auto& input : inputs) { From fe496d8b61f6b8e63301f6a46906f82ebc598c45 Mon Sep 17 00:00:00 2001 From: sunxiaoxia2022 Date: Wed, 15 May 2024 16:53:43 +0800 Subject: [PATCH 08/18] move sub_compilemodels and sub_infer_requests to MessageManage --- .../openvino/runtime/iasync_infer_request.hpp | 2 - .../openvino/runtime/icompiled_model.hpp | 3 -- .../intel_cpu/src/async_infer_request.cpp | 10 ++--- .../intel_cpu/src/async_infer_request.h | 9 ++-- src/plugins/intel_cpu/src/compiled_model.cpp | 26 ++++++----- src/plugins/intel_cpu/src/compiled_model.h | 7 +-- .../intel_cpu/src}/cpu_message.cpp | 43 ++++++++++++++++--- .../intel_cpu/src}/cpu_message.hpp | 16 ++++++- src/plugins/intel_cpu/src/infer_request.cpp | 39 ++++++++++------- src/plugins/intel_cpu/src/plugin.cpp | 13 +++--- src/plugins/intel_cpu/src/plugin.h | 3 ++ 11 files changed, 109 insertions(+), 62 deletions(-) rename src/{inference/src/dev/threading => plugins/intel_cpu/src}/cpu_message.cpp (73%) rename src/{inference/dev_api/openvino/runtime/threading => plugins/intel_cpu/src}/cpu_message.hpp (66%) diff --git a/src/inference/dev_api/openvino/runtime/iasync_infer_request.hpp b/src/inference/dev_api/openvino/runtime/iasync_infer_request.hpp index bfca199464b..14c6fa2d657 100644 --- a/src/inference/dev_api/openvino/runtime/iasync_infer_request.hpp +++ b/src/inference/dev_api/openvino/runtime/iasync_infer_request.hpp @@ -277,8 +277,6 @@ private: m_sync_callback_executor; //!< Used to run post inference callback in synchronous pipline mutable std::mutex m_mutex; std::function m_callback; - - std::vector> m_sub_infer_requests; }; } // namespace ov diff --git a/src/inference/dev_api/openvino/runtime/icompiled_model.hpp b/src/inference/dev_api/openvino/runtime/icompiled_model.hpp index 7efdee598b4..eca22b3b003 100644 --- a/src/inference/dev_api/openvino/runtime/icompiled_model.hpp +++ b/src/inference/dev_api/openvino/runtime/icompiled_model.hpp @@ -145,9 +145,6 @@ private: std::shared_ptr m_task_executor = nullptr; //!< Holds a task executor std::shared_ptr m_callback_executor = nullptr; //!< Holds a callback executor - bool m_subCompileModel = true; - std::vector> m_sub_compilemodels; - friend ov::CoreImpl; protected: diff --git a/src/plugins/intel_cpu/src/async_infer_request.cpp b/src/plugins/intel_cpu/src/async_infer_request.cpp index b77410579e9..403e7120384 100644 --- a/src/plugins/intel_cpu/src/async_infer_request.cpp +++ b/src/plugins/intel_cpu/src/async_infer_request.cpp @@ -3,6 +3,7 @@ // #include "async_infer_request.h" +#include "cpu_message.hpp" ov::intel_cpu::AsyncInferRequest::AsyncInferRequest( const std::shared_ptr& request, @@ -13,14 +14,13 @@ ov::intel_cpu::AsyncInferRequest::AsyncInferRequest( } ov::intel_cpu::AsyncInferRequest::~AsyncInferRequest() { + if (m_sub_infers) { + auto message = ov::threading::message_manager(); + message->stop_server_thread(); + } stop_and_wait(); } void ov::intel_cpu::AsyncInferRequest::throw_if_canceled() const { check_cancelled_state(); } - -void ov::intel_cpu::AsyncInferRequest::setSubInferRequest( - const std::vector>& requests) { - m_sub_infer_requests = requests; -} diff --git a/src/plugins/intel_cpu/src/async_infer_request.h b/src/plugins/intel_cpu/src/async_infer_request.h index 0b69fa11e3a..b0162e8b6b6 100644 --- a/src/plugins/intel_cpu/src/async_infer_request.h +++ b/src/plugins/intel_cpu/src/async_infer_request.h @@ -17,15 +17,12 @@ public: const std::shared_ptr& callback_executor); ~AsyncInferRequest(); - void setSubInferRequest(const std::vector>& requests); - - std::vector> getSubInferRequest() const { - return m_sub_infer_requests; + void setSubInfer(bool sub_infer) { + m_sub_infers = sub_infer; } void throw_if_canceled() const; -private: - std::vector> m_sub_infer_requests; + bool m_sub_infers = false; }; } // namespace intel_cpu diff --git a/src/plugins/intel_cpu/src/compiled_model.cpp b/src/plugins/intel_cpu/src/compiled_model.cpp index 5f8f727762f..bf63b45658b 100644 --- a/src/plugins/intel_cpu/src/compiled_model.cpp +++ b/src/plugins/intel_cpu/src/compiled_model.cpp @@ -19,6 +19,7 @@ #include "openvino/runtime/threading/cpu_streams_executor.hpp" #include "transformations/utils/utils.hpp" #include "openvino/runtime/threading/cpu_streams_info.hpp" +#include "cpu_message.hpp" #include "cpu/x64/cpu_isa_traits.hpp" #include @@ -56,7 +57,6 @@ CompiledModel::CompiledModel(const std::shared_ptr& model, if (!core) OPENVINO_THROW("Unable to get API version. Core is unavailable"); - ov::threading::IStreamsExecutor::Ptr stream_executor = nullptr; IStreamsExecutor::Config executor_confg; if (cfg.exclusiveAsyncRequests) { // special case when all InferRequests are muxed into a single queue @@ -70,12 +70,11 @@ CompiledModel::CompiledModel(const std::shared_ptr& model, false, true}; } else { - executor_confg = m_cfg.subStreamExecConfig.get_name() != "StreamExecutor" + executor_confg = m_cfg.subStreamExecConfig.get_name() != "StreamsExecutor" ? std::move(m_cfg.subStreamExecConfig) : std::move(m_cfg.streamExecutorConfig); } - stream_executor = m_plugin->get_executor_manager()->get_idle_cpu_streams_executor(executor_confg); - m_task_executor = stream_executor; + m_task_executor = m_plugin->get_executor_manager()->get_idle_cpu_streams_executor(executor_confg); } if (0 != m_cfg.streamExecutorConfig.get_streams()) { m_callback_executor = m_plugin->get_executor_manager()->get_idle_cpu_streams_executor( @@ -120,12 +119,15 @@ CompiledModel::CompiledModel(const std::shared_ptr& model, std::cout << "xxxxxxxxx: m_subCompileModel: " << m_cfg.enableSubStreams << ", " << m_cfg.streamExecutorConfig.get_sub_streams() << "\n"; if (m_cfg.enableSubStreams) { m_cfg.enableSubStreams = false; + m_subCompileModel = true; + std::vector> sub_models; + auto message = message_manager(); for (int i = 0; i < m_cfg.streamExecutorConfig.get_sub_streams(); i++) { auto streams_info_table = m_cfg.streamExecutorConfig.get_streams_info_table(); std::vector> info_table; info_table.push_back(streams_info_table[i + 1]); info_table[0][NUMBER_OF_STREAMS] = 1; - auto streamExecutorConfig = IStreamsExecutor::Config{"CPUSubStreamsExecutor" + std::to_string(i), + auto streamExecutorConfig = IStreamsExecutor::Config{"CPUStreamsExecutor", 1, 1, ov::hint::SchedulingCoreType::ANY_CORE, @@ -134,8 +136,9 @@ CompiledModel::CompiledModel(const std::shared_ptr& model, info_table, m_cfg.streamsRankTable[i]}; m_cfg.subStreamExecConfig = std::move(streamExecutorConfig); - m_sub_compilemodels.push_back(std::make_shared(model, plugin, m_cfg, loaded_from_cache)); + sub_models.push_back(std::make_shared(model, plugin, m_cfg, loaded_from_cache)); } + message->setSubCompileModels(sub_models); } // init sub stream threads of executor // int sub_streams = m_cfg.streamExecutorConfig.get_sub_streams(); @@ -202,12 +205,15 @@ std::shared_ptr CompiledModel::create_infer_request() co std::make_shared(std::static_pointer_cast(internal_request), get_task_executor(), get_callback_executor()); - if (m_sub_compilemodels.size() > 0) { + if (m_subCompileModel) { + auto message = message_manager(); + auto sub_models = message->getSubCompileModels(); std::vector> requests; - for (int i = 0; i < m_sub_compilemodels.size(); i++) { - requests.push_back(m_sub_compilemodels[i]->create_infer_request()); + for (int i = 0; i < sub_models.size(); i++) { + requests.push_back(sub_models[i]->create_infer_request()); } - async_infer_request->setSubInferRequest(requests); + message->setSubInferRequest(requests); + async_infer_request->setSubInfer(true); } return async_infer_request; } diff --git a/src/plugins/intel_cpu/src/compiled_model.h b/src/plugins/intel_cpu/src/compiled_model.h index e3a25f2bd87..a02c3d1d12d 100644 --- a/src/plugins/intel_cpu/src/compiled_model.h +++ b/src/plugins/intel_cpu/src/compiled_model.h @@ -74,12 +74,7 @@ private: */ GraphGuard::Lock get_graph() const; - std::vector> get_sub_compilemodles() const { - return m_sub_compilemodels; - } - - // bool m_subCompileModel = true; - std::vector> m_sub_compilemodels; + bool m_subCompileModel = false; }; } // namespace intel_cpu diff --git a/src/inference/src/dev/threading/cpu_message.cpp b/src/plugins/intel_cpu/src/cpu_message.cpp similarity index 73% rename from src/inference/src/dev/threading/cpu_message.cpp rename to src/plugins/intel_cpu/src/cpu_message.cpp index 256d422a1e3..30ac9ac7c23 100644 --- a/src/inference/src/dev/threading/cpu_message.cpp +++ b/src/plugins/intel_cpu/src/cpu_message.cpp @@ -1,7 +1,7 @@ // Copyright (C) 2018-2024 Intel Corporation // SPDX-License-Identifier: Apache-2.0 // -#include "openvino/runtime/threading/cpu_message.hpp" +#include "cpu_message.hpp" #include #include #include @@ -12,7 +12,8 @@ namespace ov { namespace threading { -MessageManage::MessageManage() {} +MessageManage::MessageManage() { +} void MessageManage::send_message(MessageInfo msg_info) { { @@ -49,7 +50,7 @@ void MessageManage::server_wait(int streams_num) { while (!_isServerStopped) { std::vector msgQueue; { - // std::cout << "server_wait ........\n"; + // std::cout << "server_wait ........" << _isServerStopped << "\n"; std::unique_lock lock(_msgMutex); while (_messageQueue.empty()) { _msgCondVar.wait(lock); @@ -77,6 +78,8 @@ void MessageManage::server_wait(int streams_num) { _inferCondVar.notify_one(); count = 0; } + } else if (msg_type == QUIT) { + _isServerStopped = true; } } } @@ -85,14 +88,40 @@ void MessageManage::server_wait(int streams_num) { } } +void MessageManage::setSubCompileModels(std::vector> models) { + m_sub_compilemodels = models; + std::cout << __FUNCTION__ << ": " << m_sub_compilemodels.size() << "\n"; +} + +std::vector> MessageManage::getSubCompileModels() { + return m_sub_compilemodels; + std::cout << __FUNCTION__ << ": " << m_sub_compilemodels.size() << "\n"; +} + +void MessageManage::setSubInferRequest(std::vector> requests) { + m_sub_infer_requests = requests; + std::cout << __FUNCTION__ << ": " << m_sub_infer_requests.size() << "\n"; +} + +std::vector> MessageManage::getSubInferRequest() { + return m_sub_infer_requests; + std::cout << __FUNCTION__ << ": " << m_sub_infer_requests.size() << "\n"; +} + MessageManage::~MessageManage() { - _isServerStopped = true; - _msgCondVar.notify_one(); + std::cout << "~MessageManage\n"; +} + +void MessageManage::stop_server_thread() { + MessageInfo msg_info; + msg_info.msg_type = ov::threading::MsgType::QUIT; + send_message(msg_info); if (_serverThread.joinable()) { _serverThread.join(); } + m_sub_infer_requests.clear(); + m_sub_compilemodels.clear(); } - namespace { class MessageManageHolder { @@ -117,7 +146,7 @@ public: } // namespace -std::shared_ptr message_manager(){ +std::shared_ptr message_manager() { static MessageManageHolder message_manage; return message_manage.get(); } diff --git a/src/inference/dev_api/openvino/runtime/threading/cpu_message.hpp b/src/plugins/intel_cpu/src/cpu_message.hpp similarity index 66% rename from src/inference/dev_api/openvino/runtime/threading/cpu_message.hpp rename to src/plugins/intel_cpu/src/cpu_message.hpp index 78d6606722d..b78406c72ef 100644 --- a/src/inference/dev_api/openvino/runtime/threading/cpu_message.hpp +++ b/src/plugins/intel_cpu/src/cpu_message.hpp @@ -21,11 +21,13 @@ #include "openvino/runtime/common.hpp" #include "openvino/runtime/threading/istreams_executor.hpp" #include "openvino/runtime/threading/itask_executor.hpp" +#include "compiled_model.h" +#include "openvino/runtime/iasync_infer_request.hpp" namespace ov { namespace threading { -enum MsgType { TP, START_INFER, CALL_BACK }; +enum MsgType { TP, START_INFER, CALL_BACK, QUIT }; struct MessageInfo { MsgType msg_type; @@ -44,8 +46,20 @@ public: void server_wait(int streams_num); + void stop_server_thread(); + ~MessageManage(); + + void setSubCompileModels(std::vector> models); + std::vector> getSubCompileModels(); + + void setSubInferRequest(std::vector> requests); + std::vector> getSubInferRequest(); + private: + int sub_streams; + std::vector> m_sub_compilemodels; + std::vector> m_sub_infer_requests; std::thread _serverThread; bool _isServerStopped = false; std::vector _messageQueue; diff --git a/src/plugins/intel_cpu/src/infer_request.cpp b/src/plugins/intel_cpu/src/infer_request.cpp index ba890d965d5..fff5aa9afa1 100644 --- a/src/plugins/intel_cpu/src/infer_request.cpp +++ b/src/plugins/intel_cpu/src/infer_request.cpp @@ -18,7 +18,7 @@ #include "proxy_mem_mgr.h" #include "utils/general_utils.h" #include "utils/ngraph_utils.hpp" -#include "openvino/runtime/threading/cpu_message.hpp" +#include "cpu_message.hpp" using OvString = ov::element_type_traits::value_type; @@ -109,13 +109,14 @@ void SyncInferRequest::infer() { OV_ITT_SCOPED_TASK(itt::domains::intel_cpu, m_profiling_task); auto graphLock = m_compiled_model->get_graph(); m_graph = &(graphLock._graph); - auto streams_executor = m_graph->context->getCPUStreamExecutor(); + auto message = ov::threading::message_manager(); + // auto streams_executor = m_graph->context->getCPUStreamExecutor(); throw_if_canceled(); - auto requests = m_asyncRequest->getSubInferRequest(); - // std::cout << "[ infer ] " << requests.size() << "\n"; - if (requests.size() > 0) { - auto message = ov::threading::message_manager(); + // std::cout << "[ infer ] " << m_asyncRequest->m_sub_infers << "\n"; + if (m_asyncRequest->m_sub_infers) { + // auto message = ov::threading::message_manager(); + auto requests = message->getSubInferRequest();//m_asyncRequest->getSubInferRequest(); message->server_wait(requests.size()); ov::threading::MessageInfo msg_info; msg_info.msg_type = ov::threading::MsgType::START_INFER; @@ -330,6 +331,16 @@ void SyncInferRequest::change_default_ptr() { } std::vector> SyncInferRequest::query_state() const { + if (m_asyncRequest->m_sub_infers) { + auto message = ov::threading::message_manager(); + auto requests = message->getSubInferRequest(); + std::vector> states; + for (auto request : requests) { + auto cur = request->query_state(); + states.insert(states.end(), cur.begin(), cur.end()); + } + return states; + } return {m_memory_states.begin(), m_memory_states.end()}; } @@ -629,26 +640,24 @@ SyncInferRequest::OutputControlBlock::OutputControlBlock(const ov::element::Type void SyncInferRequest::sub_streams_infer() { std::map, ov::SoPtr> input_tensors; - auto requests = m_asyncRequest->getSubInferRequest(); + auto message = ov::threading::message_manager(); + auto requests = message->getSubInferRequest(); auto inputs = m_asyncRequest->get_inputs(); auto outputs = m_asyncRequest->get_outputs(); - auto message = ov::threading::message_manager(); + size_t requests_num = requests.size(); size_t requests_count = 0; // std::cout << "[ sub_streams_infer ] inputs: " << inputs.size() << " requests: " << requests_num << "\n"; if (requests.size() > 0) { - for (const auto& input : inputs) { - auto tensor = m_asyncRequest->get_tensor(input); - input_tensors.insert({input, tensor}); - } for (const auto& output : outputs) { auto tensor = requests[0]->get_tensor(output); - m_asyncRequest->set_tensor(output, tensor); + set_tensor(output, tensor); } for (size_t i = 0; i < requests_num; i++) { - for (auto& input : input_tensors) { - requests[i]->set_tensor(input.first, input.second); + for (auto& input : inputs) { + auto tensor = m_asyncRequest->get_tensor(input); + requests[i]->set_tensor(input, tensor); } requests[i]->set_callback([i, requests, message](const std::exception_ptr& ptr) { diff --git a/src/plugins/intel_cpu/src/plugin.cpp b/src/plugins/intel_cpu/src/plugin.cpp index 0aa82a94d2e..efe40a557e7 100644 --- a/src/plugins/intel_cpu/src/plugin.cpp +++ b/src/plugins/intel_cpu/src/plugin.cpp @@ -131,13 +131,12 @@ Plugin::Plugin() : deviceFullName(getDeviceFullName()), specialSetup(new CPUSpec }); auto& ov_version = ov::get_openvino_version(); m_compiled_model_runtime_properties["OV_VERSION"] = std::string(ov_version.buildNumber); + m_msg_manager = ov::threading::message_manager(); } Plugin::~Plugin() { executor_manager()->clear("CPU"); executor_manager()->clear("CPUStreamsExecutor"); - executor_manager()->clear("CPUSubStreamsExecutor0"); - executor_manager()->clear("CPUSubStreamsExecutor1"); executor_manager()->clear("CPUMainStreamExecutor"); executor_manager()->clear("CPUCallbackExecutor"); } @@ -284,11 +283,11 @@ std::shared_ptr Plugin::compile_model(const std::shared_ptr< conf.readProperties(config, modelType); calculate_streams(conf, cloned_model); - if (conf.streamExecutorConfig.get_sub_stream_mode() == - IStreamsExecutor::Config::StreamsMode::SUB_STREAMS_FOR_SOCKET) { - int num_sub_streams = conf.streamExecutorConfig.get_sub_streams(); - transformations.SetSubStreamNum(num_sub_streams); - } + // if (conf.streamExecutorConfig.get_sub_stream_mode() == + // IStreamsExecutor::Config::StreamsMode::SUB_STREAMS_FOR_SOCKET) { + // int num_sub_streams = conf.streamExecutorConfig.get_sub_streams(); + // transformations.SetSubStreamNum(num_sub_streams); + // } transformations.PostLpt(); transformations.Snippets(); diff --git a/src/plugins/intel_cpu/src/plugin.h b/src/plugins/intel_cpu/src/plugin.h index 3332d873fdc..84c0c144d02 100644 --- a/src/plugins/intel_cpu/src/plugin.h +++ b/src/plugins/intel_cpu/src/plugin.h @@ -6,6 +6,7 @@ #include "compiled_model.h" #include "cpu_streams_calculation.hpp" +#include "cpu_message.hpp" namespace ov { namespace intel_cpu { @@ -43,6 +44,8 @@ public: OPENVINO_THROW_NOT_IMPLEMENTED("Not Implemented get_default_context is not supported by CPU plugin!"); }; + std::shared_ptr m_msg_manager; + private: ov::Any get_ro_property(const std::string& name, const ov::AnyMap& options) const; From 9aaca1aa29ea0117073d1b8e301ac74e58a6d744 Mon Sep 17 00:00:00 2001 From: sunxiaoxia2022 Date: Thu, 16 May 2024 09:55:14 +0800 Subject: [PATCH 09/18] fix memory leak --- src/plugins/intel_cpu/src/async_infer_request.cpp | 1 + src/plugins/intel_cpu/src/cpu_message.cpp | 15 ++++++++------- src/plugins/intel_cpu/src/cpu_message.hpp | 2 ++ src/plugins/intel_cpu/src/infer_request.cpp | 10 +++------- 4 files changed, 14 insertions(+), 14 deletions(-) diff --git a/src/plugins/intel_cpu/src/async_infer_request.cpp b/src/plugins/intel_cpu/src/async_infer_request.cpp index 403e7120384..628cc6156eb 100644 --- a/src/plugins/intel_cpu/src/async_infer_request.cpp +++ b/src/plugins/intel_cpu/src/async_infer_request.cpp @@ -17,6 +17,7 @@ ov::intel_cpu::AsyncInferRequest::~AsyncInferRequest() { if (m_sub_infers) { auto message = ov::threading::message_manager(); message->stop_server_thread(); + message->clear(); } stop_and_wait(); } diff --git a/src/plugins/intel_cpu/src/cpu_message.cpp b/src/plugins/intel_cpu/src/cpu_message.cpp index b437d990843..fbe7fa22004 100644 --- a/src/plugins/intel_cpu/src/cpu_message.cpp +++ b/src/plugins/intel_cpu/src/cpu_message.cpp @@ -12,8 +12,9 @@ namespace ov { namespace threading { -MessageManage::MessageManage() { -} +MessageManage::MessageManage() {} + +MessageManage::~MessageManage() {} void MessageManage::send_message(MessageInfo msg_info) { { @@ -118,7 +119,7 @@ void MessageManage::server_wait(int streams_num) { } } } - std::cout << "-------- server_wait end ---------\n"; + // std::cout << "-------- server_wait end ---------\n"; }); } } @@ -139,10 +140,6 @@ std::vector> MessageManage::getSubInferR return m_sub_infer_requests; } -MessageManage::~MessageManage() { - std::cout << "~MessageManage\n"; -} - void MessageManage::stop_server_thread() { MessageInfo msg_info; msg_info.msg_type = ov::threading::MsgType::QUIT; @@ -150,9 +147,13 @@ void MessageManage::stop_server_thread() { if (_serverThread.joinable()) { _serverThread.join(); } +} + +void MessageManage::clear() { m_sub_infer_requests.clear(); m_sub_compilemodels.clear(); } + namespace { class MessageManageHolder { diff --git a/src/plugins/intel_cpu/src/cpu_message.hpp b/src/plugins/intel_cpu/src/cpu_message.hpp index e9febcfd70a..f2c87c58a51 100644 --- a/src/plugins/intel_cpu/src/cpu_message.hpp +++ b/src/plugins/intel_cpu/src/cpu_message.hpp @@ -50,6 +50,8 @@ public: void stop_server_thread(); + void clear(); + ~MessageManage(); void setSubCompileModels(std::vector> models); diff --git a/src/plugins/intel_cpu/src/infer_request.cpp b/src/plugins/intel_cpu/src/infer_request.cpp index 25bf4202ce2..d1734fef5ae 100644 --- a/src/plugins/intel_cpu/src/infer_request.cpp +++ b/src/plugins/intel_cpu/src/infer_request.cpp @@ -110,14 +110,11 @@ void SyncInferRequest::infer() { auto graphLock = m_compiled_model->get_graph(); m_graph = &(graphLock._graph); auto message = ov::threading::message_manager(); - // auto streams_executor = m_graph->context->getCPUStreamExecutor(); throw_if_canceled(); // std::cout << "[ infer ] " << m_asyncRequest->m_sub_infers << "\n"; if (m_asyncRequest->m_sub_infers) { - // auto message = ov::threading::message_manager(); - auto requests = message->getSubInferRequest();//m_asyncRequest->getSubInferRequest(); - message->server_wait(requests.size()); + message->server_wait(message->getSubInferRequest().size()); ov::threading::MessageInfo msg_info; msg_info.msg_type = ov::threading::MsgType::START_INFER; ov::threading::Task task = [&] { @@ -333,9 +330,8 @@ void SyncInferRequest::change_default_ptr() { std::vector> SyncInferRequest::query_state() const { if (m_asyncRequest->m_sub_infers) { auto message = ov::threading::message_manager(); - auto requests = message->getSubInferRequest(); std::vector> states; - for (auto request : requests) { + for (auto request : message->getSubInferRequest()) { auto cur = request->query_state(); states.insert(states.end(), cur.begin(), cur.end()); } @@ -659,7 +655,7 @@ void SyncInferRequest::sub_streams_infer() { requests[i]->set_tensor(input, tensor); } - requests[i]->set_callback([i, requests, message](const std::exception_ptr& ptr) { + requests[i]->set_callback([i, message](const std::exception_ptr& ptr) { // std::cout << "set_callback------ " << i << "\n"; ov::threading::MessageInfo msg_info; msg_info.msg_type = ov::threading::MsgType::CALL_BACK; From 12a2c4d69bd11a0eb71c9f92de3bc613346c8e52 Mon Sep 17 00:00:00 2001 From: sunxiaoxia2022 Date: Thu, 16 May 2024 15:08:53 +0800 Subject: [PATCH 10/18] move cpu_message to threading path --- samples/cpp/benchmark_app/infer_request_wrap.hpp | 1 - .../openvino/runtime/threading}/cpu_message.hpp | 8 ++++---- src/inference/src/dev/iasync_infer_request.cpp | 4 ++-- .../src/dev/threading}/cpu_message.cpp | 12 +++++++++--- src/plugins/intel_cpu/src/async_infer_request.cpp | 3 ++- src/plugins/intel_cpu/src/compiled_model.cpp | 10 +++++----- src/plugins/intel_cpu/src/infer_request.cpp | 8 ++++---- src/plugins/intel_cpu/src/plugin.h | 2 +- 8 files changed, 27 insertions(+), 21 deletions(-) rename src/{plugins/intel_cpu/src => inference/dev_api/openvino/runtime/threading}/cpu_message.hpp (86%) rename src/{plugins/intel_cpu/src => inference/src/dev/threading}/cpu_message.cpp (92%) diff --git a/samples/cpp/benchmark_app/infer_request_wrap.hpp b/samples/cpp/benchmark_app/infer_request_wrap.hpp index fa45fde4a61..0c09a6d6ada 100644 --- a/samples/cpp/benchmark_app/infer_request_wrap.hpp +++ b/samples/cpp/benchmark_app/infer_request_wrap.hpp @@ -151,7 +151,6 @@ public: const double latency, const std::exception_ptr& ptr = nullptr) { std::unique_lock lock(_mutex); - // std::cout << "put_idle_request: id: " << id << " lat_group_id: " << lat_group_id << " latency: " << latency << "\n"; if (ptr) { inferenceException = ptr; } else { diff --git a/src/plugins/intel_cpu/src/cpu_message.hpp b/src/inference/dev_api/openvino/runtime/threading/cpu_message.hpp similarity index 86% rename from src/plugins/intel_cpu/src/cpu_message.hpp rename to src/inference/dev_api/openvino/runtime/threading/cpu_message.hpp index f2c87c58a51..971c044e674 100644 --- a/src/plugins/intel_cpu/src/cpu_message.hpp +++ b/src/inference/dev_api/openvino/runtime/threading/cpu_message.hpp @@ -21,7 +21,7 @@ #include "openvino/runtime/common.hpp" #include "openvino/runtime/threading/istreams_executor.hpp" #include "openvino/runtime/threading/itask_executor.hpp" -#include "compiled_model.h" +#include "openvino/runtime/compiled_model.hpp" #include "openvino/runtime/iasync_infer_request.hpp" namespace ov { @@ -54,15 +54,15 @@ public: ~MessageManage(); - void setSubCompileModels(std::vector> models); - std::vector> getSubCompileModels(); + void setSubCompileModels(std::vector> models); + std::vector> getSubCompileModels(); void setSubInferRequest(std::vector> requests); std::vector> getSubInferRequest(); private: int sub_streams; - std::vector> m_sub_compilemodels; + std::vector> m_sub_compilemodels; std::vector> m_sub_infer_requests; std::thread _serverThread; bool _isServerStopped = false; diff --git a/src/inference/src/dev/iasync_infer_request.cpp b/src/inference/src/dev/iasync_infer_request.cpp index c914772c23f..9e914a7c38b 100644 --- a/src/inference/src/dev/iasync_infer_request.cpp +++ b/src/inference/src/dev/iasync_infer_request.cpp @@ -208,12 +208,12 @@ std::vector ov::IAsyncInferRequest::get_profiling_info() cons } ov::SoPtr ov::IAsyncInferRequest::get_tensor(const ov::Output& port) const { - // check_state(); + check_state(); return m_sync_request->get_tensor(port); } void ov::IAsyncInferRequest::set_tensor(const ov::Output& port, const ov::SoPtr& tensor) { - // check_state(); + check_state(); return m_sync_request->set_tensor(port, tensor); } diff --git a/src/plugins/intel_cpu/src/cpu_message.cpp b/src/inference/src/dev/threading/cpu_message.cpp similarity index 92% rename from src/plugins/intel_cpu/src/cpu_message.cpp rename to src/inference/src/dev/threading/cpu_message.cpp index fbe7fa22004..8d6f19c445b 100644 --- a/src/plugins/intel_cpu/src/cpu_message.cpp +++ b/src/inference/src/dev/threading/cpu_message.cpp @@ -1,7 +1,7 @@ // Copyright (C) 2018-2024 Intel Corporation // SPDX-License-Identifier: Apache-2.0 // -#include "cpu_message.hpp" +#include "openvino/runtime/threading/cpu_message.hpp" #include #include #include @@ -124,11 +124,11 @@ void MessageManage::server_wait(int streams_num) { } } -void MessageManage::setSubCompileModels(std::vector> models) { +void MessageManage::setSubCompileModels(std::vector> models) { m_sub_compilemodels = models; } -std::vector> MessageManage::getSubCompileModels() { +std::vector> MessageManage::getSubCompileModels() { return m_sub_compilemodels; } @@ -150,7 +150,13 @@ void MessageManage::stop_server_thread() { } void MessageManage::clear() { + for (size_t i = 0; i < m_sub_infer_requests.size(); i++) { + std::cout << "infer count_" << i << " : " << m_sub_infer_requests[i].use_count() << "\n"; + } m_sub_infer_requests.clear(); + for (size_t i = 0; i < m_sub_compilemodels.size(); i++) { + std::cout << "model count_" << i << " : " << m_sub_compilemodels[i].use_count() << "\n"; + } m_sub_compilemodels.clear(); } diff --git a/src/plugins/intel_cpu/src/async_infer_request.cpp b/src/plugins/intel_cpu/src/async_infer_request.cpp index 628cc6156eb..f02de94f25e 100644 --- a/src/plugins/intel_cpu/src/async_infer_request.cpp +++ b/src/plugins/intel_cpu/src/async_infer_request.cpp @@ -3,7 +3,8 @@ // #include "async_infer_request.h" -#include "cpu_message.hpp" + +#include "openvino/runtime/threading/cpu_message.hpp" ov::intel_cpu::AsyncInferRequest::AsyncInferRequest( const std::shared_ptr& request, diff --git a/src/plugins/intel_cpu/src/compiled_model.cpp b/src/plugins/intel_cpu/src/compiled_model.cpp index bf63b45658b..580187ba569 100644 --- a/src/plugins/intel_cpu/src/compiled_model.cpp +++ b/src/plugins/intel_cpu/src/compiled_model.cpp @@ -19,7 +19,7 @@ #include "openvino/runtime/threading/cpu_streams_executor.hpp" #include "transformations/utils/utils.hpp" #include "openvino/runtime/threading/cpu_streams_info.hpp" -#include "cpu_message.hpp" +#include "openvino/runtime/threading/cpu_message.hpp" #include "cpu/x64/cpu_isa_traits.hpp" #include @@ -116,11 +116,11 @@ CompiledModel::CompiledModel(const std::shared_ptr& model, } else { CompiledModel::get_graph(); } - std::cout << "xxxxxxxxx: m_subCompileModel: " << m_cfg.enableSubStreams << ", " << m_cfg.streamExecutorConfig.get_sub_streams() << "\n"; + // std::cout << "m_subCompileModel: " << m_cfg.enableSubStreams << ", " << m_cfg.streamExecutorConfig.get_sub_streams() << "\n"; if (m_cfg.enableSubStreams) { m_cfg.enableSubStreams = false; m_subCompileModel = true; - std::vector> sub_models; + std::vector> sub_models; auto message = message_manager(); for (int i = 0; i < m_cfg.streamExecutorConfig.get_sub_streams(); i++) { auto streams_info_table = m_cfg.streamExecutorConfig.get_streams_info_table(); @@ -209,8 +209,8 @@ std::shared_ptr CompiledModel::create_infer_request() co auto message = message_manager(); auto sub_models = message->getSubCompileModels(); std::vector> requests; - for (int i = 0; i < sub_models.size(); i++) { - requests.push_back(sub_models[i]->create_infer_request()); + for (auto model : sub_models) { + requests.push_back(model->create_infer_request()); } message->setSubInferRequest(requests); async_infer_request->setSubInfer(true); diff --git a/src/plugins/intel_cpu/src/infer_request.cpp b/src/plugins/intel_cpu/src/infer_request.cpp index d1734fef5ae..f8d4b173a70 100644 --- a/src/plugins/intel_cpu/src/infer_request.cpp +++ b/src/plugins/intel_cpu/src/infer_request.cpp @@ -18,7 +18,7 @@ #include "proxy_mem_mgr.h" #include "utils/general_utils.h" #include "utils/ngraph_utils.hpp" -#include "cpu_message.hpp" +#include "openvino/runtime/threading/cpu_message.hpp" using OvString = ov::element_type_traits::value_type; @@ -638,8 +638,8 @@ void SyncInferRequest::sub_streams_infer() { std::map, ov::SoPtr> input_tensors; auto message = ov::threading::message_manager(); auto requests = message->getSubInferRequest(); - auto inputs = m_asyncRequest->get_inputs(); - auto outputs = m_asyncRequest->get_outputs(); + auto inputs = get_inputs(); + auto outputs = get_outputs(); size_t requests_num = requests.size(); // std::cout << "[ sub_streams_infer ] inputs: " << inputs.size() << " requests: " << requests_num << "\n"; @@ -651,7 +651,7 @@ void SyncInferRequest::sub_streams_infer() { } for (size_t i = 0; i < requests_num; i++) { for (auto& input : inputs) { - auto tensor = m_asyncRequest->get_tensor(input); + auto tensor = get_tensor(input); requests[i]->set_tensor(input, tensor); } diff --git a/src/plugins/intel_cpu/src/plugin.h b/src/plugins/intel_cpu/src/plugin.h index 84c0c144d02..e475c5105ca 100644 --- a/src/plugins/intel_cpu/src/plugin.h +++ b/src/plugins/intel_cpu/src/plugin.h @@ -6,7 +6,7 @@ #include "compiled_model.h" #include "cpu_streams_calculation.hpp" -#include "cpu_message.hpp" +#include "openvino/runtime/threading/cpu_message.hpp" namespace ov { namespace intel_cpu { From 3905ade0f2a40fffbbc62d153285787b91666860 Mon Sep 17 00:00:00 2001 From: sunxiaoxia2022 Date: Thu, 16 May 2024 16:28:05 +0800 Subject: [PATCH 11/18] change some name --- .../runtime/threading/cpu_message.hpp | 25 ++++---- .../src/dev/threading/cpu_message.cpp | 58 +++++++++---------- .../intel_cpu/src/async_infer_request.cpp | 2 +- .../intel_cpu/src/async_infer_request.h | 6 +- src/plugins/intel_cpu/src/compiled_model.cpp | 12 ++-- src/plugins/intel_cpu/src/compiled_model.h | 2 +- src/plugins/intel_cpu/src/infer_request.cpp | 12 ++-- src/plugins/intel_cpu/src/plugin.h | 2 +- 8 files changed, 58 insertions(+), 61 deletions(-) diff --git a/src/inference/dev_api/openvino/runtime/threading/cpu_message.hpp b/src/inference/dev_api/openvino/runtime/threading/cpu_message.hpp index 971c044e674..79d16398fd5 100644 --- a/src/inference/dev_api/openvino/runtime/threading/cpu_message.hpp +++ b/src/inference/dev_api/openvino/runtime/threading/cpu_message.hpp @@ -27,7 +27,7 @@ namespace ov { namespace threading { -enum MsgType { TP, START_INFER, CALL_BACK, REDUCE, QUIT }; +enum MsgType { START_INFER, TENSOR_PARALLEL, CALL_BACK, REDUCE, QUIT }; struct MessageInfo { MsgType msg_type; @@ -35,10 +35,11 @@ struct MessageInfo { void* buf; Task task; }; -class OPENVINO_RUNTIME_API MessageManage { +class OPENVINO_RUNTIME_API MessageManager { public: - MessageManage(); - void send_message(MessageInfo msg_info); + MessageManager(); + + void send_message(const MessageInfo& msg_info); std::vector wait_message(int cur_rank, int streams_num); @@ -52,18 +53,18 @@ public: void clear(); - ~MessageManage(); + ~MessageManager(); - void setSubCompileModels(std::vector> models); - std::vector> getSubCompileModels(); + void set_sub_compiled_models(std::vector> models); + std::vector> get_sub_compiled_models(); - void setSubInferRequest(std::vector> requests); - std::vector> getSubInferRequest(); + void set_sub_infer_requests(std::vector> requests); + std::vector> get_sub_infer_requests(); private: int sub_streams; - std::vector> m_sub_compilemodels; - std::vector> m_sub_infer_requests; + std::vector> _sub_compiled_models; + std::vector> _sub_infer_requests; std::thread _serverThread; bool _isServerStopped = false; std::vector _messageQueue; @@ -79,6 +80,6 @@ private: std::condition_variable _reduceCondVar; }; -OPENVINO_RUNTIME_API std::shared_ptr message_manager(); +OPENVINO_RUNTIME_API std::shared_ptr message_manager(); } // namespace threading } // namespace ov \ No newline at end of file diff --git a/src/inference/src/dev/threading/cpu_message.cpp b/src/inference/src/dev/threading/cpu_message.cpp index 8d6f19c445b..c2efa2555e4 100644 --- a/src/inference/src/dev/threading/cpu_message.cpp +++ b/src/inference/src/dev/threading/cpu_message.cpp @@ -2,6 +2,7 @@ // SPDX-License-Identifier: Apache-2.0 // #include "openvino/runtime/threading/cpu_message.hpp" + #include #include #include @@ -12,20 +13,21 @@ namespace ov { namespace threading { -MessageManage::MessageManage() {} +MessageManager::MessageManager() {} -MessageManage::~MessageManage() {} +MessageManager::~MessageManager() {} -void MessageManage::send_message(MessageInfo msg_info) { +void MessageManager::send_message(const MessageInfo& msg_info) { { std::lock_guard lock(_msgMutex); _messageQueue.push_back(msg_info); - // std::cout << "send : " << msg_info.msg_type << ", " << (msg_info.rank.size() > 0 ? msg_info.rank[0] : -1) << "\n"; + // std::cout << "send : " << msg_info.msg_type << ", " << (msg_info.rank.size() > 0 ? msg_info.rank[0] : -1) + // << "\n"; } _msgCondVar.notify_all(); } -std::vector MessageManage::wait_message(int cur_rank, int streams_num) { +std::vector MessageManager::wait_message(int cur_rank, int streams_num) { std::vector messages_total; std::unique_lock lock(_readMutex); _readCondVar.wait(lock, [&] { @@ -37,12 +39,12 @@ std::vector MessageManage::wait_message(int cur_rank, int streams_n return messages_total; } -void MessageManage::infer_wait() { +void MessageManager::infer_wait() { std::unique_lock lock(_inferMutex); _inferCondVar.wait(lock); } -void MessageManage::reduce_wait(int cur_rank, int streams_num) { +void MessageManager::reduce_wait(int cur_rank, int streams_num) { std::unique_lock lock(_reduceMutex); while (_reduceQueue[cur_rank] < streams_num) { // std::cout << "reduce_wait_" << cur_rank << " " << _reduceQueue[cur_rank] << " end\n"; @@ -51,7 +53,7 @@ void MessageManage::reduce_wait(int cur_rank, int streams_num) { _reduceQueue[cur_rank] = 0; } -void MessageManage::server_wait(int streams_num) { +void MessageManager::server_wait(int streams_num) { if (!_serverThread.joinable()) { _readQueue.assign(streams_num, std::vector()); _reduceQueue.assign(streams_num, 0); @@ -77,7 +79,7 @@ void MessageManage::server_wait(int streams_num) { if (msg_type == START_INFER) { Task task = std::move(rec_info.task); task(); - } else if (msg_type == TP) { + } else if (msg_type == TENSOR_PARALLEL) { // Resend _readQueue that failed last time bool stop = false; while (!stop) { @@ -124,23 +126,23 @@ void MessageManage::server_wait(int streams_num) { } } -void MessageManage::setSubCompileModels(std::vector> models) { - m_sub_compilemodels = models; +void MessageManager::set_sub_compiled_models(std::vector> models) { + _sub_compiled_models = models; } -std::vector> MessageManage::getSubCompileModels() { - return m_sub_compilemodels; +std::vector> MessageManager::get_sub_compiled_models() { + return _sub_compiled_models; } -void MessageManage::setSubInferRequest(std::vector> requests) { - m_sub_infer_requests = requests; +void MessageManager::set_sub_infer_requests(std::vector> requests) { + _sub_infer_requests = requests; } -std::vector> MessageManage::getSubInferRequest() { - return m_sub_infer_requests; +std::vector> MessageManager::get_sub_infer_requests() { + return _sub_infer_requests; } -void MessageManage::stop_server_thread() { +void MessageManager::stop_server_thread() { MessageInfo msg_info; msg_info.msg_type = ov::threading::MsgType::QUIT; send_message(msg_info); @@ -149,22 +151,16 @@ void MessageManage::stop_server_thread() { } } -void MessageManage::clear() { - for (size_t i = 0; i < m_sub_infer_requests.size(); i++) { - std::cout << "infer count_" << i << " : " << m_sub_infer_requests[i].use_count() << "\n"; - } - m_sub_infer_requests.clear(); - for (size_t i = 0; i < m_sub_compilemodels.size(); i++) { - std::cout << "model count_" << i << " : " << m_sub_compilemodels[i].use_count() << "\n"; - } - m_sub_compilemodels.clear(); +void MessageManager::clear() { + _sub_infer_requests.clear(); + _sub_compiled_models.clear(); } namespace { class MessageManageHolder { std::mutex _mutex; - std::weak_ptr _manager; + std::weak_ptr _manager; public: MessageManageHolder(const MessageManageHolder&) = delete; @@ -172,11 +168,11 @@ public: MessageManageHolder() = default; - std::shared_ptr get() { + std::shared_ptr get() { std::lock_guard lock(_mutex); auto manager = _manager.lock(); if (!manager) { - _manager = manager = std::make_shared(); + _manager = manager = std::make_shared(); } return manager; } @@ -184,7 +180,7 @@ public: } // namespace -std::shared_ptr message_manager() { +std::shared_ptr message_manager() { static MessageManageHolder message_manage; return message_manage.get(); } diff --git a/src/plugins/intel_cpu/src/async_infer_request.cpp b/src/plugins/intel_cpu/src/async_infer_request.cpp index f02de94f25e..383aa7dc460 100644 --- a/src/plugins/intel_cpu/src/async_infer_request.cpp +++ b/src/plugins/intel_cpu/src/async_infer_request.cpp @@ -15,7 +15,7 @@ ov::intel_cpu::AsyncInferRequest::AsyncInferRequest( } ov::intel_cpu::AsyncInferRequest::~AsyncInferRequest() { - if (m_sub_infers) { + if (m_has_sub_infers) { auto message = ov::threading::message_manager(); message->stop_server_thread(); message->clear(); diff --git a/src/plugins/intel_cpu/src/async_infer_request.h b/src/plugins/intel_cpu/src/async_infer_request.h index b0162e8b6b6..c7b6b5d0898 100644 --- a/src/plugins/intel_cpu/src/async_infer_request.h +++ b/src/plugins/intel_cpu/src/async_infer_request.h @@ -17,12 +17,12 @@ public: const std::shared_ptr& callback_executor); ~AsyncInferRequest(); - void setSubInfer(bool sub_infer) { - m_sub_infers = sub_infer; + void setSubInfer(bool has_sub_infer) { + m_has_sub_infers = has_sub_infer; } void throw_if_canceled() const; - bool m_sub_infers = false; + bool m_has_sub_infers = false; }; } // namespace intel_cpu diff --git a/src/plugins/intel_cpu/src/compiled_model.cpp b/src/plugins/intel_cpu/src/compiled_model.cpp index 580187ba569..45d05542d51 100644 --- a/src/plugins/intel_cpu/src/compiled_model.cpp +++ b/src/plugins/intel_cpu/src/compiled_model.cpp @@ -116,10 +116,10 @@ CompiledModel::CompiledModel(const std::shared_ptr& model, } else { CompiledModel::get_graph(); } - // std::cout << "m_subCompileModel: " << m_cfg.enableSubStreams << ", " << m_cfg.streamExecutorConfig.get_sub_streams() << "\n"; + // std::cout << "m_hans_sub_compiled_models: " << m_cfg.enableSubStreams << ", " << m_cfg.streamExecutorConfig.get_sub_streams() << "\n"; if (m_cfg.enableSubStreams) { m_cfg.enableSubStreams = false; - m_subCompileModel = true; + m_hans_sub_compiled_models = true; std::vector> sub_models; auto message = message_manager(); for (int i = 0; i < m_cfg.streamExecutorConfig.get_sub_streams(); i++) { @@ -138,7 +138,7 @@ CompiledModel::CompiledModel(const std::shared_ptr& model, m_cfg.subStreamExecConfig = std::move(streamExecutorConfig); sub_models.push_back(std::make_shared(model, plugin, m_cfg, loaded_from_cache)); } - message->setSubCompileModels(sub_models); + message->set_sub_compiled_models(sub_models); } // init sub stream threads of executor // int sub_streams = m_cfg.streamExecutorConfig.get_sub_streams(); @@ -205,14 +205,14 @@ std::shared_ptr CompiledModel::create_infer_request() co std::make_shared(std::static_pointer_cast(internal_request), get_task_executor(), get_callback_executor()); - if (m_subCompileModel) { + if (m_hans_sub_compiled_models) { auto message = message_manager(); - auto sub_models = message->getSubCompileModels(); + auto sub_models = message->get_sub_compiled_models(); std::vector> requests; for (auto model : sub_models) { requests.push_back(model->create_infer_request()); } - message->setSubInferRequest(requests); + message->set_sub_infer_requests(requests); async_infer_request->setSubInfer(true); } return async_infer_request; diff --git a/src/plugins/intel_cpu/src/compiled_model.h b/src/plugins/intel_cpu/src/compiled_model.h index a02c3d1d12d..c257c6147d1 100644 --- a/src/plugins/intel_cpu/src/compiled_model.h +++ b/src/plugins/intel_cpu/src/compiled_model.h @@ -74,7 +74,7 @@ private: */ GraphGuard::Lock get_graph() const; - bool m_subCompileModel = false; + bool m_hans_sub_compiled_models = false; }; } // namespace intel_cpu diff --git a/src/plugins/intel_cpu/src/infer_request.cpp b/src/plugins/intel_cpu/src/infer_request.cpp index f8d4b173a70..e6aa2bfbc42 100644 --- a/src/plugins/intel_cpu/src/infer_request.cpp +++ b/src/plugins/intel_cpu/src/infer_request.cpp @@ -112,9 +112,9 @@ void SyncInferRequest::infer() { auto message = ov::threading::message_manager(); throw_if_canceled(); - // std::cout << "[ infer ] " << m_asyncRequest->m_sub_infers << "\n"; - if (m_asyncRequest->m_sub_infers) { - message->server_wait(message->getSubInferRequest().size()); + // std::cout << "[ infer ] " << m_asyncRequest->m_has_sub_infers << "\n"; + if (m_asyncRequest->m_has_sub_infers) { + message->server_wait(message->get_sub_infer_requests().size()); ov::threading::MessageInfo msg_info; msg_info.msg_type = ov::threading::MsgType::START_INFER; ov::threading::Task task = [&] { @@ -328,10 +328,10 @@ void SyncInferRequest::change_default_ptr() { } std::vector> SyncInferRequest::query_state() const { - if (m_asyncRequest->m_sub_infers) { + if (m_asyncRequest->m_has_sub_infers) { auto message = ov::threading::message_manager(); std::vector> states; - for (auto request : message->getSubInferRequest()) { + for (auto request : message->get_sub_infer_requests()) { auto cur = request->query_state(); states.insert(states.end(), cur.begin(), cur.end()); } @@ -637,7 +637,7 @@ SyncInferRequest::OutputControlBlock::OutputControlBlock(const ov::element::Type void SyncInferRequest::sub_streams_infer() { std::map, ov::SoPtr> input_tensors; auto message = ov::threading::message_manager(); - auto requests = message->getSubInferRequest(); + auto requests = message->get_sub_infer_requests(); auto inputs = get_inputs(); auto outputs = get_outputs(); diff --git a/src/plugins/intel_cpu/src/plugin.h b/src/plugins/intel_cpu/src/plugin.h index e475c5105ca..c2d24e98ee6 100644 --- a/src/plugins/intel_cpu/src/plugin.h +++ b/src/plugins/intel_cpu/src/plugin.h @@ -44,7 +44,7 @@ public: OPENVINO_THROW_NOT_IMPLEMENTED("Not Implemented get_default_context is not supported by CPU plugin!"); }; - std::shared_ptr m_msg_manager; + std::shared_ptr m_msg_manager; private: ov::Any get_ro_property(const std::string& name, const ov::AnyMap& options) const; From bece8e53de190d35e7c1a98bd94ecfb68526f21e Mon Sep 17 00:00:00 2001 From: sunxiaoxia2022 Date: Fri, 17 May 2024 10:15:17 +0800 Subject: [PATCH 12/18] clone sub cfg to create sub compiled model --- .../dev/threading/cpu_streams_executor.cpp | 4 +- src/plugins/intel_cpu/src/compiled_model.cpp | 53 +++++++++---------- src/plugins/intel_cpu/src/compiled_model.h | 2 +- src/plugins/intel_cpu/src/config.h | 1 - 4 files changed, 26 insertions(+), 34 deletions(-) diff --git a/src/inference/src/dev/threading/cpu_streams_executor.cpp b/src/inference/src/dev/threading/cpu_streams_executor.cpp index 13964f1b60e..ce8da0a94e9 100644 --- a/src/inference/src/dev/threading/cpu_streams_executor.cpp +++ b/src/inference/src/dev/threading/cpu_streams_executor.cpp @@ -71,7 +71,6 @@ struct CPUStreamsExecutor::Impl { ((_impl->_config.get_streams() + _impl->_usedNumaNodes.size() - 1) / _impl->_usedNumaNodes.size())) : _impl->_usedNumaNodes.at(_streamId % _impl->_usedNumaNodes.size()); - // std::cout << "[ Stream ] " << _impl->_config.get_name() << " : " << _streamId << ", " << _impl->_config.get_streams() << "\n"; #if OV_THREAD == OV_THREAD_TBB || OV_THREAD == OV_THREAD_TBB_AUTO if (is_cpu_map_available() && _impl->_config.get_streams_info_table().size() > 0) { init_stream(); @@ -325,8 +324,7 @@ struct CPUStreamsExecutor::Impl { _exectorMgr = executor_manager(); auto numaNodes = get_available_numa_nodes(); int streams_num = _config.get_streams(); - int sub_streams_num = 0;//_config.get_sub_streams(); - // std::cout << "[ Impl ] " << _config.get_name() << " : " << streams_num << "\n"; + int sub_streams_num = 0; if (streams_num != 0) { std::copy_n(std::begin(numaNodes), std::min(streams_num, numaNodes.size()), diff --git a/src/plugins/intel_cpu/src/compiled_model.cpp b/src/plugins/intel_cpu/src/compiled_model.cpp index 45d05542d51..2c5ba010a10 100644 --- a/src/plugins/intel_cpu/src/compiled_model.cpp +++ b/src/plugins/intel_cpu/src/compiled_model.cpp @@ -62,18 +62,13 @@ CompiledModel::CompiledModel(const std::shared_ptr& model, // special case when all InferRequests are muxed into a single queue m_task_executor = m_plugin->get_executor_manager()->get_executor("CPU"); } else { - if (m_cfg.enableSubStreams) { - executor_confg = IStreamsExecutor::Config{"CPUMainStreamExecutor", - 1, - 1, - ov::hint::SchedulingCoreType::ANY_CORE, - false, - true}; - } else { - executor_confg = m_cfg.subStreamExecConfig.get_name() != "StreamsExecutor" - ? std::move(m_cfg.subStreamExecConfig) - : std::move(m_cfg.streamExecutorConfig); - } + executor_confg = m_cfg.enableSubStreams ? IStreamsExecutor::Config{"CPUMainStreamExecutor", + 1, + 1, + ov::hint::SchedulingCoreType::ANY_CORE, + false, + true} + : m_cfg.streamExecutorConfig; m_task_executor = m_plugin->get_executor_manager()->get_idle_cpu_streams_executor(executor_confg); } if (0 != m_cfg.streamExecutorConfig.get_streams()) { @@ -116,27 +111,27 @@ CompiledModel::CompiledModel(const std::shared_ptr& model, } else { CompiledModel::get_graph(); } - // std::cout << "m_hans_sub_compiled_models: " << m_cfg.enableSubStreams << ", " << m_cfg.streamExecutorConfig.get_sub_streams() << "\n"; + // std::cout << "m_has_sub_compiled_models: " << m_cfg.enableSubStreams << ", " << m_cfg.streamExecutorConfig.get_sub_streams() << "\n"; if (m_cfg.enableSubStreams) { m_cfg.enableSubStreams = false; - m_hans_sub_compiled_models = true; + m_has_sub_compiled_models = true; + auto sub_cfg = m_cfg; std::vector> sub_models; + auto streams_info_table = m_cfg.streamExecutorConfig.get_streams_info_table(); auto message = message_manager(); for (int i = 0; i < m_cfg.streamExecutorConfig.get_sub_streams(); i++) { - auto streams_info_table = m_cfg.streamExecutorConfig.get_streams_info_table(); - std::vector> info_table; - info_table.push_back(streams_info_table[i + 1]); - info_table[0][NUMBER_OF_STREAMS] = 1; - auto streamExecutorConfig = IStreamsExecutor::Config{"CPUStreamsExecutor", - 1, - 1, - ov::hint::SchedulingCoreType::ANY_CORE, - false, - true, - info_table, - m_cfg.streamsRankTable[i]}; - m_cfg.subStreamExecConfig = std::move(streamExecutorConfig); - sub_models.push_back(std::make_shared(model, plugin, m_cfg, loaded_from_cache)); + std::vector> sub_streams_table; + sub_streams_table.push_back(streams_info_table[i + 1]); + sub_streams_table[0][NUMBER_OF_STREAMS] = 1; + sub_cfg.streamExecutorConfig = IStreamsExecutor::Config{"CPUStreamsExecutor", + 1, + 1, + ov::hint::SchedulingCoreType::ANY_CORE, + false, + true, + sub_streams_table, + sub_cfg.streamsRankTable[i]}; + sub_models.push_back(std::make_shared(model, plugin, sub_cfg, loaded_from_cache)); } message->set_sub_compiled_models(sub_models); } @@ -205,7 +200,7 @@ std::shared_ptr CompiledModel::create_infer_request() co std::make_shared(std::static_pointer_cast(internal_request), get_task_executor(), get_callback_executor()); - if (m_hans_sub_compiled_models) { + if (m_has_sub_compiled_models) { auto message = message_manager(); auto sub_models = message->get_sub_compiled_models(); std::vector> requests; diff --git a/src/plugins/intel_cpu/src/compiled_model.h b/src/plugins/intel_cpu/src/compiled_model.h index c257c6147d1..75a9d047db2 100644 --- a/src/plugins/intel_cpu/src/compiled_model.h +++ b/src/plugins/intel_cpu/src/compiled_model.h @@ -74,7 +74,7 @@ private: */ GraphGuard::Lock get_graph() const; - bool m_hans_sub_compiled_models = false; + bool m_has_sub_compiled_models = false; }; } // namespace intel_cpu diff --git a/src/plugins/intel_cpu/src/config.h b/src/plugins/intel_cpu/src/config.h index 9728ddbb2ee..1df14213486 100644 --- a/src/plugins/intel_cpu/src/config.h +++ b/src/plugins/intel_cpu/src/config.h @@ -58,7 +58,6 @@ struct Config { size_t rtCacheCapacity = 0ul; #endif ov::threading::IStreamsExecutor::Config streamExecutorConfig; - ov::threading::IStreamsExecutor::Config subStreamExecConfig; int streams = 1; bool streamsChanged = false; int threads = 0; From 4ef10ab6f9b4a89d89b997ce8464ed836995ebab Mon Sep 17 00:00:00 2001 From: sunxiaoxia2022 Date: Fri, 17 May 2024 10:28:46 +0800 Subject: [PATCH 13/18] code style --- src/inference/src/dev/threading/cpu_message.cpp | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/src/inference/src/dev/threading/cpu_message.cpp b/src/inference/src/dev/threading/cpu_message.cpp index c2efa2555e4..a26290aafdf 100644 --- a/src/inference/src/dev/threading/cpu_message.cpp +++ b/src/inference/src/dev/threading/cpu_message.cpp @@ -70,7 +70,8 @@ void MessageManager::server_wait(int streams_num) { _msgCondVar.wait(lock); } std::swap(_messageQueue, msgQueue); - // std::cout << "server_wait receive: " << msgQueue[0].msg_type << " rank:" << msgQueue[0].rank.size() + // std::cout << "server_wait receive: " << msgQueue[0].msg_type << " rank:" << + // msgQueue[0].rank.size() // << " / " << msgQueue.size() << "\n"; } From a450308d6fd1e82fed6f6a250972770aabd66d07 Mon Sep 17 00:00:00 2001 From: sunxiaoxia2022 Date: Fri, 17 May 2024 10:38:18 +0800 Subject: [PATCH 14/18] remove unused value --- .../dev_api/openvino/runtime/threading/cpu_message.hpp | 1 - src/inference/src/dev/threading/cpu_message.cpp | 4 ++-- 2 files changed, 2 insertions(+), 3 deletions(-) diff --git a/src/inference/dev_api/openvino/runtime/threading/cpu_message.hpp b/src/inference/dev_api/openvino/runtime/threading/cpu_message.hpp index 79d16398fd5..a81ea9f8f71 100644 --- a/src/inference/dev_api/openvino/runtime/threading/cpu_message.hpp +++ b/src/inference/dev_api/openvino/runtime/threading/cpu_message.hpp @@ -62,7 +62,6 @@ public: std::vector> get_sub_infer_requests(); private: - int sub_streams; std::vector> _sub_compiled_models; std::vector> _sub_infer_requests; std::thread _serverThread; diff --git a/src/inference/src/dev/threading/cpu_message.cpp b/src/inference/src/dev/threading/cpu_message.cpp index a26290aafdf..398fc9fd3ee 100644 --- a/src/inference/src/dev/threading/cpu_message.cpp +++ b/src/inference/src/dev/threading/cpu_message.cpp @@ -32,7 +32,7 @@ std::vector MessageManager::wait_message(int cur_rank, int streams_ std::unique_lock lock(_readMutex); _readCondVar.wait(lock, [&] { // std::cout << "wait_" << cur_rank << " : " << _readQueue[cur_rank].size() << " / " << streams_num << "\n"; - return _readQueue[cur_rank].size() >= streams_num; + return static_cast(_readQueue[cur_rank].size()) >= streams_num; }); std::swap(_readQueue[cur_rank], messages_total); // std::cout << "wait_" << cur_rank << " " << _readQueue[cur_rank].size() << " end\n"; @@ -86,7 +86,7 @@ void MessageManager::server_wait(int streams_num) { while (!stop) { stop = true; for (int i = 0; i < streams_num; i++) { - if (_readQueue[i].size() == streams_num) { + if (static_cast(_readQueue[i].size()) == streams_num) { stop = false; } } From 5d521d0585dbadf6c1e754f237dee7b15bce4e4f Mon Sep 17 00:00:00 2001 From: sunxiaoxia2022 Date: Mon, 20 May 2024 17:01:26 +0800 Subject: [PATCH 15/18] remove inputs of server_wait --- .../runtime/threading/cpu_message.hpp | 9 ++-- .../src/dev/threading/cpu_message.cpp | 45 +++++++++++-------- src/plugins/intel_cpu/src/infer_request.cpp | 4 +- 3 files changed, 34 insertions(+), 24 deletions(-) diff --git a/src/inference/dev_api/openvino/runtime/threading/cpu_message.hpp b/src/inference/dev_api/openvino/runtime/threading/cpu_message.hpp index a81ea9f8f71..abfff73efd7 100644 --- a/src/inference/dev_api/openvino/runtime/threading/cpu_message.hpp +++ b/src/inference/dev_api/openvino/runtime/threading/cpu_message.hpp @@ -41,13 +41,13 @@ public: void send_message(const MessageInfo& msg_info); - std::vector wait_message(int cur_rank, int streams_num); + std::vector wait_message(int stream_id); void infer_wait(); - void reduce_wait(int cur_rank, int streams_num); + void reduce_wait(int stream_id); - void server_wait(int streams_num); + void server_wait(); void stop_server_thread(); @@ -61,7 +61,10 @@ public: void set_sub_infer_requests(std::vector> requests); std::vector> get_sub_infer_requests(); + int get_num_sub_streams(); + private: + int _num_sub_streams; std::vector> _sub_compiled_models; std::vector> _sub_infer_requests; std::thread _serverThread; diff --git a/src/inference/src/dev/threading/cpu_message.cpp b/src/inference/src/dev/threading/cpu_message.cpp index 398fc9fd3ee..a7491d6200f 100644 --- a/src/inference/src/dev/threading/cpu_message.cpp +++ b/src/inference/src/dev/threading/cpu_message.cpp @@ -27,14 +27,14 @@ void MessageManager::send_message(const MessageInfo& msg_info) { _msgCondVar.notify_all(); } -std::vector MessageManager::wait_message(int cur_rank, int streams_num) { +std::vector MessageManager::wait_message(int stream_id) { std::vector messages_total; std::unique_lock lock(_readMutex); _readCondVar.wait(lock, [&] { - // std::cout << "wait_" << cur_rank << " : " << _readQueue[cur_rank].size() << " / " << streams_num << "\n"; - return static_cast(_readQueue[cur_rank].size()) >= streams_num; + // std::cout << "wait_" << cur_rank << " : " << _readQueue[cur_rank].size() << " / " << _num_sub_streams << "\n"; + return static_cast(_readQueue[stream_id].size()) >= _num_sub_streams; }); - std::swap(_readQueue[cur_rank], messages_total); + std::swap(_readQueue[stream_id], messages_total); // std::cout << "wait_" << cur_rank << " " << _readQueue[cur_rank].size() << " end\n"; return messages_total; } @@ -44,21 +44,22 @@ void MessageManager::infer_wait() { _inferCondVar.wait(lock); } -void MessageManager::reduce_wait(int cur_rank, int streams_num) { +void MessageManager::reduce_wait(int stream_id) { std::unique_lock lock(_reduceMutex); - while (_reduceQueue[cur_rank] < streams_num) { + while (_reduceQueue[stream_id] < _num_sub_streams) { // std::cout << "reduce_wait_" << cur_rank << " " << _reduceQueue[cur_rank] << " end\n"; _reduceCondVar.wait(lock); } - _reduceQueue[cur_rank] = 0; + _reduceQueue[stream_id] = 0; } -void MessageManager::server_wait(int streams_num) { +void MessageManager::server_wait() { if (!_serverThread.joinable()) { - _readQueue.assign(streams_num, std::vector()); - _reduceQueue.assign(streams_num, 0); + assert(_num_sub_streams); + _readQueue.assign(_num_sub_streams, std::vector()); + _reduceQueue.assign(_num_sub_streams, 0); MsgType msg_type; - _serverThread = std::thread([&, streams_num]() { + _serverThread = std::thread([&]() { int count = 0; int reduce_count = 0; while (!_isServerStopped) { @@ -85,8 +86,8 @@ void MessageManager::server_wait(int streams_num) { bool stop = false; while (!stop) { stop = true; - for (int i = 0; i < streams_num; i++) { - if (static_cast(_readQueue[i].size()) == streams_num) { + for (int i = 0; i < _num_sub_streams; i++) { + if (static_cast(_readQueue[i].size()) == _num_sub_streams) { stop = false; } } @@ -94,25 +95,25 @@ void MessageManager::server_wait(int streams_num) { _readCondVar.notify_all(); } } - for (int i = 0; i < streams_num; i++) { + for (int i = 0; i < _num_sub_streams; i++) { std::lock_guard lock(_readMutex); _readQueue[i].push_back(rec_info); } _readCondVar.notify_all(); } else if (msg_type == CALL_BACK) { // CALL_BACK count++; - // std::cout << "server_wait CALL_BACK: " << count << "/" << streams_num << "\n"; - if (count == streams_num) { + // std::cout << "server_wait CALL_BACK: " << count << "/" << _num_sub_streams << "\n"; + if (count == _num_sub_streams) { _inferCondVar.notify_one(); count = 0; } } else if (msg_type == REDUCE) { // REDUCE reduce_count++; - // std::cout << "server_wait REDUCE: " << reduce_count << "/" << streams_num << "\n"; - if (reduce_count == streams_num) { + // std::cout << "server_wait REDUCE: " << reduce_count << "/" << _num_sub_streams << "\n"; + if (reduce_count == _num_sub_streams) { { std::lock_guard lock(_reduceMutex); - _reduceQueue.assign(streams_num, reduce_count); + _reduceQueue.assign(_num_sub_streams, reduce_count); } _reduceCondVar.notify_all(); reduce_count = 0; @@ -129,6 +130,8 @@ void MessageManager::server_wait(int streams_num) { void MessageManager::set_sub_compiled_models(std::vector> models) { _sub_compiled_models = models; + _num_sub_streams = static_cast(_sub_compiled_models.size()); + assert(_num_sub_streams); } std::vector> MessageManager::get_sub_compiled_models() { @@ -143,6 +146,10 @@ std::vector> MessageManager::get_sub_inf return _sub_infer_requests; } +int MessageManager::get_num_sub_streams() { + return _num_sub_streams; +} + void MessageManager::stop_server_thread() { MessageInfo msg_info; msg_info.msg_type = ov::threading::MsgType::QUIT; diff --git a/src/plugins/intel_cpu/src/infer_request.cpp b/src/plugins/intel_cpu/src/infer_request.cpp index e6aa2bfbc42..c5af7c9be44 100644 --- a/src/plugins/intel_cpu/src/infer_request.cpp +++ b/src/plugins/intel_cpu/src/infer_request.cpp @@ -114,7 +114,7 @@ void SyncInferRequest::infer() { throw_if_canceled(); // std::cout << "[ infer ] " << m_asyncRequest->m_has_sub_infers << "\n"; if (m_asyncRequest->m_has_sub_infers) { - message->server_wait(message->get_sub_infer_requests().size()); + message->server_wait(); ov::threading::MessageInfo msg_info; msg_info.msg_type = ov::threading::MsgType::START_INFER; ov::threading::Task task = [&] { @@ -655,7 +655,7 @@ void SyncInferRequest::sub_streams_infer() { requests[i]->set_tensor(input, tensor); } - requests[i]->set_callback([i, message](const std::exception_ptr& ptr) { + requests[i]->set_callback([message](const std::exception_ptr& ptr) { // std::cout << "set_callback------ " << i << "\n"; ov::threading::MessageInfo msg_info; msg_info.msg_type = ov::threading::MsgType::CALL_BACK; From 6ba1b2caf7cc75033a9a9f40642ca0f08edd65cc Mon Sep 17 00:00:00 2001 From: sunxiaoxia2022 Date: Mon, 20 May 2024 17:30:00 +0800 Subject: [PATCH 16/18] code style --- src/inference/src/dev/threading/cpu_message.cpp | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/src/inference/src/dev/threading/cpu_message.cpp b/src/inference/src/dev/threading/cpu_message.cpp index a7491d6200f..915b8064a0f 100644 --- a/src/inference/src/dev/threading/cpu_message.cpp +++ b/src/inference/src/dev/threading/cpu_message.cpp @@ -31,11 +31,12 @@ std::vector MessageManager::wait_message(int stream_id) { std::vector messages_total; std::unique_lock lock(_readMutex); _readCondVar.wait(lock, [&] { - // std::cout << "wait_" << cur_rank << " : " << _readQueue[cur_rank].size() << " / " << _num_sub_streams << "\n"; + // std::cout << "wait_" << stream_id << " : " << _readQueue[stream_id].size() << " / " << _num_sub_streams << + // "\n"; return static_cast(_readQueue[stream_id].size()) >= _num_sub_streams; }); std::swap(_readQueue[stream_id], messages_total); - // std::cout << "wait_" << cur_rank << " " << _readQueue[cur_rank].size() << " end\n"; + // std::cout << "wait_" << stream_id << " " << _readQueue[stream_id].size() << " end\n"; return messages_total; } @@ -47,7 +48,7 @@ void MessageManager::infer_wait() { void MessageManager::reduce_wait(int stream_id) { std::unique_lock lock(_reduceMutex); while (_reduceQueue[stream_id] < _num_sub_streams) { - // std::cout << "reduce_wait_" << cur_rank << " " << _reduceQueue[cur_rank] << " end\n"; + // std::cout << "reduce_wait_" << stream_id << " " << _reduceQueue[stream_id] << " end\n"; _reduceCondVar.wait(lock); } _reduceQueue[stream_id] = 0; From 5f0f37a6840ec5dad012340a2fd0d95b2ffa29d2 Mon Sep 17 00:00:00 2001 From: sunxiaoxia2022 Date: Mon, 20 May 2024 19:29:50 +0800 Subject: [PATCH 17/18] remove unused code --- .../runtime/threading/istreams_executor.hpp | 13 ------------- src/plugins/intel_cpu/src/graph.cpp | 3 --- 2 files changed, 16 deletions(-) diff --git a/src/inference/dev_api/openvino/runtime/threading/istreams_executor.hpp b/src/inference/dev_api/openvino/runtime/threading/istreams_executor.hpp index e96acbe8964..180a2b77a48 100644 --- a/src/inference/dev_api/openvino/runtime/threading/istreams_executor.hpp +++ b/src/inference/dev_api/openvino/runtime/threading/istreams_executor.hpp @@ -38,19 +38,6 @@ public: */ using Ptr = std::shared_ptr; - enum MsgType{ - TP, - START_INFER, - CALL_BACK - }; - - struct MessageInfo{ - MsgType msg_type; - std::vector rank; - void* buf; - Task task; - }; - /** * @brief Defines inference thread binding type */ diff --git a/src/plugins/intel_cpu/src/graph.cpp b/src/plugins/intel_cpu/src/graph.cpp index 93e7528e266..bfd6f949563 100644 --- a/src/plugins/intel_cpu/src/graph.cpp +++ b/src/plugins/intel_cpu/src/graph.cpp @@ -1094,9 +1094,6 @@ void Graph::PullOutputData(std::unordered_map>& void Graph::InferStatic(SyncInferRequest* request) { dnnl::stream stream(getEngine()); - const auto& cpuExecutor = context->getCPUStreamExecutor(); - // std::cout << "execute stream: " << cpuExecutor->get_stream_id() << "\n"; - for (const auto& node : executableGraphNodes) { VERBOSE(node, getConfig().debugCaps.verbose); PERF(node, getConfig().collectPerfCounters); From a131c1371f97646f055321522f0c8e97253d4fff Mon Sep 17 00:00:00 2001 From: sunxiaoxia2022 Date: Thu, 23 May 2024 10:35:28 +0800 Subject: [PATCH 18/18] change while to predicate --- src/inference/src/dev/threading/cpu_message.cpp | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/src/inference/src/dev/threading/cpu_message.cpp b/src/inference/src/dev/threading/cpu_message.cpp index 915b8064a0f..02722b2c342 100644 --- a/src/inference/src/dev/threading/cpu_message.cpp +++ b/src/inference/src/dev/threading/cpu_message.cpp @@ -68,9 +68,9 @@ void MessageManager::server_wait() { { // std::cout << "server_wait ........" << _isServerStopped << "\n"; std::unique_lock lock(_msgMutex); - while (_messageQueue.empty()) { - _msgCondVar.wait(lock); - } + _msgCondVar.wait(lock, [&] { + return !_messageQueue.empty(); + }); std::swap(_messageQueue, msgQueue); // std::cout << "server_wait receive: " << msgQueue[0].msg_type << " rank:" << // msgQueue[0].rank.size()