前文

日志库spdlog(一) 源码面前,了无秘密 spdlog · Code Analysis 正如前文所说,本文将聚焦于spdlog在并发方面的工作,同时提点一下值得一提的技术。

线程安全

C++支持多线程编程,既然支持并发(并发或者并行是有区别的,但常规而言并发用的更多,读者明白就好),那spdlog库就天然存在重入数据竞争的问题。 比如,在registry层,存储一个name to logger的映射,此时如果并发写入,显然会有UB。 再比如,如果多个线程同时对一个logger写入日志,那此时如果没有保护,显然输出会乱掉。

spdlog为了支持并发,几乎在每一层都放了mutex来保护数据。

registry层

registry层主要有三个锁,如下:

1std::mutex logger_map_mutex_, flusher_mutex_;
2std::recursive_mutex tp_mutex_;

前两个锁结合代码都很好理解:前者用于保护name to logger结构,后者用于保护定时会写的结构(支持定时刷盘,多线程实现)。 在对logger执行操作时,此时不希望logger被更改,应当上锁,以达到快照功能:

1SPDLOG_INLINE void registry::set_level(level::level_enum log_level) {
2    // 如果要设置日志等级,此时就不允许其他线程修改logger
3    std::lock_guard<std::mutex> lock(logger_map_mutex_);
4    for (auto &l : loggers_) {
5        l.second->set_level(log_level);
6    }
7    global_log_level_ = log_level;
8}

最后一个则是用到了std::recursive_mutex,这是C++的可重入锁,也就是说可以被多个owner持有,内部存在计数,被锁上多少次就要被解锁多少次。 其作用暂且按下不表,等将异步时再说明。

logger层

logger负责整合sink,实际上的并发安全是由sink保证的。当然这里logger初始化的时候会有些风险,显然你不应该在两个线程内同时去初始化同一个logger

sink层

该层负责真正写入时的并发安全,但sink基类却没有锁,这就比较奇怪了,回过头我们再看看两种sink的实现:

 1// ansicolor_sink.h
 2template <typename ConsoleMutex>
 3class ansicolor_sink : public sink {
 4    ...
 5    using mutex_t = typename ConsoleMutex::mutex_t;
 6    mutex_t &mutex_;
 7    ...
 8}
 9
10using ansicolor_stdout_sink_mt = ansicolor_stdout_sink<details::console_mutex>;
11using ansicolor_stdout_sink_st = ansicolor_stdout_sink<details::console_nullmutex>;
12
13// basic_file_sink.h
14template <typename Mutex>
15class basic_file_sink final : public base_sink<Mutex> {
16};
17
18using basic_file_sink_mt = basic_file_sink<std::mutex>;
19using basic_file_sink_st = basic_file_sink<details::null_mutex>;
20
21// base_sink.h
22template <typename Mutex>
23class SPDLOG_API base_sink : public sink {
24    ...
25    Mutex mutex_;
26    ...
27};

对于不同的sink,都会定义一个mutex,以此实现多线程下的写入安全。 在使用部分,则是在log、flush等多处使用,保证临界区写入安全。 对于一个sink,提供了_mt _st两类后缀类型实例,传入了正常的mutex锁以及只实现了lock、unlock两个空方法的结构,有兴趣可以看看。 这样对于任意一个类型的sink,不必两种类型都实现一遍,空方法编译器是会优化掉的。

小结

spdlog在并发上做的很严谨,正如上文所说,锁对性能也有一定影响,锁基本也只加在了需要的地方。 显然,针对单线程场景,我们使用single thread类型的sink,针对多线程场景,则使用multi thread类型的,便捷切换,直接改个mt即可,你大概率也不会有任何心智负担。

异步

从之前的分析中可以看到,一套日志输出还是比较繁琐的,层级多,而且还需要主线程老老实实等待每一步做完,性能显然不佳。 对于这个问题,spdlog也提供了异步的解决方案,这里就是线程池。 以该程序为例

 1void async_logging_example() {
 2    // 业务线程只负责把日志放入队列;后台线程负责真正写文件。
 3    constexpr std::size_t queue_size = 8192;
 4    constexpr std::size_t backend_threads = 1;
 5    // 线程池任务队列大小和线程数量
 6    spdlog::init_thread_pool(queue_size, backend_threads);
 7
 8    auto logger = spdlog::create_async<spdlog::sinks::basic_file_sink_mt>(
 9        "async_file", "logs/async-log.txt", true);
10    logger->set_pattern("[%H:%M:%S.%e] [%n] [thread %t] %v");
11
12    constexpr int producer_count = 4;
13    constexpr int messages_per_producer = 100;
14    std::vector<std::thread> producers;
15    producers.reserve(producer_count);
16
17    for (int producer_id = 0; producer_id < producer_count; ++producer_id) {
18        producers.emplace_back([logger, producer_id] {
19            for (int message_id = 0; message_id < messages_per_producer; ++message_id) {
20                logger->info("producer={}, message={}", producer_id, message_id);
21            }
22        });
23    }
24
25    for (auto &producer: producers) {
26        producer.join();
27    }
28
29    logger->flush();
30}

