ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

【Linux】手写日志与固定线程池:任务队列、工作线程和安全退出

【Linux】手写日志与固定线程池:任务队列、工作线程和安全退出 个人主页爱和冰阔乐专栏传送门《数据结构与算法》 、C、 《Linux操作系统》学习方向C方向学习爱好者⭐人生格言得知坦然 失之淡然博主简介文章目录前言本文使用的接口与范围一、先把日志模块准备好1.1 为什么不只使用 printf1.2 生成日志与刷新日志分开1.3 为什么日志输出也要加锁1.4 时间函数的线程安全二、实现一个简单日志器2.1 输出策略2.2 Logger 与临时消息对象2.3 简单测试三、为什么需要线程池3.1 来一个任务就创建一个线程的问题3.2 线程池仍然是生产者消费者模型3.3 固定线程数怎么选四、固定线程池完整实现4.1 线程池需要哪些状态4.2 ThreadPool.hpp4.3 为什么使用 unique_lock4.4 为什么任务必须在锁外执行4.5 为什么捕获任务异常五、线程池怎样安全退出5.1 工作线程的三种状态5.2 Stop 的执行顺序5.3 排空退出和立即退出5.4 不要让工作线程销毁自己的线程池六、运行线程池6.1 测试代码6.2 可以继续扩展哪些功能总结参考资料前言前面已经把互斥锁、条件变量和生产者消费者模型写完了这一篇把它们真正组合起来先实现一个够用的日志模块再实现固定数量线程的线程池。这里不追求一次写出工业级日志库和通用线程池。我们主要把几个关键问题想清楚任务为什么不能来一个就创建一个线程工作线程没有任务时怎么等待任务应该在锁内还是锁外执行以及线程池析构时怎样让所有线程正常退出。本文使用的接口与范围示例使用 C17 的std::thread、std::condition_variable和std::filesystem编译时需要添加-pthread。文中的线程池用于理解任务队列、等待谓词和停止流程不是直接替代成熟线程池库的生产版本。若编译器较旧std::filesystem可能还需要额外链接选项应以实际工具链为准。一、先把日志模块准备好1.1 为什么不只使用printf以前写测试代码时直接使用printf或cout没有问题。但项目运行时间一长只在显示器上打印很难排查问题终端关闭后消息就没了也不知道某条输出来自哪个文件、哪一行、哪个线程。一条实用日志至少应包含时间日志等级具体内容。文件名、行号、进程 ID 和线程 ID 不是每个场景都必须有但调试并发程序时非常有用。常见等级可以这样理解等级适合记录的内容DEBUG调试阶段需要的细节INFO程序正常运行过程中的关键事件WARNING出现异常迹象但程序还能继续工作ERROR当前操作失败需要排查FATAL关键功能无法继续程序通常需要退出等级名称只是一套约定。真正重要的是团队对每个等级的使用范围保持一致不能把所有消息都打成ERROR。1.2 生成日志与刷新日志分开我原稿把日志分成两个动作这个思路是对的组装一条完整日志把日志刷新到终端或文件。第二步可以用策略模式实现。日志对象只负责形成字符串具体输出位置由ConsoleLogStrategy或FileLogStrategy决定。这样切换输出位置时不需要修改每个LOG(...)调用。1.3 为什么日志输出也要加锁多个工作线程可能同时写日志。如果每条日志由多次operator拼接线程间输出可能交叉最终一行同时包含两条消息。因此刷新策略要保护真正的输出动作。锁的粒度应是一条完整日志而不是每写一个字符就加一次锁。对于文件输出目录只需要初始化一次后续以追加方式打开文件。1.4 时间函数的线程安全localtime可能返回指向共享静态对象的指针多个线程同时调用时会互相覆盖。Linux/POSIX 环境可以使用localtime_r让调用者提供struct tm存储空间。inlinestd::stringCurrentTime(){std::time_t nowstd::time(nullptr);std::tm local_time{};localtime_r(now,local_time);charbuffer[32];std::snprintf(buffer,sizeof(buffer),%04d-%02d-%02d %02d:%02d:%02d,local_time.tm_year1900,local_time.tm_mon1,local_time.tm_mday,local_time.tm_hour,local_time.tm_min,local_time.tm_sec);returnbuffer;}二、实现一个简单日志器2.1 输出策略#pragmaonce#includefilesystem#includefstream#includeiostream#includememory#includemutex#includestringclassLogStrategy{public:virtual~LogStrategy()default;virtualvoidWrite(conststd::stringmessage)0;};classConsoleLogStrategyfinal:publicLogStrategy{public:voidWrite(conststd::stringmessage)override{std::lock_guardstd::mutexlock(_mutex);std::coutmessage\n;}private:std::mutex _mutex;};classFileLogStrategyfinal:publicLogStrategy{public:explicitFileLogStrategy(std::string file./log/thread_pool.log):_file(std::move(file)){std::filesystem::pathpath(_file);if(path.has_parent_path())std::filesystem::create_directories(path.parent_path());}voidWrite(conststd::stringmessage)override{std::lock_guardstd::mutexlock(_mutex);std::ofstreamout(_file,std::ios::app);if(out)outmessage\n;}private:std::string _file;std::mutex _mutex;};基类析构函数必须是虚函数否则以后通过基类指针销毁派生策略时会出问题。策略对象内部各自持有一把锁保证同一种输出目标上的一条日志不会被拆开。2.2Logger与临时消息对象#pragmaonce#includecstdio#includectime#includememory#includemutex#includesstream#includestring#includesys/syscall.h#includeunistd.henumclassLogLevel{Debug,Info,Warning,Error,Fatal};inlineconstchar*ToString(LogLevel level){switch(level){caseLogLevel::Debug:returnDEBUG;caseLogLevel::Info:returnINFO;caseLogLevel::Warning:returnWARNING;caseLogLevel::Error:returnERROR;caseLogLevel::Fatal:returnFATAL;}returnUNKNOWN;}classLogger{public:Logger():_strategy(std::make_sharedConsoleLogStrategy()){}voidUseConsole(){std::lock_guardstd::mutexlock(_strategy_mutex);_strategystd::make_sharedConsoleLogStrategy();}voidUseFile(conststd::stringfile){autonextstd::make_sharedFileLogStrategy(file);std::lock_guardstd::mutexlock(_strategy_mutex);_strategystd::move(next);}classMessage{public:Message(Loggerlogger,LogLevel level,constchar*file,intline):_logger(logger){_stream[CurrentTime()][ToString(level)][pid:getpid()][tid:syscall(SYS_gettid)][file:line] ;}templateclassTMessageoperator(constTvalue){_streamvalue;return*this;}~Message(){_logger.Flush(_stream.str());}private:Logger_logger;std::ostringstream _stream;};Messageoperator()(LogLevel level,constchar*file,intline){returnMessage(*this,level,file,line);}private:voidFlush(conststd::stringmessage){std::shared_ptrLogStrategystrategy;{std::lock_guardstd::mutexlock(_strategy_mutex);strategy_strategy;}strategy-Write(message);}private:std::mutex _strategy_mutex;std::shared_ptrLogStrategy_strategy;};inlineLogger logger;#defineLOG(level)logger(level,__FILE__,__LINE__)使用临时Message对象的好处是可以保留熟悉的流式写法LOG(LogLevel::Info)task task_id finished;整条语句结束时临时对象析构已经拼好的字符串一次性交给刷新策略。这里没有让多个线程共同修改同一个字符串流每个线程在栈上创建自己的Message。策略切换也要同步。代码先把当前shared_ptr复制到局部变量再释放策略锁并真正输出避免文件 I/O 长时间占着_strategy_mutex。2.3 简单测试intmain(){LOG(LogLevel::Debug)console message;logger.UseFile(./log/thread_pool.log);LOG(LogLevel::Info)file message;LOG(LogLevel::Warning)queue is almost full;}这个实现没有日志轮转、异步写盘、过滤器和批量刷新因此不应替代spdlog、glog等成熟库。自己写一遍的目的是为后面的线程池准备输出工具并把多线程日志中最基本的锁范围想明白。三、为什么需要线程池3.1 来一个任务就创建一个线程的问题如果每收到一个短任务就pthread_create处理完再pthread_join或分离线程线程创建、栈空间准备、内核调度对象管理和销毁成本会反复发生。请求突然增多时线程数量也可能跟着失控。线程池的做法是提前创建固定数量的工作线程。任务到来后只放进队列空闲线程从队列取任务执行不再为每个任务重新创建线程。我原稿用餐厅预制菜类比池化资源提前准备好需要时直接复用。这个类比抓住了核心但线程池不是“任务提前算好”而是执行任务的线程提前创建好。3.2 线程池仍然是生产者消费者模型提交任务的线程是生产者任务队列是交易场所工作线程是消费者外部线程把任务放入队列条件变量唤醒一个工作线程工作线程取走任务在队列锁外执行任务再回到队列等待下一个任务。线程池只限制工作线程数量。若任务提交速度长期大于处理速度无界队列仍会不断增长。真正用于服务端时要考虑有界队列、拒绝策略或上游限流。3.3 固定线程数怎么选没有一个线程数适合所有程序。CPU 密集任务通常从硬件并发数附近开始测试任务经常等待磁盘或网络时可以适当多一些线程。但线程越多不等于越快过多线程会增加切换、栈内存和缓存失效成本。最可靠的办法仍然是结合任务类型、机器配置和实际负载压测而不是抄一个固定公式。四、固定线程池完整实现4.1 线程池需要哪些状态一个最小可用线程池需要工作线程数组任务队列保护队列和状态的互斥锁队列为空时供工作线程等待的条件变量是否停止接收任务的标记。退出时采用“排空队列再结束”的策略调用Stop后不再接收新任务已经进入队列的任务继续执行队列为空以后工作线程退出循环。4.2ThreadPool.hpp#pragmaonce#includecondition_variable#includecstddef#includefunctional#includemutex#includequeue#includestdexcept#includethread#includeutility#includevectorclassThreadPool{public:usingTaskstd::functionvoid();explicitThreadPool(std::size_t thread_count){if(thread_count0)throwstd::invalid_argument(thread_count must be positive);_workers.reserve(thread_count);for(std::size_t i0;ithread_count;i)_workers.emplace_back(ThreadPool::WorkerLoop,this,i);}~ThreadPool(){Stop();}ThreadPool(constThreadPool)delete;ThreadPooloperator(constThreadPool)delete;voidSubmit(Task task){{std::lock_guardstd::mutexlock(_mutex);if(_stopping)throwstd::runtime_error(submit on stopped thread pool);_tasks.push(std::move(task));}_condition.notify_one();}voidStop(){{std::lock_guardstd::mutexlock(_mutex);if(_stopping)return;_stoppingtrue;}_condition.notify_all();for(std::threadworker:_workers){if(worker.joinable())worker.join();}}private:voidWorkerLoop(std::size_t worker_id){LOG(LogLevel::Info)worker worker_id started;while(true){Task task;{std::unique_lockstd::mutexlock(_mutex);_condition.wait(lock,[this]{return_stopping||!_tasks.empty();});if(_stopping_tasks.empty())break;taskstd::move(_tasks.front());_tasks.pop();}try{task();}catch(conststd::exceptionerror){LOG(LogLevel::Error)worker worker_id task failed: error.what();}catch(...){LOG(LogLevel::Error)worker worker_id task failed: unknown exception;}}LOG(LogLevel::Info)worker worker_id stopped;}private:std::vectorstd::thread_workers;std::queueTask_tasks;std::mutex _mutex;std::condition_variable _condition;bool_stopping{false};};4.3 为什么使用unique_lockstd::condition_variable::wait需要在等待时释放锁醒来后再加锁因此它接收std::unique_lockstd::mutex。普通lock_guard不提供这种可解锁、再加锁的控制接口。谓词写成_stopping || !_tasks.empty()表示有任务或者线程池开始停止时工作线程都应该醒来继续判断。4.4 为什么任务必须在锁外执行工作线程只在锁内完成三件事检查状态、取队头任务、删除队头任务。拿到局部task后立刻离开作用域释放队列锁然后再执行任务。如果把task()放在锁内某个任务执行 3 秒其他工作线程就会 3 秒都拿不到队列锁。虽然创建了多个线程任务仍然接近串行执行。队列锁保护的是线程池内部状态不应该保护任务自己的业务过程。任务若访问其他共享数据应使用业务自己的同步手段。4.5 为什么捕获任务异常若异常一直逃出工作线程入口程序通常会调用std::terminate。一个任务抛错不应该直接带走整个线程池所以这里在工作循环中捕获异常并写日志。但不能空着catch (...)什么也不做。至少要记录任务失败否则线程池表面还在运行错误却完全丢失。五、线程池怎样安全退出5.1 工作线程的三种状态线程池运行期间工作线程大致处于三种状态队列为空在条件变量下等待拿到队列锁正在取任务已经释放队列锁正在执行任务。退出不能直接把这些线程“一股脑取消”。异步取消可能让线程停在持锁、分配内存或修改业务状态的中间位置资源很难收拾。5.2Stop的执行顺序本文采用下面的退出顺序在锁内把_stopping设为truenotify_all唤醒所有仍在等待的工作线程工作线程继续取完队列中的旧任务观察到“停止且队列为空”后退出循环Stop对每个工作线程执行join。只修改标记却不广播空闲工作线程可能永远睡在条件变量上析构函数会卡在join。只广播却不修改谓词线程醒来后又会发现队列为空重新等待。5.3 排空退出和立即退出线程池通常要明确两种语义退出方式已排队任务适合场景排空退出全部执行完正常停机、希望不丢任务立即退出放弃尚未开始的任务故障止损、任务已失去意义本文代码是排空退出。若要立即退出工作线程的判断和任务队列清理都要一起修改不能只把某个条件改掉。5.4 不要让工作线程销毁自己的线程池如果某个任务直接析构它所在的线程池Stop可能尝试join当前线程自己形成死锁或错误。线程池的所有权应由更外层对象管理停止动作也应在工作线程之外发起。六、运行线程池6.1 测试代码#includechrono#includeiostream#includethread#includeLog.hpp#includeThreadPool.hppintmain(){logger.UseConsole();ThreadPoolpool(3);for(inttask_id0;task_id8;task_id){pool.Submit([task_id]{LOG(LogLevel::Info)task task_id begin;std::this_thread::sleep_for(std::chrono::milliseconds(120));LOG(LogLevel::Info)task task_id end;});}pool.Stop();std::coutall tasks finished\n;}g-stdc17-O2-pthreadmain.cc-othread_pool ./thread_pool运行这段程序时不同任务的开始顺序和结束顺序可能变化这是正常的调度结果。观察重点不是固定输出顺序而是每条日志是否完整、8 个任务是否都执行结束以及三个工作线程能否在Stop后退出。6.2 可以继续扩展哪些功能这个版本跑通以后可以逐步增加有界任务队列与提交超时返回future的任务提交接口任务优先级运行中任务数、排队数和拒绝数统计日志文件轮转与异步写入明确区分排空退出和立即退出。扩展时不要一次把所有功能塞进去。先固定线程池最基本的状态机再加监控和策略否则退出路径很容易失控。总结日志模块负责把并发程序中的关键状态留下来线程池负责复用一组工作线程处理任务。线程池本质上仍是生产者消费者模型提交方生产任务工作线程消费任务队列和条件变量协调二者速度。实现时最容易写错的地方不是pthread_create而是锁的范围和退出过程队列状态在锁内修改任务在锁外执行停止时先修改谓词再唤醒全部线程最后逐个join。下一篇收尾线程安全专题继续讨论单例模式、可重入函数、死锁和 STL/智能指针的线程安全边界。参考资料【Linux】信号量到底在数什么从PV操作到RingQueue环形队列生产者消费者模型【Linux】多线程抢票为什么会出错从 mutex、futex 到 RAII讲清线程互斥【Linux】epoll 真的是 O(1) 吗从阻塞 I/O、select/poll 到 LT/ET 与 Reactor 一次讲透
返回列表