微信程序员压测20种并发模型,性能最强的竟是?
👉目录
1 缘起
2 前置说明
3 预备工作
4 BenchMark 工具
5 20种不同的并发模型
6 参考
01
02
.├── BenchMark│ ├── benchmark.cpp│ ├── client.hpp│ ├── clientmanager.hpp│ ├── makefile│ ├── percentile.hpp│ ├── stat.hpp│ └── timer.hpp├── common│ ├── cmdline.cpp│ ├── cmdline.h│ ├── codec.hpp│ ├── conn.hpp│ ├── coroutine.cpp│ ├── coroutine.h│ ├── epollctl.hpp│ ├── packet.hpp│ └── utils.hpp├── ConcurrencyModel│ ├── Epoll│ ├── EpollReactorProcessPoolCoroutine│ ├── EpollReactorProcessPoolMS│ ├── EpollReactorSingleProcess│ ├── EpollReactorSingleProcessCoroutine│ ├── EpollReactorSingleProcessET│ ├── EpollReactorThreadPool│ ├── EpollReactorThreadPoolHSHA│ ├── EpollReactorThreadPoolMS│ ├── LeaderAndFollower│ ├── MultiProcess│ ├── MultiThread│ ├── Poll│ ├── PollReactorSingleProcess│ ├── ProcessPool1│ ├── ProcessPool2│ ├── Select│ ├── SelectReactorSingleProcess│ ├── SingleProcess│ └── ThreadPool├── readme.md└── test├── codectest.cpp├── coroutinetest.cpp├── makefile├── packettest.cpp├── unittestcore.hpp└── unittestentry.cpp
BenchMark 是基准性能压测工具的代码目录。 ConcurrencyModel 是20种不同并发模型的代码目录,这个目录下有 20 个不同的子目录,每个子目录都代表着一种并发模型的实现示例。 common 是公共代码的目录。 test 目录为单元测试代码的目录。
03
template <typename Function, typename... Args>int CoroutineCreate(Schedule& schedule, Function&& f, Args&&... args) {int id = 0;for (id = 0; id < schedule.coroutineCnt; id++) {if (schedule.coroutines[id]->state == Idle) break;}if (id >= schedule.coroutineCnt) {return kInvalidRoutineId;}Coroutine* routine = schedule.coroutines[id];std::function<void()> entry = std::bind(std::forward<Function>(f), std::forward<Args>(args)...);CoroutineInit(schedule, routine, entry);return id;}
04
root@centos BenchMark $ ./BenchMark -hBenchMark -ip 0.0.0.0 -port 1688 -thread_count 1 -max_req_count 100000 -pkt_size 1024 -client_count 200 -run_time 60 -rate_limit 10000 -debugoptions:-h,--help print usage-ip,--ip service listen ip-port,--port service listen port-thread_count,--thread_count run thread count-max_req_count,--max_req_count one connection max req count-pkt_size,--pkt_size size of send packet, unit is byte-client_count,--client_count count of client-run_time,--run_time run time, unit is second-rate_limit,--rate_limit rate limit, unit is qps/second-debug,--debug debug mode, more info print
支持对监听在指定的 ip 和 port 的服务发起压测。 支持多线程压测,并可以指定使用的线程数。 支持指定客户端连接建立成功之后,最多可以发起多少次请求。(ps:这个选项值如果设置为1,则请求就退化成通过短连接来完成) 支持指定请求包的大小,单位为字节。 支持指定每个线程下发起请求的客户端并发连接数。 支持指定总的压测时间,单位秒。 支持指定压测能产生的最大的流量负载,单位 qps。 支持 debug 模式。
接口的 pct50、pct95、pct99 和 pct999 的耗时数据。 请求成功数、请求失败数、尝试建立连接数、连接失败数、读失败数和写失败数。 客户端连接数和请求成功的 qps 数。 请求失败率和连接失败率。
05
#include <sys/socket.h>#include <unistd.h>#include <iostream>#include "../../common/cmdline.h"#include "../../common/utils.hpp"using namespace std;using namespace MyEcho;void handlerClient(int client_fd) {string msg;while (true) {if (not RecvMsg(client_fd, msg)) {return;}if (not SendMsg(client_fd, msg)) {return;}}}void usage() {cout << "SingleProcess -ip 0.0.0.0 -port 1688" << endl;cout << "options:" << endl;cout << " -h,--help print usage" << endl;cout << " -ip,--ip listen ip" << endl;cout << " -port,--port listen port" << endl;cout << endl;}int main(int argc, char* argv[]) {string ip;int64_t port;CmdLine::StrOptRequired(&ip, "ip");CmdLine::Int64OptRequired(&port, "port");CmdLine::SetUsage(usage);CmdLine::Parse(argc, argv);int sock_fd = CreateListenSocket(ip, port, false);if (sock_fd < 0) {return -1;}while (true) {int client_fd = accept(sock_fd, NULL, 0);if (client_fd < 0) {perror("accept failed");continue;}handlerClient(client_fd);close(client_fd);}return 0;}
#include <signal.h>#include <sys/socket.h>#include <unistd.h>#include <iostream>#include "../../common/cmdline.h"#include "../../common/utils.hpp"using namespace std;using namespace MyEcho;void handlerClient(int client_fd) {string msg;while (true) {if (not RecvMsg(client_fd, msg)) {return;}if (not SendMsg(client_fd, msg)) {return;}}}void childExitSignalHandler() {struct sigaction act;act.sa_handler = SIG_IGN; //设置信号处理函数,这里忽略子进程的退出信号sigemptyset(&act.sa_mask); //信号屏蔽设置为空act.sa_flags = 0; //标志位设置为0sigaction(SIGCHLD, &act, NULL);}void usage() {cout << "MultiProcess -ip 0.0.0.0 -port 1688" << endl;cout << "options:" << endl;cout << " -h,--help print usage" << endl;cout << " -ip,--ip listen ip" << endl;cout << " -port,--port listen port" << endl;cout << endl;}int main(int argc, char* argv[]) {string ip;int64_t port;CmdLine::StrOptRequired(&ip, "ip");CmdLine::Int64OptRequired(&port, "port");CmdLine::SetUsage(usage);CmdLine::Parse(argc, argv);int sock_fd = CreateListenSocket(ip, port, false);if (sock_fd < 0) {return -1;}childExitSignalHandler(); // 这里需要忽略子进程退出信号,否则会导致大量的僵尸进程,服务后续无法再创建子进程while (true) {int client_fd = accept(sock_fd, NULL, 0);if (client_fd < 0) {perror("accept failed");continue;}pid_t pid = fork();if (pid == -1) {close(client_fd);perror("fork failed");continue;}if (pid == 0) { // 子进程handlerClient(client_fd);close(client_fd);exit(0); // 处理完请求,子进程直接退出} else {close(client_fd); // 父进程直接关闭客户端连接,否则文件描述符会泄露}}return 0;}
#include <sys/socket.h>#include <unistd.h>#include <iostream>#include <thread>#include "../../common/cmdline.h"#include "../../common/utils.hpp"using namespace std;using namespace MyEcho;void handlerClient(int client_fd) {string msg;while (true) {if (not RecvMsg(client_fd, msg)) {break;}if (not SendMsg(client_fd, msg)) {break;}}close(client_fd);}void usage() {cout << "MultiThread -ip 0.0.0.0 -port 1688" << endl;cout << "options:" << endl;cout << " -h,--help print usage" << endl;cout << " -ip,--ip listen ip" << endl;cout << " -port,--port listen port" << endl;cout << endl;}int main(int argc, char* argv[]) {string ip;int64_t port;CmdLine::StrOptRequired(&ip, "ip");CmdLine::Int64OptRequired(&port, "port");CmdLine::SetUsage(usage);CmdLine::Parse(argc, argv);int sock_fd = CreateListenSocket(ip, port, false);if (sock_fd < 0) {return -1;}while (true) {int client_fd = accept(sock_fd, NULL, 0);if (client_fd < 0) {perror("accept failed");continue;}std::thread(handlerClient, client_fd).detach(); // 这里需要调用detach,让创建的线程独立运行}return 0;}
#include <sys/socket.h>#include <unistd.h>#include <iostream>#include "../../common/cmdline.h"#include "../../common/utils.hpp"using namespace std;using namespace MyEcho;void handlerClient(int client_fd, int64_t& count) {string msg;while (true) {if (not RecvMsg(client_fd, msg)) {return;}if (not SendMsg(client_fd, msg)) {return;}count++;}}void handler(int worker_id, int sock_fd) {int64_t count = 0;while (true) {int client_fd = accept(sock_fd, NULL, 0);if (client_fd < 0) {perror("accept failed");continue;}handlerClient(client_fd, count);close(client_fd);count++;if (count >= 10000) {cout << "worker_id[" << worker_id << "] deal_1w_request" << endl;count = 0;}}}void usage() {cout << "ProcessPool1 -ip 0.0.0.0 -port 1688" << endl;cout << "options:" << endl;cout << " -h,--help print usage" << endl;cout << " -ip,--ip listen ip" << endl;cout << " -port,--port listen port" << endl;cout << endl;}int main(int argc, char* argv[]) {string ip;int64_t port;CmdLine::StrOptRequired(&ip, "ip");CmdLine::Int64OptRequired(&port, "port");CmdLine::SetUsage(usage);CmdLine::Parse(argc, argv);int sock_fd = CreateListenSocket(ip, port, false);if (sock_fd < 0) {return -1;}for (int i = 0; i < GetNProcs(); i++) {pid_t pid = fork();if (pid < 0) {perror("fork failed");continue;}if (0 == pid) {handler(i, sock_fd); // 子进程陷入死循环,处理客户端请求exit(0);}}while (true) sleep(1); // 父进程陷入死循环return 0;}
#include <sys/socket.h>#include <unistd.h>#include <iostream>#include "../../common/cmdline.h"#include "../../common/utils.hpp"using namespace std;using namespace MyEcho;void handlerClient(int client_fd, int64_t& count) {string msg;while (true) {if (not RecvMsg(client_fd, msg)) {return;}if (not SendMsg(client_fd, msg)) {return;}count++;}}void handler(int worker_id, string ip, int64_t port) {// 开启SO_REUSEPORT选项int sock_fd = CreateListenSocket(ip, port, true);if (sock_fd < 0) {return;}int64_t count = 0;while (true) {int client_fd = accept(sock_fd, NULL, 0);if (client_fd < 0) {perror("accept failed");continue;}handlerClient(client_fd, count);close(client_fd);if (count >= 10000) {cout << "worker_id[" << worker_id << "] deal_1w_request" << endl;count = 0;}}}void usage() {cout << "ProcessPool2 -ip 0.0.0.0 -port 1688" << endl;cout << "options:" << endl;cout << " -h,--help print usage" << endl;cout << " -ip,--ip listen ip" << endl;cout << " -port,--port listen port" << endl;cout << endl;}int main(int argc, char* argv[]) {string ip;int64_t port;CmdLine::StrOptRequired(&ip, "ip");CmdLine::Int64OptRequired(&port, "port");CmdLine::SetUsage(usage);CmdLine::Parse(argc, argv);for (int i = 0; i < GetNProcs(); i++) {pid_t pid = fork();if (pid < 0) {perror("fork failed");continue;}if (0 == pid) {handler(i, ip, port); // 子进程陷入死循环,处理客户端请求exit(0);}}while (true) sleep(1); // 父进程陷入死循环return 0;}
#include <arpa/inet.h>#include <netinet/in.h>#include <sys/socket.h>#include <unistd.h>#include <iostream>#include <thread>#include "../../common/cmdline.h"#include "../../common/utils.hpp"using namespace std;using namespace MyEcho;void handlerClient(int client_fd) {string msg;while (true) {if (not RecvMsg(client_fd, msg)) {return;}if (not SendMsg(client_fd, msg)) {return;}}}void handler(string ip, int64_t port) {// 开启SO_REUSEPORT选项int sock_fd = CreateListenSocket(ip, port, true);if (sock_fd < 0) {return;}while (true) {int client_fd = accept(sock_fd, NULL, 0);if (client_fd < 0) {perror("accept failed");continue;}handlerClient(client_fd);close(client_fd);}}void usage() {cout << "ThreadPool -ip 0.0.0.0 -port 1688" << endl;cout << "options:" << endl;cout << " -h,--help print usage" << endl;cout << " -ip,--ip listen ip" << endl;cout << " -port,--port listen port" << endl;cout << endl;}int main(int argc, char* argv[]) {string ip;int64_t port;CmdLine::StrOptRequired(&ip, "ip");CmdLine::Int64OptRequired(&port, "port");CmdLine::SetUsage(usage);CmdLine::Parse(argc, argv);for (int i = 0; i < GetNProcs(); i++) {std::thread(handler, ip, port).detach(); // 这里需要调用detach,让创建的线程独立运行}while (true) sleep(1); // 主线程陷入死循环return 0;}
#include <sys/socket.h>#include <unistd.h>#include <iostream>#include <mutex>#include <thread>#include "../../common/cmdline.h"#include "../../common/utils.hpp"using namespace std;using namespace MyEcho;std::mutex Mutex;void handlerClient(int client_fd) {string msg;while (true) {if (not RecvMsg(client_fd, msg)) {return;}if (not SendMsg(client_fd, msg)) {return;}}}void handler(int sock_fd) {while (true) {int client_fd;// follower等待获取锁,成为leader{std::lock_guard<std::mutex> guard(Mutex);client_fd = accept(sock_fd, NULL, 0); // 获取锁,并获取客户端的连接if (client_fd < 0) {perror("accept failed");continue;}}handlerClient(client_fd); // 处理每个客户端请求close(client_fd);}}void usage() {cout << "LeaderAndFollower -ip 0.0.0.0 -port 1688" << endl;cout << "options:" << endl;cout << " -h,--help print usage" << endl;cout << " -ip,--ip listen ip" << endl;cout << " -port,--port listen port" << endl;cout << endl;}int main(int argc, char* argv[]) {string ip;int64_t port;CmdLine::StrOptRequired(&ip, "ip");CmdLine::Int64OptRequired(&port, "port");CmdLine::SetUsage(usage);CmdLine::Parse(argc, argv);int sock_fd = CreateListenSocket(ip, port, false);if (sock_fd < 0) {return -1;}for (int i = 0; i < GetNProcs(); i++) {std::thread(handler, sock_fd).detach(); // 这里需要调用detach,让创建的线程独立运行}while (true) sleep(1); // 主进程陷入死循环return 0;}
我们可以看到单进程的并发模型性能是最差的。 多进程和多线程的并发模型,由于能使用到多个 CPU,所以性能有所提升,但是因为需要频繁的创建和销毁进程和线程,接口的 pct50 和 pct95 耗时较高。 进程池1和进程池2的并发模型,由于没有频繁的创建和销毁进程和线程的损耗,性能比多进程和多线程的并发模型高,接口的 pct50 和 pct95 耗时也更低。 进程池2并发模型接口的 pct999 耗时比进程池1并发模型的高出不少,这个是因为进程池2的并发模型是由操作系统来做负载均衡的,但这个策略并无法保证对流量负载做完美的均分,导致接口长尾的耗时较高,而进程池1的并发模型是多进程抢锁,每个进程的流量负载会更均衡,但因为有锁,所以进程池1的并发模型性能比进程池2的并发模型低一些。 线程池的并发模型和进程池2的并发模型,性能差异并不是很大,因为线程池的并发模型也是由操作系统来做负载均衡的,所以存在接口长尾的耗时较高的情况。 领导者/跟随者的并发模型和进程池1的并发模型很相似,这两个模型所有的指标都差异很小,领导者/跟随者的并发模型可以看到显式的使用锁,而进程池1的并发模型没有。
#include <stdio.h>#include <unistd.h>#include <iostream>#include <unordered_set>#include "../../common/cmdline.h"#include "../../common/utils.hpp"using namespace std;using namespace MyEcho;void updateReadSet(unordered_set<int> &read_fds, int &max_fd, int sock_fd, fd_set &read_set) {max_fd = sock_fd;FD_ZERO(&read_set);for (const auto &read_fd : read_fds) {if (read_fd > max_fd) {max_fd = read_fd;}FD_SET(read_fd, &read_set);}}void handlerClient(int client_fd) {string msg;while (true) {if (not RecvMsg(client_fd, msg)) {return;}if (not SendMsg(client_fd, msg)) {return;}}}void usage() {cout << "Select -ip 0.0.0.0 -port 1688" << endl;cout << "options:" << endl;cout << " -h,--help print usage" << endl;cout << " -ip,--ip listen ip" << endl;cout << " -port,--port listen port" << endl;cout << endl;}int main(int argc, char *argv[]) {string ip;int64_t port;CmdLine::StrOptRequired(&ip, "ip");CmdLine::Int64OptRequired(&port, "port");CmdLine::SetUsage(usage);CmdLine::Parse(argc, argv);int sock_fd = CreateListenSocket(ip, port, false);if (sock_fd < 0) {return -1;}int max_fd;fd_set read_set;SetNotBlock(sock_fd);unordered_set<int> read_fds;while (true) {read_fds.insert(sock_fd);updateReadSet(read_fds, max_fd, sock_fd, read_set);int ret = select(max_fd + 1, &read_set, NULL, NULL, NULL);if (ret <= 0) {if (ret < 0) perror("select failed");continue;}for (int i = 0; i <= max_fd; i++) {if (not FD_ISSET(i, &read_set)) {continue;}if (i == sock_fd) { // 监听的sock_fd可读,则表示有新的链接LoopAccept(sock_fd, 1024, [&read_fds](int client_fd) {if (client_fd >= FD_SETSIZE) { // 大于FD_SETSIZE的值,则不支持close(client_fd);return;}read_fds.insert(client_fd); // 新增到要监听的fd集合中});continue;}handlerClient(i);read_fds.erase(i);close(i);}}return 0;}
#include <arpa/inet.h>#include <netinet/in.h>#include <poll.h>#include <stdio.h>#include <unistd.h>#include <iostream>#include <unordered_set>#include "../../common/cmdline.h"#include "../../common/utils.hpp"using namespace std;using namespace MyEcho;void updateFds(std::unordered_set<int> &client_fds, pollfd **fds, int &nfds) {if (*fds != nullptr) {delete[](*fds);}nfds = client_fds.size();*fds = new pollfd[nfds];int index = 0;for (const auto &client_fd : client_fds) {(*fds)[index].fd = client_fd;(*fds)[index].events = POLLIN;(*fds)[index].revents = 0;index++;}}void handlerClient(int client_fd) {string msg;while (true) {if (not RecvMsg(client_fd, msg)) {return;}if (not SendMsg(client_fd, msg)) {return;}}}void usage() {cout << "Poll -ip 0.0.0.0 -port 1688" << endl;cout << "options:" << endl;cout << " -h,--help print usage" << endl;cout << " -ip,--ip listen ip" << endl;cout << " -port,--port listen port" << endl;cout << endl;}int main(int argc, char *argv[]) {string ip;int64_t port;CmdLine::StrOptRequired(&ip, "ip");CmdLine::Int64OptRequired(&port, "port");CmdLine::SetUsage(usage);CmdLine::Parse(argc, argv);int sock_fd = CreateListenSocket(ip, port, false);if (sock_fd < 0) {return -1;}int nfds = 0;pollfd *fds = nullptr;std::unordered_set<int> client_fds;client_fds.insert(sock_fd);SetNotBlock(sock_fd);while (true) {updateFds(client_fds, &fds, nfds);int ret = poll(fds, nfds, -1);if (ret <= 0) {if (ret < 0) perror("poll failed");continue;}for (int i = 0; i < nfds; i++) {if (not(fds[i].revents & POLLIN)) {continue;}int current_fd = fds[i].fd;if (current_fd == sock_fd) {LoopAccept(sock_fd, 1024, [&client_fds](int client_fd) {client_fds.insert(client_fd); // 新增到要监听的fd集合中});continue;}handlerClient(current_fd);client_fds.erase(current_fd);close(current_fd);}}return 0;}
#include <arpa/inet.h> #include <assert.h>#include <netinet/in.h>#include <stdio.h>#include <stdlib.h>#include <sys/epoll.h>#include <sys/socket.h>#include <unistd.h>#include <iostream>#include "../../common/cmdline.h"#include "../../common/epollctl.hpp"using namespace std;using namespace MyEcho;void handlerClient(int client_fd) {string msg;while (true) {if (not RecvMsg(client_fd, msg)) {return;}if (not SendMsg(client_fd, msg)) {return;}}}void usage() {cout << "Epoll -ip 0.0.0.0 -port 1688 -la" << endl;cout << "options:" << endl;cout << " -h,--help print usage" << endl;cout << " -ip,--ip listen ip" << endl;cout << " -port,--port listen port" << endl;cout << " -la,--la loop accept" << endl;cout << endl;}int main(int argc, char *argv[]) {string ip;int64_t port;bool is_loop_accept;CmdLine::StrOptRequired(&ip, "ip");CmdLine::Int64OptRequired(&port, "port");CmdLine::BoolOpt(&is_loop_accept, "la");CmdLine::SetUsage(usage);CmdLine::Parse(argc, argv);int sock_fd = CreateListenSocket(ip, port, false);if (sock_fd < 0) {return -1;}epoll_event events[2048];int epoll_fd = epoll_create(1);if (epoll_fd < 0) {perror("epoll_create failed");return -1;}cout << "loop_accept = " << is_loop_accept << endl;Conn conn(sock_fd, epoll_fd, false);SetNotBlock(sock_fd);AddReadEvent(&conn);while (true) {int num = epoll_wait(epoll_fd, events, 2048, -1);if (num < 0) {perror("epoll_wait failed");continue;}for (int i = 0; i < num; i++) {Conn *conn = (Conn *)events[i].data.ptr;if (conn->Fd() == sock_fd) {int max_conn = is_loop_accept ? 2048 : 1;LoopAccept(sock_fd, max_conn, [epoll_fd](int client_fd) {Conn *conn = new Conn(client_fd, epoll_fd, false);AddReadEvent(conn); // 监听可读事件,保持fd为阻塞IO});continue;}handlerClient(conn->Fd());ClearEvent(conn);delete conn;}}return 0;}
#include <fcntl.h>#include <stdio.h>#include <unistd.h>#include <iostream>#include <unordered_map>#include <unordered_set>#include "../../common/cmdline.h"#include "../../common/conn.hpp"using namespace std;using namespace MyEcho;void updateSet(unordered_set<int> &read_fds, unordered_set<int> &write_fds, int &max_fd, int sock_fd, fd_set &read_set,fd_set &write_set) {max_fd = sock_fd;FD_ZERO(&read_set);FD_ZERO(&write_set);for (const auto &read_fd : read_fds) {if (read_fd > max_fd) {max_fd = read_fd;}FD_SET(read_fd, &read_set);}for (const auto &write_fd : write_fds) {if (write_fd > max_fd) {max_fd = write_fd;}FD_SET(write_fd, &write_set);}}void usage() {cout << "SelectReactorSingleProcess -ip 0.0.0.0 -port 1688" << endl;cout << "options:" << endl;cout << " -h,--help print usage" << endl;cout << " -ip,--ip listen ip" << endl;cout << " -port,--port listen port" << endl;cout << endl;}int main(int argc, char *argv[]) {string ip;int64_t port;CmdLine::StrOptRequired(&ip, "ip");CmdLine::Int64OptRequired(&port, "port");CmdLine::SetUsage(usage);CmdLine::Parse(argc, argv);int sock_fd = CreateListenSocket(ip, port, false);if (sock_fd < 0) {return -1;}int max_fd;fd_set read_set;fd_set write_set;SetNotBlock(sock_fd);unordered_set<int> read_fds;unordered_set<int> write_fds;unordered_map<int, Conn *> conns;while (true) {read_fds.insert(sock_fd);updateSet(read_fds, write_fds, max_fd, sock_fd, read_set, write_set);int ret = select(max_fd + 1, &read_set, &write_set, nullptr, nullptr);if (ret <= 0) {if (ret < 0) perror("select failed");continue;}unordered_set<int> temp = read_fds;for (const auto &fd : temp) {if (not FD_ISSET(fd, &read_set)) {continue;}if (fd == sock_fd) { // 监听的sock_fd可读,则表示有新的链接LoopAccept(sock_fd, 1024, [&read_fds, &conns](int client_fd) {if (client_fd >= FD_SETSIZE) { // 大于FD_SETSIZE的值,则不支持close(client_fd);return;}read_fds.insert(client_fd); // 新增到要监听的fd集合中conns[client_fd] = new Conn(client_fd, true);});continue;}// 执行到这里,表明可读Conn *conn = conns[fd];if (not conn->Read()) { // 执行读失败delete conn;conns.erase(fd);read_fds.erase(fd);close(fd);continue;}if (conn->OneMessage()) { // 判断是否要触发写事件conn->EnCode();read_fds.erase(fd);write_fds.insert(fd);}}temp = write_fds;for (const auto &fd : temp) {if (not FD_ISSET(fd, &write_set)) {continue;}// 执行到这里,表明可写Conn *conn = conns[fd];if (not conn->Write()) { // 执行写失败delete conn;conns.erase(fd);write_fds.erase(fd);close(fd);continue;}if (conn->FinishWrite()) { // 完成了请求的应答写conn->Reset();write_fds.erase(fd);read_fds.insert(fd);}}}return 0;}
#include <poll.h>#include <stdio.h>#include <unistd.h>#include <iostream>#include <unordered_map>#include <unordered_set>#include "../../common/cmdline.h"#include "../../common/conn.hpp"#include "../../common/utils.hpp"using namespace std;using namespace MyEcho;void updateFds(unordered_set<int> &read_fds, unordered_set<int> &write_fds, pollfd **fds, int &nfds) {if (*fds != nullptr) {delete[](*fds);}nfds = read_fds.size() + write_fds.size();*fds = new pollfd[nfds];int index = 0;for (const auto &read_fd : read_fds) {(*fds)[index].fd = read_fd;(*fds)[index].events = POLLIN;(*fds)[index].revents = 0;index++;}for (const auto &write_fd : write_fds) {(*fds)[index].fd = write_fd;(*fds)[index].events = POLLOUT;(*fds)[index].revents = 0;index++;}}void usage() {cout << "PollReactorSingleProcess -ip 0.0.0.0 -port 1688" << endl;cout << "options:" << endl;cout << " -h,--help print usage" << endl;cout << " -ip,--ip listen ip" << endl;cout << " -port,--port listen port" << endl;cout << endl;}int main(int argc, char *argv[]) {string ip;int64_t port;CmdLine::StrOptRequired(&ip, "ip");CmdLine::Int64OptRequired(&port, "port");CmdLine::SetUsage(usage);CmdLine::Parse(argc, argv);int sock_fd = CreateListenSocket(ip, port, false);if (sock_fd < 0) {return -1;}int nfds = 0;pollfd *fds = nullptr;unordered_set<int> read_fds;unordered_set<int> write_fds;unordered_map<int, Conn *> conns;SetNotBlock(sock_fd);while (true) {read_fds.insert(sock_fd);updateFds(read_fds, write_fds, &fds, nfds);int ret = poll(fds, nfds, -1);if (ret <= 0) {if (ret < 0) perror("poll failed");continue;}for (int i = 0; i < nfds; i++) {if (fds[i].revents & POLLIN) {int fd = fds[i].fd;if (fd == sock_fd) {LoopAccept(sock_fd, 2048, [&read_fds, &conns](int client_fd) {read_fds.insert(client_fd); // 新增到要监听的fd集合中conns[client_fd] = new Conn(client_fd, true);});continue;}// 执行到这里,表明可读Conn *conn = conns[fd];if (not conn->Read()) { // 执行读失败delete conn;conns.erase(fd);read_fds.erase(fd);close(fd);continue;}if (conn->OneMessage()) { // 判断是否要触发写事件conn->EnCode();read_fds.erase(fd);write_fds.insert(fd);}}if (fds[i].revents & POLLOUT) { // 可写int fd = fds[i].fd;Conn *conn = conns[fd];if (not conn->Write()) { // 执行写失败delete conn;conns.erase(fd);write_fds.erase(fd);close(fd);continue;}if (conn->FinishWrite()) { // 完成了请求的应答写conn->Reset();write_fds.erase(fd);read_fds.insert(fd);}}}}return 0;}
#include <arpa/inet.h>#include <assert.h>#include <fcntl.h>#include <netinet/in.h>#include <stdio.h>#include <stdlib.h>#include <sys/epoll.h>#include <sys/socket.h>#include <unistd.h>#include <iostream>#include "../../common/cmdline.h"#include "../../common/epollctl.hpp"using namespace std;using namespace MyEcho;void usage() {cout << "EpollReactorSingleProcess -ip 0.0.0.0 -port 1688 -multiio -la -writefirst" << endl;cout << "options:" << endl;cout << " -h,--help print usage" << endl;cout << " -ip,--ip listen ip" << endl;cout << " -port,--port listen port" << endl;cout << " -multiio,--multiio multi io" << endl;cout << " -la,--la loop accept" << endl;cout << " -writefirst--writefirst write first" << endl;cout << endl;}int main(int argc, char *argv[]) {string ip;int64_t port;bool is_multi_io;bool is_loop_accept;bool is_write_first;CmdLine::StrOptRequired(&ip, "ip");CmdLine::Int64OptRequired(&port, "port");CmdLine::BoolOpt(&is_multi_io, "multiio");CmdLine::BoolOpt(&is_loop_accept, "la");CmdLine::BoolOpt(&is_write_first, "writefirst");CmdLine::SetUsage(usage);CmdLine::Parse(argc, argv);cout << "is_loop_accept = " << is_loop_accept << endl;cout << "is_multi_io = " << is_multi_io << endl;cout << "is_write_first = " << is_write_first << endl;int sock_fd = CreateListenSocket(ip, port, false);if (sock_fd < 0) {return -1;}epoll_event events[2048];int epoll_fd = epoll_create(1);if (epoll_fd < 0) {perror("epoll_create failed");return -1;}Conn conn(sock_fd, epoll_fd, is_multi_io);SetNotBlock(sock_fd);AddReadEvent(&conn);while (true) {int num = epoll_wait(epoll_fd, events, 2048, -1);if (num < 0) {perror("epoll_wait failed");continue;}for (int i = 0; i < num; i++) {Conn *conn = (Conn *)events[i].data.ptr;if (conn->Fd() == sock_fd) {int max_conn = is_loop_accept ? 2048 : 1;LoopAccept(sock_fd, max_conn, [epoll_fd, is_multi_io](int client_fd) {Conn *conn = new Conn(client_fd, epoll_fd, is_multi_io);SetNotBlock(client_fd);AddReadEvent(conn); // 监听可读事件});continue;}auto releaseConn = [&conn]() {ClearEvent(conn);delete conn;};if (events[i].events & EPOLLIN) { // 可读if (not conn->Read()) { // 执行读失败releaseConn();continue;}if (conn->OneMessage()) { // 判断是否要触发写事件conn->EnCode();if (is_write_first) { // 判断是否要先写数据if (not conn->Write()) {releaseConn();continue;}}if (conn->FinishWrite()) {conn->Reset();ModToReadEvent(conn); // 修改成只监控可读事件} else {ModToWriteEvent(conn); // 修改成只监控可写事件}}}if (events[i].events & EPOLLOUT) { // 可写if (not conn->Write()) { // 执行写失败releaseConn();continue;}if (conn->FinishWrite()) { // 完成了请求的应答写conn->Reset();ModToReadEvent(conn); // 修改成只监控可读事件}}}}return 0;}
#include <arpa/inet.h>#include <assert.h>#include <fcntl.h>#include <netinet/in.h>#include <stdio.h>#include <stdlib.h>#include <sys/epoll.h>#include <sys/socket.h>#include <unistd.h>#include <iostream>#include "../../common/cmdline.h"#include "../../common/epollctl.hpp"using namespace std;using namespace MyEcho;void usage() {cout << "EpollReactorSingleProcessET -ip 0.0.0.0 -port 1688" << endl;cout << "options:" << endl;cout << " -h,--help print usage" << endl;cout << " -ip,--ip listen ip" << endl;cout << " -port,--port listen port" << endl;cout << endl;}int main(int argc, char *argv[]) {string ip;int64_t port;CmdLine::StrOptRequired(&ip, "ip");CmdLine::Int64OptRequired(&port, "port");CmdLine::SetUsage(usage);CmdLine::Parse(argc, argv);int sock_fd = CreateListenSocket(ip, port, false);if (sock_fd < 0) {return -1;}epoll_event events[2048];int epoll_fd = epoll_create(1);if (epoll_fd < 0) {perror("epoll_create failed");return -1;}Conn conn(sock_fd, epoll_fd, true);SetNotBlock(sock_fd);AddReadEvent(&conn);while (true) {int num = epoll_wait(epoll_fd, events, 2048, -1);if (num < 0) {perror("epoll_wait failed");continue;}for (int i = 0; i < num; i++) {Conn *conn = (Conn *)events[i].data.ptr;if (conn->Fd() == sock_fd) {LoopAccept(sock_fd, 2048, [epoll_fd](int client_fd) {Conn *conn = new Conn(client_fd, epoll_fd, true);SetNotBlock(client_fd);AddReadEvent(conn, true); // 监听可读事件,开启边缘模式});continue;}auto releaseConn = [&conn]() {ClearEvent(conn);delete conn;};if (events[i].events & EPOLLIN) { // 可读if (not conn->Read()) { // 执行非阻塞读releaseConn();continue;}if (conn->OneMessage()) { // 判断是否要触发写事件conn->EnCode();ModToWriteEvent(conn, true); // 修改成只监控可写事件,开启边缘模式}}if (events[i].events & EPOLLOUT) { // 可写if (not conn->Write()) { // 执行非阻塞写releaseConn();continue;}if (conn->FinishWrite()) { // 完成了了请求的应答写,则可以释放连接conn->Reset();ModToReadEvent(conn, true); // 修改成只监控可读事件}}}}return 0;}
#include <arpa/inet.h>#include <assert.h>#include <fcntl.h>#include <netinet/in.h>#include <stdio.h>#include <stdlib.h>#include <sys/epoll.h>#include <sys/socket.h>#include <unistd.h>#include <iostream>#include "../../common/cmdline.h"#include "../../common/coroutine.h"#include "../../common/epollctl.hpp"using namespace std;using namespace MyEcho;struct EventData {EventData(int fd, int epoll_fd) : fd_(fd), epoll_fd_(epoll_fd){};int fd_{0};int epoll_fd_{0};int cid_{MyCoroutine::kInvalidRoutineId};MyCoroutine::Schedule *schedule_{nullptr};};void EchoDeal(const std::string req_message, std::string &resp_message) { resp_message = req_message; }void handlerClient(EventData *event_data) {auto releaseConn = [&event_data]() {ClearEvent(event_data->epoll_fd_, event_data->fd_);delete event_data; // 释放内存};while (true) {ssize_t ret = 0;Codec codec;string *req_message{nullptr};string resp_message;while (true) { // 读操作ret = read(event_data->fd_, codec.Data(), codec.Len());if (ret == 0) {perror("peer close connection");releaseConn();return;}if (ret < 0) {if (EINTR == errno) continue; // 被中断,可以重启读操作if (EAGAIN == errno or EWOULDBLOCK == errno) {MyCoroutine::CoroutineYield(*event_data->schedule_); // 让出cpu,切换到主协程,等待下一次数据可读continue;}perror("read failed");releaseConn();return;}codec.DeCode(ret); // 解析请求数据req_message = codec.GetMessage();if (req_message) { // 解析出一个完整的请求break;}}// 执行到这里说明已经读取到一个完整的请求EchoDeal(*req_message, resp_message); // 业务handler的封装,这样协程的调用就对业务逻辑函数EchoDeal透明delete req_message;Packet pkt;codec.EnCode(resp_message, pkt);ModToWriteEvent(event_data->epoll_fd_, event_data->fd_, event_data); // 监听可写事件。size_t sendLen = 0;while (sendLen != pkt.UseLen()) { // 写操作ret = write(event_data->fd_, pkt.Data() + sendLen, pkt.UseLen() - sendLen);if (ret < 0) {if (EINTR == errno) continue; // 被中断,可以重启写操作if (EAGAIN == errno or EWOULDBLOCK == errno) {MyCoroutine::CoroutineYield(*event_data->schedule_); // 让出cpu,切换到主协程,等待下一次数据可写continue;}perror("write failed");releaseConn();return;}sendLen += ret;}ModToReadEvent(event_data->epoll_fd_, event_data->fd_, event_data); // 监听可读事件。}}void usage() {cout << "EpollReactorSingleProcessCoroutine -ip 0.0.0.0 -port 1688 -d" << endl;cout << "options:" << endl;cout << " -h,--help print usage" << endl;cout << " -ip,--ip listen ip" << endl;cout << " -port,--port listen port" << endl;cout << " -d,--d dynamic epoll time out" << endl;cout << endl;}int main(int argc, char *argv[]) {string ip;int64_t port;bool is_dynamic_time_out{false};CmdLine::StrOptRequired(&ip, "ip");CmdLine::Int64OptRequired(&port, "port");CmdLine::BoolOpt(&is_dynamic_time_out, "d");CmdLine::SetUsage(usage);CmdLine::Parse(argc, argv);int sock_fd = CreateListenSocket(ip, port, false);if (sock_fd < 0) {return -1;}epoll_event events[2048];int epoll_fd = epoll_create(1);if (epoll_fd < 0) {perror("epoll_create failed");return -1;}cout << "is_dynamic_time_out = " << is_dynamic_time_out << endl;EventData event_data(sock_fd, epoll_fd);SetNotBlock(sock_fd);AddReadEvent(epoll_fd, sock_fd, &event_data);MyCoroutine::Schedule schedule;MyCoroutine::ScheduleInit(schedule, 5000); // 协程池初始化int msec = -1;while (true) {int num = epoll_wait(epoll_fd, events, 2048, msec);if (num < 0) {perror("epoll_wait failed");continue;} else if (num == 0) { // 没有事件了,下次调用epoll_wait大概率被挂起sleep(0); // 这里直接sleep(0)让出cpu,大概率被挂起,这里主动让出cpu,可以减少一次epoll_wait的调用msec = -1; // 大概率被挂起,故这里超时时间设置为-1continue;}if (is_dynamic_time_out) msec = 0; // 下次大概率还有事件,故msec设置为0for (int i = 0; i < num; i++) {EventData *event_data = (EventData *)events[i].data.ptr;if (event_data->fd_ == sock_fd) {LoopAccept(sock_fd, 2048, [epoll_fd](int client_fd) {EventData *event_data = new EventData(client_fd, epoll_fd);SetNotBlock(client_fd);AddReadEvent(epoll_fd, client_fd, event_data); // 监听可读事件});continue;}if (event_data->cid_ == MyCoroutine::kInvalidRoutineId) { // 第一次事件,则创建协程event_data->schedule_ = &schedule;event_data->cid_ = MyCoroutine::CoroutineCreate(schedule, handlerClient, event_data); // 创建协程MyCoroutine::CoroutineResumeById(schedule, event_data->cid_);} else {MyCoroutine::CoroutineResumeById(schedule, event_data->cid_); // 唤醒之前主动让出cpu的协程}}}return 0;}
#include <arpa/inet.h>#include <assert.h>#include <fcntl.h>#include <netinet/in.h>#include <stdio.h>#include <stdlib.h>#include <sys/epoll.h>#include <sys/socket.h>#include <unistd.h>#include <algorithm>#include <iostream>#include <vector>#include "../../common/cmdline.h"#include "../../common/conn.hpp"#include "../../common/epollctl.hpp"using namespace std;using namespace MyEcho;void mainReactor(string ip, int64_t port, bool is_main_read, int64_t sub_reactor_count) {vector<int> client_unix_sockets;for (int i = 0; i < sub_reactor_count; i++) {int fd = CreateClientUnixSocket("./unix.sock." + to_string(i));assert(fd > 0);client_unix_sockets.push_back(fd);}int index = 0;int sock_fd = CreateListenSocket(ip, port, true);assert(sock_fd > 0);epoll_event events[2048];int epoll_fd = epoll_create(1);assert(epoll_fd > 0);Conn conn(sock_fd, epoll_fd, true);SetNotBlock(sock_fd);AddReadEvent(&conn);auto getClientUnixSocketFd = [&index, &client_unix_sockets, sub_reactor_count]() {index++;index %= sub_reactor_count;return client_unix_sockets[index];};while (true) {int num = epoll_wait(epoll_fd, events, 2048, -1);if (num < 0) {perror("epoll_wait failed");continue;}for (int i = 0; i < num; i++) {Conn *conn = (Conn *)events[i].data.ptr;if (conn->Fd() == sock_fd) { // 有客户端的连接到来了LoopAccept(sock_fd, 2048, [is_main_read, epoll_fd, getClientUnixSocketFd](int client_fd) {SetNotBlock(client_fd);if (is_main_read) {Conn *conn = new Conn(client_fd, epoll_fd, true);AddReadEvent(conn); // 在mainReactor线程中监听可读事件} else {SendFd(getClientUnixSocketFd(), client_fd);close(client_fd);}});continue;}// 客户端有数据可读,则把连接迁移到subReactor线程中管理ClearEvent(conn, false);SendFd(getClientUnixSocketFd(), conn->Fd());close(conn->Fd());delete conn;}}}void subReactor(int server_unix_socket, int64_t sub_reactor_count) {epoll_event events[2048];int epoll_fd = epoll_create(1);if (epoll_fd < 0) {perror("epoll_create failed");return;}Conn conn(server_unix_socket, epoll_fd, true);SetNotBlock(server_unix_socket);AddReadEvent(&conn);while (true) {int num = epoll_wait(epoll_fd, events, 2048, -1);if (num < 0) {perror("epoll_wait failed");continue;}for (int i = 0; i < num; i++) {Conn *conn = (Conn *)events[i].data.ptr;auto releaseConn = [&conn]() {ClearEvent(conn);delete conn;};if (conn->Fd() == server_unix_socket) {// 接受从mainReactor过来的连接LoopAccept(server_unix_socket, 1024, [epoll_fd](int main_reactor_client_fd) {Conn *conn = new Conn(main_reactor_client_fd, epoll_fd, true);conn->SetUnixSocket();AddReadEvent(conn);cout << "accept mainReactor unix_socet connect. pid = " << getpid() << endl;});continue;}if (conn->IsUnixSocket()) {int client_fd = 0;// 接收从mainReactor传递过来的客户端连接fdif (0 == RecvFd(conn->Fd(), client_fd)) {Conn *conn = new Conn(client_fd, epoll_fd, true);AddReadEvent(conn);}continue;}// 执行到这里就是真正的客户端的读写事件if (events[i].events & EPOLLIN) { // 可读if (not conn->Read()) { // 执行非阻塞读releaseConn();continue;}if (conn->OneMessage()) { // 判断是否要触发写事件conn->EnCode();ModToWriteEvent(conn); // 修改成只监控可写事件}}if (events[i].events & EPOLLOUT) { // 可写if (not conn->Write()) { // 执行非阻塞写releaseConn();continue;}if (conn->FinishWrite()) { // 完成了请求的应答写,则可以释放连接conn->Reset();ModToReadEvent(conn); // 修改成只监控可读事件}}}}}void createServerUnixSocket(vector<int> &server_unix_sockets, int64_t sub_reactor_count) {for (int i = 0; i < sub_reactor_count; i++) {string path = "./unix.sock." + to_string(i);remove(path.c_str());int server_unix_socket = CreateListenUnixSocket(path);assert(server_unix_socket > 0);server_unix_sockets.push_back(server_unix_socket);}}void createSubReactor(vector<int> &server_unix_sockets, int64_t sub_reactor_count) {for (int i = 0; i < sub_reactor_count; i++) {pid_t pid = fork();assert(pid != -1);if (pid == 0) { // 子进程int fd = server_unix_sockets[i];// 关闭不需要的fd,避免fd泄漏for_each(server_unix_sockets.begin(), server_unix_sockets.end(), [fd](int server_unix_socket_fd) {if (server_unix_socket_fd != fd) {close(server_unix_socket_fd);}});cout << "subReactor pid = " << getpid() << endl;subReactor(fd, sub_reactor_count);exit(0);}}}void createMainReactor(string ip, int64_t port, bool is_main_read, int64_t sub_reactor_count,int64_t main_reactor_count) {for (int i = 0; i < main_reactor_count; i++) {pid_t pid = fork();assert(pid != -1);if (pid == 0) { // 子进程cout << "mainReactor pid = " << getpid() << endl;mainReactor(ip, port, is_main_read, sub_reactor_count);exit(0);}}}void usage() {cout << "EpollReactorProcessPoolMS -ip 0.0.0.0 -port 1688 -main 3 -sub 8 -mainread" << endl;cout << "options:" << endl;cout << " -h,--help print usage" << endl;cout << " -ip,--ip listen ip" << endl;cout << " -port,--port listen port" << endl;cout << " -main,--main mainReactor count" << endl;cout << " -sub,--sub subReactor count" << endl;cout << " -mainread,--mainread mainReactor read" << endl;cout << endl;}int main(int argc, char *argv[]) {string ip;int64_t port;int64_t main_reactor_count;int64_t sub_reactor_count;bool is_main_read;CmdLine::StrOptRequired(&ip, "ip");CmdLine::Int64OptRequired(&port, "port");CmdLine::Int64OptRequired(&main_reactor_count, "main");CmdLine::Int64OptRequired(&sub_reactor_count, "sub");CmdLine::BoolOpt(&is_main_read, "mainread");CmdLine::SetUsage(usage);CmdLine::Parse(argc, argv);main_reactor_count = main_reactor_count > GetNProcs() ? GetNProcs() : main_reactor_count;sub_reactor_count = sub_reactor_count > GetNProcs() ? GetNProcs() : sub_reactor_count;vector<int> server_unix_sockets;createServerUnixSocket(server_unix_sockets, sub_reactor_count);createSubReactor(server_unix_sockets, sub_reactor_count); // 创建SubReactor进程// 不再需要这些fd,需要及时关闭,避免fd泄漏for_each(server_unix_sockets.begin(), server_unix_sockets.end(), [](int fd) { close(fd); });createMainReactor(ip, port, is_main_read, sub_reactor_count, main_reactor_count); // 创建MainRector进程while (true) sleep(1); // 主进程陷入死循环return 0;}
#include <arpa/inet.h>#include <assert.h>#include <fcntl.h>#include <netinet/in.h>#include <stdio.h>#include <stdlib.h>#include <sys/epoll.h>#include <sys/socket.h>#include <unistd.h>#include <iostream>#include "../../common/cmdline.h"#include "../../common/coroutine.h"#include "../../common/epollctl.hpp"using namespace std;using namespace MyEcho;struct EventData {EventData(int fd, int epoll_fd) : fd_(fd), epoll_fd_(epoll_fd){};int fd_{0};int epoll_fd_{0};int cid_{MyCoroutine::kInvalidRoutineId};MyCoroutine::Schedule *schedule_{nullptr};};void EchoDeal(const std::string req_message, std::string &resp_message) { resp_message = req_message; }void handlerClient(EventData *event_data) {auto releaseConn = [&event_data]() {ClearEvent(event_data->epoll_fd_, event_data->fd_);delete event_data; // 释放内存};while (true) {ssize_t ret = 0;Codec codec;string *req_message{nullptr};string resp_message;while (true) { // 读操作ret = read(event_data->fd_, codec.Data(), codec.Len()); // 一次最多读取100字节if (ret == 0) {perror("peer close connection");releaseConn();return;}if (ret < 0) {if (EINTR == errno) continue; // 被中断,可以重启读操作if (EAGAIN == errno or EWOULDBLOCK == errno) {MyCoroutine::CoroutineYield(*event_data->schedule_); // 让出cpu,切换到主协程,等待下一次数据可读continue;}perror("read failed");releaseConn();return;}codec.DeCode(ret); // 解析请求数据req_message = codec.GetMessage();if (req_message) { // 解析出一个完整的请求break;}}// 执行到这里说明已经读取到一个完整的请求EchoDeal(*req_message, resp_message); // 业务handler的封装,这样协程的调用就对业务逻辑函数EchoDeal透明delete req_message;Packet pkt;codec.EnCode(resp_message, pkt);ModToWriteEvent(event_data->epoll_fd_, event_data->fd_, event_data); // 监听可写事件。size_t sendLen = 0;while (sendLen != pkt.UseLen()) { // 写操作ret = write(event_data->fd_, pkt.Data() + sendLen, pkt.UseLen() - sendLen);if (ret < 0) {if (EINTR == errno) continue; // 被中断,可以重启写操作if (EAGAIN == errno or EWOULDBLOCK == errno) {MyCoroutine::CoroutineYield(*event_data->schedule_); // 让出cpu,切换到主协程,等待下一次数据可写continue;}perror("write failed");releaseConn();return;}sendLen += ret;}ModToReadEvent(event_data->epoll_fd_, event_data->fd_, event_data); // 监听可读事件。}}int handler(string ip, int64_t port) {int sock_fd = CreateListenSocket(ip, port, true);if (sock_fd < 0) {return -1;}epoll_event events[2048];int epoll_fd = epoll_create(1);if (epoll_fd < 0) {perror("epoll_create failed");return -1;}EventData event_data(sock_fd, epoll_fd);SetNotBlock(sock_fd);AddReadEvent(epoll_fd, sock_fd, &event_data);MyCoroutine::Schedule schedule;MyCoroutine::ScheduleInit(schedule, 5000); // 协程池初始化int msec = -1;while (true) {int num = epoll_wait(epoll_fd, events, 2048, msec);if (num < 0) {perror("epoll_wait failed");continue;} else if (num == 0) { // 没有事件了,下次调用epoll_wait大概率被挂起sleep(0); // 这里直接sleep(0)让出cpu,大概率被挂起,这里主动让出cpu,可以减少一次epoll_wait的调用msec = -1; // 大概率被挂起,故这里超时时间设置为-1continue;}msec = 0; // 下次大概率还有事件,故msec设置为0for (int i = 0; i < num; i++) {EventData *event_data = (EventData *)events[i].data.ptr;if (event_data->fd_ == sock_fd) {LoopAccept(sock_fd, 2048, [epoll_fd](int client_fd) {EventData *event_data = new EventData(client_fd, epoll_fd);SetNotBlock(client_fd);AddReadEvent(epoll_fd, client_fd, event_data); // 监听可读事件});continue;}if (event_data->cid_ == MyCoroutine::kInvalidRoutineId) { // 第一次事件,则创建协程event_data->schedule_ = &schedule;event_data->cid_ = MyCoroutine::CoroutineCreate(schedule, handlerClient, event_data); // 创建协程MyCoroutine::CoroutineResumeById(schedule, event_data->cid_);} else {MyCoroutine::CoroutineResumeById(schedule, event_data->cid_); // 唤醒之前主动让出cpu的协程}}}return 0;}void usage() {cout << "EpollReactorProcessPoolCoroutine -ip 0.0.0.0 -port 1688 -poolsize 8" << endl;cout << "options:" << endl;cout << " -h,--help print usage" << endl;cout << " -ip,--ip listen ip" << endl;cout << " -port,--port listen port" << endl;cout << " -poolsize,--poolsize pool size" << endl;cout << endl;}int main(int argc, char *argv[]) {string ip;int64_t port;int64_t pool_size;CmdLine::StrOptRequired(&ip, "ip");CmdLine::Int64OptRequired(&port, "port");CmdLine::Int64OptRequired(&pool_size, "poolsize");CmdLine::SetUsage(usage);CmdLine::Parse(argc, argv);pool_size = pool_size > GetNProcs() ? GetNProcs() : pool_size;for (int i = 0; i < pool_size; i++) {pid_t pid = fork();assert(pid != -1);if (0 == pid) {handler(ip, port); // 子进程陷入死循环,处理客户端请求exit(0);}}while (true) sleep(1); // 父进程陷入死循环return 0;}
#include <arpa/inet.h>#include <assert.h>#include <fcntl.h>#include <netinet/in.h>#include <stdio.h>#include <stdlib.h>#include <sys/epoll.h>#include <sys/socket.h>#include <unistd.h>#include <iostream>#include <thread>#include "../../common/cmdline.h"#include "../../common/conn.hpp"#include "../../common/epollctl.hpp"using namespace std;using namespace MyEcho;void handler(string ip, int64_t port) {int sock_fd = CreateListenSocket(ip, port, true);if (sock_fd < 0) {return;}epoll_event events[2048];int epoll_fd = epoll_create(1);if (epoll_fd < 0) {perror("epoll_create failed");return;}Conn conn(sock_fd, epoll_fd, true);SetNotBlock(sock_fd);AddReadEvent(&conn);while (true) {int num = epoll_wait(epoll_fd, events, 2048, -1);if (num < 0) {perror("epoll_wait failed");continue;}for (int i = 0; i < num; i++) {Conn *conn = (Conn *)events[i].data.ptr;if (conn->Fd() == sock_fd) {LoopAccept(sock_fd, 2048, [epoll_fd](int client_fd) {Conn *conn = new Conn(client_fd, epoll_fd, true);SetNotBlock(client_fd);AddReadEvent(conn); // 监听可读事件});continue;}auto releaseConn = [&conn]() {ClearEvent(conn);delete conn;};if (events[i].events & EPOLLIN) { // 可读if (not conn->Read()) { // 执行读失败releaseConn();continue;}if (conn->OneMessage()) { // 判断是否要触发写事件conn->EnCode();ModToWriteEvent(conn); // 修改成只监控可写事件}}if (events[i].events & EPOLLOUT) { // 可写if (not conn->Write()) { // 执行写失败releaseConn();continue;}if (conn->FinishWrite()) { // 完成了请求的应答写,则可以释放连接conn->Reset();ModToReadEvent(conn); // 修改成只监控可读事件}}}}}void usage() {cout << "EpollReactorThreadPool -ip 0.0.0.0 -port 1688 -poolsize 8" << endl;cout << "options:" << endl;cout << " -h,--help print usage" << endl;cout << " -ip,--ip listen ip" << endl;cout << " -port,--port listen port" << endl;cout << " -poolsize,--poolsize pool size" << endl;cout << endl;}int main(int argc, char *argv[]) {string ip;int64_t port;int64_t pool_size;CmdLine::StrOptRequired(&ip, "ip");CmdLine::Int64OptRequired(&port, "port");CmdLine::Int64OptRequired(&pool_size, "poolsize");CmdLine::SetUsage(usage);CmdLine::Parse(argc, argv);pool_size = pool_size > GetNProcs() ? GetNProcs() : pool_size;for (int i = 0; i < pool_size; i++) {std::thread(handler, ip, port).detach(); // 这里需要调用detach,让创建的线程独立运行}while (true) sleep(1); // 主线程陷入死循环return 0;}
#include <arpa/inet.h>#include <fcntl.h>#include <netinet/in.h>#include <stdio.h>#include <stdlib.h>#include <sys/epoll.h>#include <sys/socket.h>#include <unistd.h>#include <condition_variable>#include <iostream>#include <mutex>#include <queue>#include <thread>#include "../../common/cmdline.h"#include "../../common/conn.hpp"#include "../../common/epollctl.hpp"using namespace std;using namespace MyEcho;std::mutex Mutex;std::condition_variable Cond;std::queue<Conn *> Queue;void pushInQueue(Conn *conn) {{std::unique_lock<std::mutex> locker(Mutex);Queue.push(conn);}Cond.notify_one();}Conn *getQueueData() {std::unique_lock<std::mutex> locker(Mutex);Cond.wait(locker, []() -> bool { return Queue.size() > 0; });Conn *conn = Queue.front();Queue.pop();return conn;}void workerHandler(bool is_direct) {while (true) {Conn *conn = getQueueData();conn->EnCode();if (is_direct) { // 直接把数据发送给客户端,而不是通过I/O线程来发送bool success = true;while (not conn->FinishWrite()) {if (not conn->Write()) {success = false;break;}}if (not success) {ClearEvent(conn);delete conn;} else {conn->Reset();ReStartReadEvent(conn); // 修改成只监控可读事件,携带oneshot选项}} else {ModToWriteEvent(conn); // 监听写事件,数据通过I/O线程来发送}}}void ioHandler(string ip, int64_t port) {int sock_fd = CreateListenSocket(ip, port, true);if (sock_fd < 0) {return;}epoll_event events[2048];int epoll_fd = epoll_create(1);if (epoll_fd < 0) {perror("epoll_create failed");return;}Conn conn(sock_fd, epoll_fd, true);SetNotBlock(sock_fd);AddReadEvent(&conn);int msec = -1;while (true) {int num = epoll_wait(epoll_fd, events, 2048, msec);if (num < 0) {perror("epoll_wait failed");continue;}for (int i = 0; i < num; i++) {Conn *conn = (Conn *)events[i].data.ptr;if (conn->Fd() == sock_fd) {LoopAccept(sock_fd, 2048, [epoll_fd](int client_fd) {Conn *conn = new Conn(client_fd, epoll_fd, true);SetNotBlock(client_fd);AddReadEvent(conn, false, true); // 监听可读事件,开启oneshot});continue;}auto releaseConn = [&conn]() {ClearEvent(conn);delete conn;};if (events[i].events & EPOLLIN) { // 可读if (not conn->Read()) { // 执行非阻塞readreleaseConn();continue;}if (conn->OneMessage()) {pushInQueue(conn); // 入共享输入队列,有锁} else {ReStartReadEvent(conn); // 还没收到完整的请求,则重新启动可读事件的监听,携带oneshot选项}}if (events[i].events & EPOLLOUT) { // 可写if (not conn->Write()) { // 执行非阻塞writereleaseConn();continue;}if (conn->FinishWrite()) { // 完成了请求的应答写,则可以释放连接closeconn->Reset();ReStartReadEvent(conn); // 修改成只监控可读事件,携带oneshot选项}}}}}void usage() {cout << "EpollReactorThreadPoolHSHA -ip 0.0.0.0 -port 1688 -io 3 -worker 8 -direct" << endl;cout << "options:" << endl;cout << " -h,--help print usage" << endl;cout << " -ip,--ip listen ip" << endl;cout << " -port,--port listen port" << endl;cout << " -io,--io io thread count" << endl;cout << " -worker,--worker worker thread count" << endl;cout << " -direct,--direct direct send response data by worker thread" << endl;cout << endl;}int main(int argc, char *argv[]) {string ip;int64_t port;int64_t io_count;int64_t worker_count;bool is_direct;CmdLine::StrOptRequired(&ip, "ip");CmdLine::Int64OptRequired(&port, "port");CmdLine::Int64OptRequired(&io_count, "io");CmdLine::Int64OptRequired(&worker_count, "worker");CmdLine::BoolOpt(&is_direct, "direct");CmdLine::SetUsage(usage);CmdLine::Parse(argc, argv);cout << "is_direct=" << is_direct << endl;io_count = io_count > GetNProcs() ? GetNProcs() : io_count;worker_count = worker_count > GetNProcs() ? GetNProcs() : worker_count;for (int i = 0; i < worker_count; i++) { // 创建worker线程std::thread(workerHandler, is_direct).detach(); // 这里需要调用detach,让创建的线程独立运行}for (int i = 0; i < io_count; i++) { // 创建io线程std::thread(ioHandler, ip, port).detach(); // 这里需要调用detach,让创建的线程独立运行}while (true) sleep(1); // 主线程陷入死循环return 0;}
#include <arpa/inet.h>#include <assert.h>#include <fcntl.h>#include <netinet/in.h>#include <stdio.h>#include <stdlib.h>#include <sys/epoll.h>#include <sys/socket.h>#include <unistd.h>#include <iostream>#include <thread>#include "../../common/cmdline.h"#include "../../common/conn.hpp"#include "../../common/epollctl.hpp"using namespace std;using namespace MyEcho;int *EpollFd;int createEpoll() { int epoll_fd = epoll_create(1); assert(epoll_fd > 0); return epoll_fd;}void addToSubReactor(int &index, int sub_reactor_count, int client_fd) { index++; index %= sub_reactor_count; // 轮询的方式添加到subReactor线程中 Conn *conn = new Conn(client_fd, EpollFd[index], true); AddReadEvent(conn); // 监听可读事件}void mainReactor(string ip, int64_t port, int64_t sub_reactor_count, bool is_main_read) { int sock_fd = CreateListenSocket(ip, port, true); if (sock_fd < 0) { return; } epoll_event events[2048]; int epoll_fd = epoll_create(1); if (epoll_fd < 0) { perror("epoll_create failed"); return; } int index = 0; Conn conn(sock_fd, epoll_fd, true); SetNotBlock(sock_fd); AddReadEvent(&conn); while (true) { int num = epoll_wait(epoll_fd, events, 2048, -1); if (num < 0) { perror("epoll_wait failed"); continue; } for (int i = 0; i < num; i++) { Conn *conn = (Conn *)events[i].data.ptr; if (conn->Fd() == sock_fd) { // 有客户端的连接到来了 LoopAccept(sock_fd, 2048, [&index, is_main_read, epoll_fd, sub_reactor_count](int client_fd) { SetNotBlock(client_fd); if (is_main_read) { Conn *conn = new Conn(client_fd, epoll_fd, true); AddReadEvent(conn); // 在mainReactor线程中监听可读事件 } else { addToSubReactor(index, sub_reactor_count, client_fd); } }); continue; } // 客户端有数据可读,则把连接迁移到subReactor线程中管理 ClearEvent(conn, false); addToSubReactor(index, sub_reactor_count, conn->Fd()); delete conn; } }}void subReactor(int thread_id) { epoll_event events[2048]; int epoll_fd = EpollFd[thread_id]; while (true) { int num = epoll_wait(epoll_fd, events, 2048, -1); if (num < 0) { perror("epoll_wait failed"); continue; } for (int i = 0; i < num; i++) { Conn *conn = (Conn *)events[i].data.ptr; auto releaseConn = [&conn]() { ClearEvent(conn); delete conn; }; if (events[i].events & EPOLLIN) { // 可读 if (not conn->Read()) { // 执行非阻塞读 releaseConn(); continue; } if (conn->OneMessage()) { // 判断是否要触发写事件 conn->EnCode(); ModToWriteEvent(conn); // 修改成只监控可写事件 } } if (events[i].events & EPOLLOUT) { // 可写 if (not conn->Write()) { // 执行非阻塞写 releaseConn(); continue; } if (conn->FinishWrite()) { // 完成了请求的应答写,则可以释放连接 conn->Reset(); ModToReadEvent(conn); // 修改成只监控可读事件 } } } }}void usage() { cout << "EpollReactorThreadPoolMS -ip 0.0.0.0 -port 1688 -main 3 -sub 8 -mainread" << endl; cout << "options:" << endl; cout << " -h,--help print usage" << endl; cout << " -ip,--ip listen ip" << endl; cout << " -port,--port listen port" << endl; cout << " -main,--main mainReactor count" << endl; cout << " -sub,--sub subReactor count" << endl; cout << " -mainread,--mainread mainReactor read" << endl; cout << endl;}int main(int argc, char *argv[]) { string ip; int64_t port; int64_t main_reactor_count; int64_t sub_reactor_count; bool is_main_read; CmdLine::StrOptRequired(&ip, "ip"); CmdLine::Int64OptRequired(&port, "port"); CmdLine::Int64OptRequired(&main_reactor_count, "main"); CmdLine::Int64OptRequired(&sub_reactor_count, "sub"); CmdLine::BoolOpt(&is_main_read, "mainread"); CmdLine::SetUsage(usage); CmdLine::Parse(argc, argv); cout << "is_main_read=" << is_main_read << endl; main_reactor_count = main_reactor_count > GetNProcs() ? GetNProcs() : main_reactor_count; sub_reactor_count = sub_reactor_count > GetNProcs() ? GetNProcs() : sub_reactor_count; EpollFd = new int[sub_reactor_count]; for (int i = 0; i < sub_reactor_count; i++) { EpollFd[i] = createEpoll(); std::thread(subReactor, i).detach(); // 这里需要调用detach,让创建的线程独立运行 } for (int i = 0; i < main_reactor_count; i++) { std::thread(mainReactor, ip, port, sub_reactor_count, is_main_read) .detach(); // 这里需要调用detach,让创建的线程独立运行 } while (true) sleep(1); // 主线程陷入死循环 return 0;}06
初学者可以通过阅读本书快速掌握 Linux C/C++ 后端研发的核心技能,并直接从事相关岗位的研发工作。对于初级、中级或者高级后端研发工程师来说,本书也能够帮助他们快速提升技术水平,完善自身的技术知识体系,并在实践中掌握后端研发的最佳实践。无论您是想要入门 Linux 后端研发,还是想要深入了解这个领域的读者,本书都将为您提供有价值的学习资源。
📢📢欢迎加入腾讯云开发者社群,享前沿资讯、大咖干货,找兴趣搭子,交同城好友,更有鹅厂招聘机会、限量周边好礼等你来~
(长按图片立即扫码)