前文
日志库spdlog(一) 源码面前,了无秘密 正如前文所说,本文将聚焦于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}
至上而下分析到这里,我们也有了一些异步的运行猜想了:
- 业务线程创建async_msg,投递到消息队列里:
thread_pool::post_log - 再由每个线程来处理消息,根据不同的消息类型,反过来调用消息类型里所携带的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
为了减少篇幅,这里就不贴具体实现,无非是一个简单的环状数组,需要注意以下几点:
- 多放一个位置,用于标识full的状态;
- 一切操作后都要+1对max_size来取模
小结
至此,我们可以看到spdlog是如何支持异步的,又是如何做到并发安全的。 其中所用到的知识也并没有任何黑魔法,没有跳出C++语法的范畴,其中有许多不错的技巧值得一提。 想想还是将其放在下一章好了,下一节做个简短的实验,来直观感受一下异步和同步之间的区别。
benchmark
考虑这样一个场景:同样对单文件输出1kw条日志,内容完全相同,一个同步,一个异步,异步使用线程池队列大小为65536,线程数量为1(通过同一个sink写同一个文件,没必要开俩线程)
所得结果如下:
其中caller-path是异步提交消息过程,end-to-end其实是后台线程写入过程。
同步的性能居然比异步要强,这怎么可能呢? 其实是合理的,回想一下同步和异步的过程:只要写入文件的消耗比创建深拷贝、提交到任务队列等所消耗的时间更少,那同步注定比异步快。
这里不禁要问,那异步绕这么大一圈,有什么意义呢?
其实是有的,上述场景其实不太好,设想这样一个场景,一个多线程程序,N个线程同时写一个文件,如果此时还是同步,那所有人都会抢锁,业务代码甚至会被卡死。
此时合理的做法就是异步提交,业务线程只负责投递消息,由线程池来处理无效的锁等待、IO等待等过程。
可以看到,提交的速度快了不是一星半点。
这里的性能测试不太严谨,只是测测最简单的场景。 但不难发现,不要盲目使用异步,只有确定业务代码需要异步后才可以上,否则只是空谈。
总结
工作太消耗精力,居然卡了大半个月,好在最后benchmark部分让codex抬了一手,不然更是懒得去细看怎么写benchmark。 我想这一篇写下来最大的感触就是写并发、异步程序实在是不容易,不仅要考虑锁的粒度,还要考虑数据的生命周期。 另外一个则是异步也未必比同步快,最大的优势就是不用管忙等待,对于性能敏感的业务代码而言,异步仍旧是上上之选。 日志也并非太重要的数据,所以异步场景下为了不阻塞,甚至允许覆盖/丢弃,当然这里也要考虑清楚你的应用场景,万一线上出了个bug,刚好丢了重要日志,那就彻底麻了。
下一篇会谈谈各个技术点,也是我的舒适区,开搞!