我们分析上述每一步,大体上就能理解spdlog是如何通过异步来提高性能的。

线程池

第一步就是初始化一个线程池,上面设定了任务队列的大小(在这里称呼为消息队列比较准确),以及线程数量。 线程池的构造函数如下:

 1// thread_pool-inl.h
 2SPDLOG_INLINE thread_pool::thread_pool(size_t q_max_items,
 3                                       size_t threads_n,
 4                                       std::function<void()> on_thread_start,
 5                                       std::function<void()> on_thread_stop)
 6    : q_(q_max_items) {
 7    if (threads_n == 0 || threads_n > 1000) {
 8        throw_spdlog_ex(
 9            "thread_pool(): invalid threads_n param (valid "
10            "range is 1-1000)");
11    }
12    for (size_t i = 0; i < threads_n; i++) {
13        // 套用匿名函数创建线程
14        // 还用function塞俩回调进去了
15        threads_.emplace_back([this, on_thread_start, on_thread_stop] {
16            on_thread_start();
17            this->thread_pool::worker_loop_();
18            on_thread_stop();
19        });
20    }
21}
22
23void SPDLOG_INLINE thread_pool::worker_loop_() {
24    // 核心处理loop
25    while (process_next_msg_()) {
26    }
27}
28
29// process next message in the queue
30// returns true if this thread should still be active (while no terminated msg was received)
31bool SPDLOG_INLINE thread_pool::process_next_msg_() {
32    // 从消息队列中取出一个消息
33    async_msg incoming_async_msg;
34    q_.dequeue(incoming_async_msg);
35
36    switch (incoming_async_msg.msg_type) {
37        // 再根据msg type来做具体的行为
38        case async_msg_type::log: {
39            incoming_async_msg.worker_ptr->backend_sink_it_(incoming_async_msg);
40            return true;
41        }
42        case async_msg_type::flush: {
43            incoming_async_msg.worker_ptr->backend_flush_();
44            return true;
45        }
46
47        case async_msg_type::terminate: {
48            return false;
49        }
50
51        default: {
52            assert(false);
53        }
54    }
55
56    return true;
57}

至上而下分析到这里,我们也有了一些异步的运行猜想了:

  1. 业务线程创建async_msg,投递到消息队列里:thread_pool::post_log
  2. 再由每个线程来处理消息,根据不同的消息类型,反过来调用消息类型里所携带的logger指针继续我们的sink调用

async_logger

回看示例程序,在创建logger时,使用的是spdlog::create_async这个函数:

 1// async.h
 2template <async_overflow_policy OverflowPolicy = async_overflow_policy::block>
 3struct async_factory_impl {
 4    template <typename Sink, typename... SinkArgs>
 5    static std::shared_ptr<async_logger> create(std::string logger_name, SinkArgs &&...args) {
 6        // ...
 7        auto new_logger = std::make_shared<async_logger>(std::move(logger_name), std::move(sink),
 8                                                         std::move(tp), OverflowPolicy);
 9        registry_inst.initialize_logger(new_logger);
10        // ...
11        return new_logger;
12    }
13};
14
15using async_factory = async_factory_impl<async_overflow_policy::block>;
16using async_factory_nonblock = async_factory_impl<async_overflow_policy::overrun_oldest>;
17
18template <typename Sink, typename... SinkArgs>
19inline std::shared_ptr<logger> create_async(std::string logger_name,
20                                                    SinkArgs &&...sink_args) {
21    return async_factory::create<Sink>(std::move(logger_name),
22                                       std::forward<SinkArgs>(sink_args)...);
23}

返回的logger类型也是async_logger,我们不禁要问,和普通的logger有什么不同呢,这才是真正的关键:

