多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

第四篇 线程池与任务队列

第四篇 线程池与任务队列 原项目qinguoyi/TinyWebServer复刻仓库L2501031968/ccTinyWebServer完整 20 章教程仓库内docs/TinyWebServer-Recreation.md第 7 章 线程池与任务队列7.0 本章操作顺序本章按照以下顺序修改避免中途出现“引用了还不存在的文件”之类的编译错误新建threadpool/threadpool.h修改http/http_conn.cpp中的addfd()和modfd()在webserver.cpp顶部增加监听套接字的注册函数修改webserver.h给WebServer增加线程池成员修改webserver.cpp的构造函数和析构函数修改webserver.cpp的dealwithread()和dealwithwrite()修改makefile编译并验证单连接和并发连接本章结束后事件循环和数据处理的职责会变成主线程 等待 epoll 事件 接收新连接 把读写任务放入队列 工作线程 从队列取出连接 执行 read_once - process 或者执行 write7.1 本章目标目前每个请求都在epoll事件循环所在线程中串行处理。只要某个请求处理得慢其他连接就必须等待。本章要完成创建一个固定大小的线程池使用std::list保存待处理任务使用互斥锁保护任务队列使用信号量通知工作线程使用EPOLLONESHOT防止同一个连接被重复加入队列让主线程负责事件监听让工作线程负责连接上的读写学会安全停止并回收线程本章暂时不实现动态调整线程数量连接超时和定时器写日志配置文件数据库连接池7.2 为什么不能只增加std::thread如果每来一个连接就创建一个线程连接数量增加时会产生大量线程创建和销毁开销。线程池则预先创建固定数量的工作线程并让它们循环等待任务任务进入队列 - 信号量计数加一 - 某个阻塞在信号量上的工作线程被唤醒 - 工作线程从队列取出任务 - 执行任务 - 继续等待下一个任务这里有两个共享资源m_workqueue 由多个线程共同访问需要 locker m_queuestat 负责阻塞和唤醒线程使用 sem7.3 本章采用的线程模型本章采用 Reactor 风格的简化版本epoll 报告 EPOLLIN - 主线程把“读任务”放进队列 - 工作线程执行 read_once() 和 process() epoll 报告 EPOLLOUT - 主线程把“写任务”放进队列 - 工作线程执行 write()任务使用state区分state 0读取并解析请求 state 1发送响应http_conn已经有公开成员m_state因此线程池可以直接设置它。7.4 新建线程池头文件在项目根目录执行mkdir-pthreadpooltouchthreadpool/threadpool.h然后把下面的完整内容写入threadpool/threadpool.h。这是一个新文件不需要在现有文件里寻找插入位置。#ifndefTHREADPOOL_H#defineTHREADPOOL_H#includeexception#includelist#includepthread.h#include../lock/locker.htemplatetypenameTclassthreadpool{public:explicitthreadpool(intthread_number8,intmax_requests10000);~threadpool();boolappend(T*request,intstate);private:staticvoid*worker(void*arg);voidrun();intm_thread_number;intm_max_requests;pthread_t*m_threads;std::listT*m_workqueue;locker m_queuelocker;sem m_queuestat;boolm_stop;};templatetypenameTthreadpoolT::threadpool(intthread_number,intmax_requests):m_thread_number(thread_number),m_max_requests(max_requests),m_threads(nullptr),m_stop(false){if(m_thread_number0||m_max_requests0)throwstd::exception();m_threadsnewpthread_t[m_thread_number];for(inti0;im_thread_number;i){if(pthread_create(m_threads[i],nullptr,worker,this)!0){m_stoptrue;for(intj0;ji;j)m_queuestat.post();for(intj0;ji;j)pthread_join(m_threads[j],nullptr);delete[]m_threads;m_threadsnullptr;throwstd::exception();}}}templatetypenameTthreadpoolT::~threadpool(){m_queuelocker.lock();m_stoptrue;m_queuelocker.unlock();for(inti0;im_thread_number;i)m_queuestat.post();for(inti0;im_thread_number;i)pthread_join(m_threads[i],nullptr);delete[]m_threads;m_threadsnullptr;}templatetypenameTboolthreadpoolT::append(T*request,intstate){m_queuelocker.lock();if(m_stop||static_castint(m_workqueue.size())m_max_requests){m_queuelocker.unlock();returnfalse;}request-m_statestate;m_workqueue.push_back(request);m_queuelocker.unlock();m_queuestat.post();returntrue;}templatetypenameTvoid*threadpoolT::worker(void*arg){threadpool*poolstatic_castthreadpool*(arg);pool-run();returnnullptr;}templatetypenameTvoidthreadpoolT::run(){while(true){m_queuestat.wait();m_queuelocker.lock();if(m_stopm_workqueue.empty()){m_queuelocker.unlock();break;}if(m_workqueue.empty()){m_queuelocker.unlock();continue;}T*requestm_workqueue.front();m_workqueue.pop_front();m_queuelocker.unlock();if(!request)continue;if(request-m_state0){if(request-read_once())request-process();elserequest-close_conn();}else{if(!request-write())request-close_conn();}}}#endif7.5 理解线程池的关键代码7.5.1 构造函数创建线程pthread_create()的第四个参数传入thispthread_create(m_threads[i],nullptr,worker,this);worker()是静态函数因此它没有this指针。通过arg参数把对象地址传入threadpool*poolstatic_castthreadpool*(arg);pool-run();7.5.2 队列必须加锁主线程会调用append()添加任务工作线程会在run()中删除任务。两个线程不能同时修改std::list。因此所有对m_workqueue的访问都放在m_queuelocker.lock();// 访问 m_workqueuem_queuelocker.unlock();7.5.3 信号量负责等待和唤醒工作线程启动后调用m_queuestat.wait();队列为空时信号量计数为0线程在这里阻塞。主线程调用append()成功后执行m_queuestat.post();信号量计数增加一个等待中的工作线程被唤醒。7.5.4 为什么增加m_stop如果线程一直处于分离状态并在while (true)中运行threadpool析构时无法确认工作线程已经退出。本章没有调用pthread_detach()而是在析构函数中设置 m_stop - 给每个线程发送一次信号 - pthread_join 等待所有线程退出 - 删除线程数组这比直接分离线程更容易保证析构顺序正确。7.6 给连接启用 EPOLLONESHOT打开http/http_conn.cpp找到最前面的addfd()。修改前voidaddfd(intepollfd,intfd){epoll_event event{};event.data.fdfd;event.eventsEPOLLIN|EPOLLRDHUP;只替换事件标志这一行event.eventsEPOLLIN|EPOLLRDHUP|EPOLLONESHOT;修改后的完整函数是voidaddfd(intepollfd,intfd){epoll_event event{};event.data.fdfd;event.eventsEPOLLIN|EPOLLRDHUP|EPOLLONESHOT;if(epoll_ctl(epollfd,EPOLL_CTL_ADD,fd,event)-1){perror(epoll_ctl add);return;}setnonblocking(fd);}EPOLLONESHOT表示一个文件描述符产生一次事件后内核会暂时把它从监控集合中停用。工作线程处理完成后必须通过modfd()再重新注册。这样做是为了防止下面的情况主线程发现 EPOLLIN - 把连接加入队列 连接还没有被工作线程读取 - 主线程再次发现 EPOLLIN - 同一个连接又被加入队列如果同一个连接被两个线程同时处理会破坏缓冲区并导致响应错乱。接着找到modfd()把事件标志改成voidmodfd(intepollfd,intfd,intev){epoll_event event{};event.data.fdfd;event.eventsev|EPOLLRDHUP|EPOLLONESHOT;if(epoll_ctl(epollfd,EPOLL_CTL_MOD,fd,event)-1)perror(epoll_ctl mod);}这样工作线程重新注册连接时仍然保留“一次触发”的语义。7.7 监听套接字不能使用 EPOLLONESHOT这里有一个非常容易踩到的坑。WebServer::eventListen()当前也调用addfd(m_epollfd,m_listenfd);修改后的addfd()会给所有文件描述符加上EPOLLONESHOT。如果监听套接字也只触发一次那么第一批连接被接收后监听套接字就不会再次上报事件后面的连接会一直等待最终表现为curl超时。监听套接字必须持续接收新连接因此不要再对它调用带EPOLLONESHOT的addfd()。打开webserver.cpp在WebServer::WebServer()之前插入一个匿名命名空间函数namespace{voidadd_listenfd(intepollfd,intfd){epoll_event event{};event.data.fdfd;event.eventsEPOLLIN;if(epoll_ctl(epollfd,EPOLL_CTL_ADD,fd,event)-1){perror(epoll_ctl add listenfd);return;}setnonblocking(fd);}}插入位置如下webserver.cpp include 区域 在这里插入 namespace { ... } WebServer::WebServer()然后在WebServer::eventListen()中找到addfd(m_epollfd,m_listenfd);替换成add_listenfd(m_epollfd,m_listenfd);最终事件注册规则是监听套接字EPOLLIN不加 EPOLLONESHOT 连接套接字EPOLLIN | EPOLLRDHUP | EPOLLONESHOT7.8 修改 WebServer 头文件打开webserver.h。7.8.1 增加头文件找到#includehttp/http_conn.h在它下面增加#includethreadpool/threadpool.h修改后#includehttp/http_conn.h#includethreadpool/threadpool.h#includesys/epoll.h#includestring7.8.2 增加线程池成员找到私有成员http_conn*m_users;std::string m_root;在m_root下面增加threadpoolhttp_conn*m_thread_pool;修改后http_conn*m_users;std::string m_root;threadpoolhttp_conn*m_thread_pool;这里不需要新增任何成员函数后续直接修改已有的dealwithread()和dealwithwrite()。7.9 在构造函数中创建线程池打开webserver.cpp找到WebServer::WebServer()。当前构造函数已经在最后计算m_rootstd::string(server_path)/root;在它下面增加m_thread_poolnewthreadpoolhttp_conn(8,10000);修改后的尾部是m_rootstd::string(server_path)/root;m_thread_poolnewthreadpoolhttp_conn(8,10000);这里的两个参数分别是8 工作线程数量 10000 任务队列允许保存的最大请求数本章先使用固定参数。后面的配置章节再把这些值改由配置文件传入。7.10 在析构函数中停止线程池找到WebServer::~WebServer()。在线程池对象还存在时工作线程可能正在运行。因此应该先停止并回收线程池再关闭描述符和删除连接数组。修改后的析构函数是WebServer::~WebServer(){deletem_thread_pool;if(m_epollfd!-1)close(m_epollfd);if(m_listenfd!-1)close(m_listenfd);delete[]m_users;}特别注意删除顺序先 delete m_thread_pool 再 delete[] m_users如果先删除m_users仍在运行的工作线程可能访问已经释放的http_conn对象。7.11 修改读事件分发找到WebServer::dealwithread()。修改前voidWebServer::dealwithread(intsockfd){if(m_users[sockfd].read_once())m_users[sockfd].process();elseclose_conn(sockfd);}整个函数替换成voidWebServer::dealwithread(intsockfd){if(!m_thread_pool-append(m_users[sockfd],0))close_conn(sockfd);}这里传入m_users[sockfd]这是连接对象在数组中的地址。工作线程稍后会通过这个指针调用read_once()process()如果队列已满或者线程池正在停止append()返回false主线程直接关闭这个连接避免任务丢失后连接一直挂着。7.12 修改写事件分发找到WebServer::dealwithwrite()。修改前voidWebServer::dealwithwrite(intsockfd){if(!m_users[sockfd].write())close_conn(sockfd);}整个函数替换成voidWebServer::dealwithwrite(intsockfd){if(!m_thread_pool-append(m_users[sockfd],1))close_conn(sockfd);}state 1告诉工作线程这次执行写任务而不是读取和解析任务。整个流程现在是主线程 EPOLLIN - append(request, 0) EPOLLOUT - append(request, 1) 工作线程 state 0 - read_once - process state 1 - write7.13 修改 Makefile打开makefile。7.13.1 增加头文件依赖找到server: main.cpp webserver.cpp webserver.h \ http/http_conn.cpp http/http_conn.h替换成server: main.cpp webserver.cpp webserver.h \ http/http_conn.cpp http/http_conn.h \ threadpool/threadpool.h lock/locker.h这样修改threadpool.h或locker.h后make会重新编译服务器。7.13.2 增加 pthread 链接参数找到编译命令$(CXX) -o server main.cpp webserver.cpp \ http/http_conn.cpp $(CXXFLAGS)替换成$(CXX) -o server main.cpp webserver.cpp \ http/http_conn.cpp $(CXXFLAGS) -pthread完整的目标可以写成server: main.cpp webserver.cpp webserver.h \ http/http_conn.cpp http/http_conn.h \ threadpool/threadpool.h lock/locker.h $(CXX) -o server main.cpp webserver.cpp \ http/http_conn.cpp $(CXXFLAGS) -pthread7.14 编译检查点执行cd/home/cc/ccTinyWebServermakeclean server本次验证的编译输出rm -f test_locker server g -o server main.cpp webserver.cpp \ http/http_conn.cpp -g -Wall -Wextra -stdc11 -pthread编译器没有错误也没有警告。如果出现undefined reference to pthread_create说明编译命令缺少-pthread如果出现threadpool/threadpool.h: No such file or directory检查是否在项目根目录创建了threadpool文件夹文件名必须是threadpool/threadpool.h7.15 运行服务器执行./server输出server is listening on port 90067.16 验证普通静态文件请求打开另一个终端执行curl-i--max-time5http://127.0.0.1:9006/本次返回HTTP/1.1 200 OK Content-Length: 209 Content-Type: text/html; charsetutf-8 Connection: close !DOCTYPE html html head meta charsetUTF-8 titleTinyWebServer/title /head body h1Static File Service/h1 pThis page is served from the root directory./p /body /html这说明请求已经被工作线程读取、解析并生成了响应。7.17 验证 404 请求执行curl-i--max-time5http://127.0.0.1:9006/missing.html返回HTTP/1.1 404 Not Found Content-Length: 49 Content-Type: text/html; charsetutf-8 Connection: close The requested file was not found on this server.这说明EPOLLONESHOT重新注册和process_write()仍然正常工作。7.18 验证十个并发请求执行for i in {1..10}; do curl -sS --max-time 5 \ -w %{size_download}\n \ -o /dev/null \ http://127.0.0.1:9006/ done本次输出209 209 209 209 209 209 209 209 209 209十个请求都由线程池处理并返回了完整的 209 字节文件。如果这里只有第一批请求成功后面的请求超时首先检查监听套接字是否被错误地加上了EPOLLONESHOT。eventListen()必须调用不带一次触发语义的add_listenfd()。7.19 本章最容易出现的几个问题7.19.1 主线程和工作线程同时读同一个连接错误结构if(m_users[sockfd].read_once())m_thread_pool-append(m_users[sockfd],0);这样主线程先读取了一部分数据工作线程随后又读取同一个连接解析状态容易混乱。本章采用的分工是主线程只把任务放进队列 工作线程执行 read_once因此dealwithread()中不要再调用read_once()。7.19.2 忘记重新注册 EPOLLONESHOT工作线程执行完成后process()或write()会调用modfd()modfd(m_epollfd,m_sockfd,EPOLLIN);或者modfd(m_epollfd,m_sockfd,EPOLLOUT);如果modfd()没有保留EPOLLONESHOT一次触发语义就会丢失。如果完全没有调用modfd()连接会一直处于停用状态。7.19.3 队列满以后直接丢弃连接本章在append()返回false时直接关闭连接。这是一种最简单的保护方式可以避免连接永远得不到处理。更完整的服务器可以在这里发送503 Service Unavailable这个功能留到后面的 HTTP 错误处理章节扩展。7.20 当前实现边界本章完成后服务器已经具备多线程处理能力但仍然有以下边界线程数量和队列长度写死为8和10000一个连接使用一次EPOLLONESHOT同一时刻只会有一个线程处理它http_conn::m_user_count还不是原子变量没有连接超时和空闲连接清理没有优雅处理SIGINT和SIGTERM队列满时直接关闭连接没有返回 503没有日志系统输出仍使用std::cout7.21 本章小结本章新增了一条完整的线程池链路epoll 事件 - WebServer::dealwithread / dealwithwrite - threadpool::append - m_workqueue - m_queuestat.post - 工作线程 m_queuestat.wait - 取出 http_conn - read_once / process / write同时使用EPOLLONESHOT保证一个连接不会被多个工作线程重复处理。下一章将实现定时器定期清理长时间没有活动的连接避免无效文件描述符一直占用服务器资源。、
返回列表