NOTE / 2019/10/29

C++ 并发编程[Part 2]

工程技术笔记SLAMVIO传感器融合

在大多数系统中,将每个任务指定给某个线程时不切实际的,不过可以利用现有的并发性,进行并发执行。线程池就提供了这样的功能,提交到线程池中的任务并发执行,提到的任务将会挂在任务队列上。队列中的每个任务都会被池中的工作线程获取,当一个任务执行完成后,到队列中获取下一个任务执行。——《C++ Concurrenty in Action》

#include <atomic>
#include <condition_variable>
#include <functional>
#include <iostream>
#include <mutex>
#include <queue>
#include <thread>
#include <vector>

class ThreadPool {
public:
    ThreadPool(const int thread_num) {
        threads_.reserve(thread_num);
        for (size_t i = 0; i < thread_num; ++i) {
            threads_.emplace_back(&ThreadPool::Thread, this);
        }
    }

    ~ThreadPool() {
        // Stop all threads in the pool.
        done_.store(true);
        thread_cond_var_.notify_all();
        for (size_t i = 0; i < threads_.size(); ++i) {
            threads_.at(i).join();
        } 

        std::cout << "[Threads stoped]\n";
    }

    void Submit(const std::function<void()>& work) {
        {
            // Add work to the queue.
            std::lock_guard<std::mutex> lg(mutex_);
            works_.push(work);
        }

        // Tell all threads to work. 
        thread_cond_var_.notify_all();
    }

private:
    void Thread() {
        done_.store(false);
        while (true) {
            // Wait for work.
            std::function<void()> work;
            {
                std::unique_lock<std::mutex> ul(mutex_);
                thread_cond_var_.wait(ul, [this] { return !works_.empty() || done_.load(); });
                if (done_.load()) {
                    break;
                }

                work = works_.front();
                works_.pop();
            }           

            // Do work.
            work();
        }
    }

    std::atomic_bool done_;
    std::mutex mutex_;
    std::condition_variable thread_cond_var_;
    std::queue<std::function<void()>> works_;
    std::vector<std::thread> threads_;
};

class Task {
public:
    
    void PrintText(const std::string& text) {
        std::this_thread::sleep_for(std::chrono::milliseconds(100));
        
        {
            std::lock_guard<std::mutex> lg(mutex_);
            std::cout << "Thread " << std::this_thread::get_id() << " do task " << text << std::endl;
        }
    }

    std::mutex mutex_;
};

int main (int argc, char** argv) {
    ThreadPool thread_pool(2);

    Task task;
    for (size_t i = 0; i < 10; ++i) {
        thread_pool.Submit(std::bind(&Task::PrintText, &task, std::to_string(i)));
        std::cout << "Submit task: " << i << std::endl;
    }

    std::this_thread::sleep_for(std::chrono::seconds(5));
    return 1;
}