ZLToolKit源码框架

主要分为Thread、Poller、Network、Util四大部分。

Thread文件夹

TaskExecutor.h

cpu负载计算,Task函数指针模板,任务执行器管理,管理任务执行线程池

ThreadLoadCounter

​ cpu负载计算器,基类,统计线程每一次的睡眠时长和工作时长,并记录样本,调用load计算cpu负载=工作时长/总时长。通过滑动窗口busy-rate计算,不是去读OS的真实cpu统计。

原理:

通过在线程进入阻塞前记一段 run_time,被唤醒后记一段 sleep_timeload() 再按最近窗口做比例计算。

//startSleep() 记录的是“刚刚结束的一段运行时间”。
void ThreadLoadCounter::startSleep() {
    lock_guard<mutex> lck(_mtx);
    _sleeping = true;
    auto current_time = getCurrentMicrosecond();
    auto run_time = current_time - _last_wake_time;//线程准备休眠时,说明上一段活跃执行区间结束了
    _last_sleep_time = current_time; //为下一段 sleep 计时做准备
    _time_list.emplace_back(run_time, false);//表示这段时间不是睡眠,而是运行
    if (_time_list.size() > _max_size) {
        _time_list.pop_front();
    }
}
//sleepWakeUp() 记录的是“刚刚结束的一段休眠时间”。
void ThreadLoadCounter::sleepWakeUp() {
    lock_guard<mutex> lck(_mtx);
    _sleeping = false;
    auto current_time = getCurrentMicrosecond();
    auto sleep_time = current_time - _last_sleep_time;//计算睡了多久,相当于空闲时间
    _last_wake_time = current_time//更新为当前时刻。开始新的运行区间计时
    _time_list.emplace_back(sleep_time, true);
    if (_time_list.size() > _max_size) {
        _time_list.pop_front();
    }
}

Util文件夹

RingBuffer

RingBuffer是由多个类组成,分为两大功能:存储和数据分发。

存储功能由类RingStorage实现,是数据存储类,它是一个循环队列,有最大容量定义,从尾部插入最新数据,当队列满了,从头部删除老数据。

在RingBuffer类里中的数据结构,以EventPoller的指针作为Key,

template <typename T>
class RingBuffer public std::enable_shared_from_this<RingBuffer<T>>  {
public:
    void write(T in, bool is_key = true) {
        if (_delegate) {
            _delegate->onWrite(std::move(in), is_key);
            return;
        }

        LOCK_GUARD(_mtx_map);
        for (auto &pr : _dispatcher_map) {
            auto &second = pr.second;
            //切换线程后触发onRead事件  [AUTO-TRANSLATED:4ca6647d]
            //Switch thread and trigger onRead event
            pr.first->async([second, in, is_key]() mutable { second->write(std::move(in), is_key); }, false);
        }
        _storage->write(std::move(in), is_key);
    }
private:
std::unordered_map<EventPoller::Ptr, typename RingReaderDispatcher::Ptr, HashOfPtr> _dispatcher_map;
};
  1. **向读者分发数据:**当RingReaderDispatcher收到write调用时,它已处于对应的EventPoller线程。RingReaderDispatcher会遍历其管理的所有读者句柄(_RingReader),并逐一调用reader->onRead(in, is_key)实时分发给该 poller 上的 reader,_storage->write(std::move(in), is_key); 更新该 poller 自己的 GOP cache
// RingReaderDispatcher::write() 的核心代码段
void write(T in, bool is_key = true) {
    // 遍历本分发器管理的所有读者,直接回调
    for (auto it = _reader_map.begin(); it != _reader_map.end();) {
        auto reader = it->second.lock();
        if (reader) {
            reader->onRead(in, is_key); // 最终会调用业务回调 _read_cb
            ++it;
        }
    }
    // 通知存储层
    _storage->write(std::move(in), is_key);
}
  1. 写入存储层_RingStorage决定帧的缓存命运

所以整体模型是:

RingBuffer::_storage
    作用:主 GOP cache,新 poller attach 时用来 clone

RingReaderDispatcher::_storage
    作用:该 poller 线程内的 GOP cache,供该 poller  reader flushGop 使用

RingReader
    作用:实际观看者/消费者,只能在自己的 poller 线程里操作
////////////////
RingBuffer::write(frame, is_key=true)

├── 1. 锁定 _mtx_map,遍历 _dispatcher_map
     ├── 对线程A的分发器:通过 pollerA->async 将任务抛到线程A
     └── 对线程B的分发器:通过 pollerB->async 将任务抛到线程B

├── 2. 写入主存储:_storage->write(frame, true)  (主存储也更新)

└── 3. 解锁
            
(异步执行,各自独立)
线程ADispatcherA->write(frame, true)

├── 遍历线程A的 reader_map,逐个回调 onRead (直接访问副本A里的历史GOP)
└── 写入副本A_storage_of_A->write(frame, true)  
            (副本A开启新GOP,淘汰旧数据)
线程BDispatcherB->write(frame, true)

├── 遍历线程B的 reader_map,逐个回调 onRead
└── 写入副本B_storage_of_B->write(frame, true)
            (副本B开启新GOP,淘汰旧数据)

c