1void async_logger::sink_it_(const details::log_msg &msg){
2    SPDLOG_TRY{if (auto pool_ptr = thread_pool_.lock()){
3        pool_ptr -> post_log(shared_from_this(), msg, overflow_policy_);
4}
5else {
6    throw_spdlog_ex("async log: thread pool doesn't exist anymore");
7}
8}

显而易见,重写了sink_it_方法!通过weak_ptr拿到线程池,并调用上文提到的post_log,把一个普通的msg(复用的同一套格式化、创建msg的逻辑)转换成async_msg。

post_log

又回到了线程池的实现中,post_log的实现如下:

 1// thread_pool-inl.h
 2void SPDLOG_INLINE thread_pool::post_log(async_logger_ptr &&worker_ptr,
 3                                         const details::log_msg &msg,
 4                                         async_overflow_policy overflow_policy) {
 5    async_msg async_m(std::move(worker_ptr), async_msg_type::log, msg);
 6    post_async_msg_(std::move(async_m), overflow_policy);
 7}
 8
 9void SPDLOG_INLINE thread_pool::post_async_msg_(async_msg &&new_msg,
10                                                async_overflow_policy overflow_policy) {
11    if (overflow_policy == async_overflow_policy::block) {
12        q_.enqueue(std::move(new_msg));
13    } else if (overflow_policy == async_overflow_policy::overrun_oldest) {
14        q_.enqueue_nowait(std::move(new_msg));
15    } else {
16        assert(overflow_policy == async_overflow_policy::discard_new);
17        q_.enqueue_if_have_room(std::move(new_msg));
18    }
19}

刚刚一直有意地忽略了policy这个概念,显然只要是队列,就会有满的时候,我们需要通过不同的策略来对其进行处理,这里就是通过调用q_的不同方法来进行处理。 spdlog提供了三种不同的策略来处理队列满的情况:

1enum class async_overflow_policy {
2    block,           // Block until message can be enqueued
3    overrun_oldest,  // Discard oldest message in the queue if full when trying to
4                     // add new item.
5    discard_new      // Discard new message if the queue is full when trying to add new item.
6};

另外值得一提的是,这里通过普通的log_msg创建了async_msg,不妨也可以看看async_msg和log_msg有什么不同:

 1// thread_pool-inl.h
 2struct async_msg : log_msg_buffer {
 3    async_msg_type msg_type{async_msg_type::log};
 4    async_logger_ptr worker_ptr;
 5    // ...
 6};
 7// log_msg_buffer
 8class SPDLOG_API log_msg_buffer : public log_msg {
 9    memory_buf_t buffer;
10    void append_source();
11    void update_string_views();
12    // ...
13};

这里很有意思,log_msg里所拥有的其实是视图,因为是同步,所以能保证数据源一直在栈上。 但在这种异步的场景就行不通了,所以有了log_msg_buffer来深拷贝一波真实数据。

再来一个问题,为什么这里会存async_logger_ptr worker_ptr;? 因为上面的sink_it_只是提交到消息队列里,并没有真正的落盘,真实落盘还是需要async_logger来辅助。

 1void async_logger::backend_sink_it_(const details::log_msg &incoming_log_msg) {
 2    for (auto &sink : sinks_) {
 3        if (sink->should_log(incoming_log_msg.level)) {
 4            SPDLOG_TRY { sink->log(incoming_log_msg); }
 5            SPDLOG_LOGGER_CATCH(incoming_log_msg.source)
 6        }
 7    }
 8
 9    if (should_flush_(incoming_log_msg)) {
10        backend_flush_();
11    }
12}

这就和我们之前看到的普通logger的行为是一样的了。在提交时,只是完成了logger创建msg的任务,还没有完成写入落盘的任务。

mpmc_blocking_queue

最后一块内容就是消息队列的实现。 这里的mpmc是multi-producer、multi-consumer的含义。 由于这个线程池的任务高度一致,所以这里也没有任务队列这一说(普通线程池缓存函数对象)。 那么这里线程同步的任务肯定是由mpmc_blocking_queue来实现:

 1template <typename T>
 2class mpmc_blocking_queue {
 3public:
 4    using item_type = T;
 5    // try to enqueue and block if no room left
 6    void enqueue(T &&item) {
 7        {
 8            std::unique_lock<std::mutex> lock(queue_mutex_);
 9            // 这里需要等待
10            pop_cv_.wait(lock, [this] { return !this->q_.full(); });
11            q_.push_back(std::move(item));
12        }
13        push_cv_.notify_one();
14    }
15
16    // enqueue immediately. overrun oldest message in the queue if no room left.
17    void enqueue_nowait(T &&item) {
18        {
19            std::unique_lock<std::mutex> lock(queue_mutex_);
20            q_.push_back(std::move(item));
21        }
22        push_cv_.notify_one();
23    }
24
25    void enqueue_if_have_room(T &&item) {
26        bool pushed = false;
27        {
28            std::unique_lock<std::mutex> lock(queue_mutex_);
29            if (!q_.full()) {
30                q_.push_back(std::move(item));
31                pushed = true;
32            }
33        }
34
35        if (pushed) {
36            push_cv_.notify_one();
37        } else {
38            ++discard_counter_;
39        }
40    }
41
42    // blocking dequeue without a timeout.
43    void dequeue(T &popped_item) {
44        {
45            std::unique_lock<std::mutex> lock(queue_mutex_);
46            push_cv_.wait(lock, [this] { return !this->q_.empty(); });
47            popped_item = std::move(q_.front());
48            q_.pop_front();
49        }
50        pop_cv_.notify_one();
51    }
52
53private:
54    std::mutex queue_mutex_;
55    std::condition_variable push_cv_;
56    std::condition_variable pop_cv_;
57    circular_q<T> q_;
58    std::atomic<size_t> discard_counter_{0};
59};

这显然就是一个常规的线程安全的消息队列,circular_q暂且按下不表,可以看到有三种方式的enqueue,对应着上面说的三种不同policy:阻塞、覆盖、丢弃。 其实借助了两个条件变量,例如阻塞场景下,等待着队列非空才会将新msg入队列,同时调用push_cv_来通知所有等待的线程有消息可以消费。 这也是条件变量用于线程同步的精妙所在,锁用于保护数据成员,而条件变量则可以借助锁来做到多线程之间的同步,减少CPU消耗。

circular_q

为了减少篇幅,这里就不贴具体实现,无非是一个简单的环状数组,需要注意以下几点:

  1. 多放一个位置,用于标识full的状态;
  2. 一切操作后都要+1对max_size来取模

小结

至此,我们可以看到spdlog是如何支持异步的,又是如何做到并发安全的。 其中所用到的知识也并没有任何黑魔法,没有跳出C++语法的范畴,其中有许多不错的技巧值得一提。 想想还是将其放在下一章好了,下一节做个简短的实验,来直观感受一下异步和同步之间的区别。

benchmark

考虑这样一个场景:同样对单文件输出1kw条日志,内容完全相同,一个同步,一个异步,异步使用线程池队列大小为65536,线程数量为1(通过同一个sink写同一个文件,没必要开俩线程) 所得结果如下:

其中caller-path是异步提交消息过程,end-to-end其实是后台线程写入过程。

同步的性能居然比异步要强,这怎么可能呢? 其实是合理的,回想一下同步和异步的过程:只要写入文件的消耗比创建深拷贝、提交到任务队列等所消耗的时间更少,那同步注定比异步快。

这里不禁要问,那异步绕这么大一圈,有什么意义呢? 其实是有的,上述场景其实不太好,设想这样一个场景,一个多线程程序,N个线程同时写一个文件,如果此时还是同步,那所有人都会抢锁,业务代码甚至会被卡死。 此时合理的做法就是异步提交,业务线程只负责投递消息,由线程池来处理无效的锁等待、IO等待等过程。

可以看到,提交的速度快了不是一星半点。

这里的性能测试不太严谨,只是测测最简单的场景。 但不难发现,不要盲目使用异步,只有确定业务代码需要异步后才可以上,否则只是空谈。

总结

工作太消耗精力,居然卡了大半个月,好在最后benchmark部分让codex抬了一手,不然更是懒得去细看怎么写benchmark。 我想这一篇写下来最大的感触就是写并发、异步程序实在是不容易,不仅要考虑锁的粒度,还要考虑数据的生命周期。 另外一个则是异步也未必比同步快,最大的优势就是不用管忙等待,对于性能敏感的业务代码而言,异步仍旧是上上之选。 日志也并非太重要的数据,所以异步场景下为了不阻塞,甚至允许覆盖/丢弃,当然这里也要考虑清楚你的应用场景,万一线上出了个bug,刚好丢了重要日志,那就彻底麻了。

下一篇会谈谈各个技术点,也是我的舒适区,开搞!