NOTE / 10/29/2019

C++ Concurrent Programming: Part 2

EngineeringTechnical NotesSLAMVIOsensor fusion

In most systems, assigning every task to a particular thread is impractical. Existing concurrency can instead be used for concurrent execution. A thread pool provides this: submitted tasks are queued, worker threads take tasks from the queue, and each worker takes another task after completing one. — C++ Concurrency in Action

The following is a simple thread-pool implementation:

#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;
}