ZLToolKit源码框架
主要分为Thread、Poller、Network、Util四大部分。
Thread文件夹
TaskExecutor.h
有cpu负载计算,Task函数指针模板,任务执行器管理,管理任务执行线程池
ThreadLoadCounter
cpu负载计算器,基类,统计线程每一次的睡眠时长和工作时长,并记录样本,调用load计算cpu负载=工作时长/总时长。通过滑动窗口busy-rate计算,不是去读OS的真实cpu统计。
原理:
通过在线程进入阻塞前记一段 run_time,被唤醒后记一段 sleep_time,load() 再按最近窗口做比例计算。
//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;
};
- **向读者分发数据:**当
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);
}
- 写入存储层:
_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. 解锁
▼
(异步执行,各自独立)
线程A:DispatcherA->write(frame, true)
│
├── 遍历线程A的 reader_map,逐个回调 onRead (直接访问副本A里的历史GOP)
└── 写入副本A:_storage_of_A->write(frame, true)
(副本A开启新GOP,淘汰旧数据)
线程B:DispatcherB->write(frame, true)
│
├── 遍历线程B的 reader_map,逐个回调 onRead
└── 写入副本B:_storage_of_B->write(frame, true)
(副本B开启新GOP,淘汰旧数据)
c