NOTE / 2019/10/27

C++ 并发编程[Part 1]

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

在某些情况下,多线程可以大幅改善程序执行时间。如在SLAM中,通常有两个线程,一个线程用来实现高频率的里程计或定位,另一个线程用来进行低频的地图构建或者量测信息处理。这样大计算量的地图构建不会阻塞高频的定位输出。从C++11,标准库开始支持多线程编程。本文,对C++标准库中的多线程模块进行总结。本文内容摘抄自《C++标准库》。

1. 高级接口async和Futhure

假设,我们要计算下面四个函数的和。

r=f1()+f2()+f3()+f4();r = f_1() +f_2() + f_3() + f_4(); \\

在单线程中,会序列化的计算f(),总的运算时间就是这四个函数运算时间的总和+加法的运算时间。

但是,如果使用多线程技术,我们可以并行计算这四个函数,然后最后一起加起来,这样程序的耗时就是那个耗时最大函数的计算时间+加法的运算时间。例如这种的运算场景,多线程便可以大幅度提高运算效率。我们可以使用async启动一个线程,使用future,在未来获取线程返回的结果。程序如下:

#include <chrono>
#include <future>
#include <iostream>
#include <thread>

double function(const double var) {
    std::this_thread::sleep_for(std::chrono::milliseconds(1000));
    return var;
}

int main(int argc, char** agrv) {
    // Compute f1 + f2 + f3 + f4 using one thread.
    std::chrono::steady_clock::time_point t1 = std::chrono::steady_clock::now();
    
    double result = function(1.) + function(2.) + function(3.) + function(4.);
    
    std::chrono::steady_clock::time_point t2 = std::chrono::steady_clock::now();
    double time_cost = std::chrono::duration_cast<std::chrono::duration<double>>(t2 - t1).count() * 1000.;
    std::cout << std::fixed << "Sigal thread result: " << result << " and time cost: " << time_cost << std::endl;

    // Compute f1 + f2 + f3 + f4 using four threads.
    t1 = std::chrono::steady_clock::now();
    
    std::future<double> f1(std::async(std::launch::async, function, 1.));
    auto f2 = std::async(std::launch::async, function, 2.);
    auto f3 = std::async(std::launch::async, function, 3.);
    result = function(4.) + f1.get() + f2.get() + f3.get();

    t2 = std::chrono::steady_clock::now();
    time_cost = std::chrono::duration_cast<std::chrono::duration<double>>(t2 - t1).count() * 1000.;
    std::cout << std::fixed << "Multi-threads result: " << result << " and time cost: " << time_cost << std::endl;

    return 1;
}

结果如下:

文章配图

可以看出,单线程处理时间几乎是多线程的4倍。

async和future、shared_future的具体使用方法请参考《C++标准库》

2. 线程同步与并发问题

在上述例子中,我们的四个线程之间都是独立进行的,并没有进行数据的交换。而在SLAM中,前端的定位线程与后端定位之间要进行地图的共享。在其他一些应用中,多线程之间的数据线程也是很常见的。

The only safe way to concurrently access the same data by multiple threads without synchronization is when ALL threads only READ the data

如果在多个线程之中对同一个数据进行读、写就会出现数据竞争问题。这种情况可以使用mutex和Lock解决。如下,我们对一个数据分别再两个线程中同时进行读写。

#include <chrono>
#include <iostream>
#include <mutex>
#include <random>
#include <thread>

double g_common_var = 0;
std::mutex g_common_var_mutex;

void Thread() {
    std::random_device rd; 
    std::mt19937 gen(rd());
    std::uniform_real_distribution<> dis(0., 1000.);

    for (size_t i = 0; i < 100; ++i) {
        double random_common = dis(gen);
        {   
            // or std::unique_lock<std::mutex> ul(g_common_var_mutex);
            std::lock_guard<std::mutex> lg(g_common_var_mutex);
            std::cout << std::fixed << "[Read] The common var in thread " << std::this_thread::get_id() << " is " << g_common_var << std::endl;
            g_common_var = random_common;
            std::cout << "[Write] The common var in thread " << std::this_thread::get_id() << " is " << g_common_var << std::endl << std::endl;
        }
        
        std::this_thread::sleep_for(std::chrono::milliseconds(500));
    }
}

int main(int argc, char** argv) {
    std::thread t1(Thread);
    std::thread t2(Thread);
 
    t1.join();
    t2.join(); 

    return 1;
}

结果如下:

文章配图

3. 数据生产与消费 - Condition variable

在SLAM中,我们经常由前端里程计部分挑选一些关键帧,发送给后端建图模块进行使用。也就是说1)关键帧在前端和后端线程中共享,2)在发送给后端的同时,我们还需要通知后端一声,让它开始处理数据。

这两项功能可以由condition variable实现,如下:

#include <atomic>
#include <chrono>
#include <condition_variable>
#include <iostream>
#include <queue>
#include <random>
#include <thread>
#include <mutex>
#include <memory>

class BackEnd {
public: 
    BackEnd() {
        thread_ptr_ = std::make_unique<std::thread>(&BackEnd::Process, this);
    }

    void AddData(const int data) {
        {
            std::lock_guard<std::mutex> lg(data_queue_mutex_);
            data_queue_.push(data);
        }
        data_queue_cond_var_.notify_one();
    }

    void Process() {
        thread_running_.store(true);
        while (thread_running_.load()) {
            int data;
            {   
                // Wait for data.
                std::unique_lock<std::mutex> ul(data_queue_mutex_);
                data_queue_cond_var_.wait(ul, [this]{ return !data_queue_.empty(); });
                data = data_queue_.front();
                data_queue_.pop();
            } // Release lock.

            // Process data.
            std::cout << "[BackEnd]: Recive data: " << data << std::endl << std::endl;
        }
    }

private:
    std::atomic<bool> thread_running_;
    
    std::mutex data_queue_mutex_;
    std::queue<int> data_queue_;
    std::condition_variable data_queue_cond_var_;

    std::unique_ptr<std::thread> thread_ptr_;
};

int main(int argc, char** argv) {
    // Start the back-end thread.
    BackEnd back_end;
    
    // This is the front-end.
    std::default_random_engine generator;
    std::uniform_int_distribution<int> distribution(0, 1000);
    for (size_t i = 0; i < 10; ++i) {
        // Produce data.
        int random_var = distribution(generator);

        // Send data to the back-end.
        back_end.AddData(random_var);

        std::cout << "[FrontEnd]: Add var " << random_var << std::endl;
        std::this_thread::sleep_for(std::chrono::milliseconds(500));
    }

    return 1;
}

结果如下:

文章配图