网络IO 之 reactor

简介: 网络IO 之 reactor

百万级高并发,可以使用epoll模型,那么io是存放在哪里的?

   reactor反应堆

 

如何实现reactor反应堆?

1、将epoll 里面的客户端fd 节点,做成event,event里面携带io的相关数据

2、accept的时候,只关心,读取数据和发送数据

一般情况下,客户端连接上服务器之后,会先发数据给服务器,因此连接后,服务器这边的做法先是 读取客户端套接字里面的信息,再是做send操作

 

编码

写一个demo

#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <errno.h>
#include <time.h>
#include <unistd.h>
#include <sys/types.h>
#include <sys/socket.h>
#include <arpa/inet.h>
#include <fcntl.h>
#include <sys/epoll.h>
//1、定义数据结构
#define BUFFER_LENGTH 1024
#define MAX_EPOLL_EVENTS 1024
typedef int (*QBCALLBACK)(int fd,int events, void *arg);
//用于用户层的event
struct qbevent{
  int fd;
  int events;
  void *arg;
  int (*callback)(int fd,int events, void *arg);
  int status;
  char buff[BUFFER_LENGTH];
  int length;
  long last_alive;
};
struct qbreactor{
  int epollfd;
  struct qbevent * events;
};
//回调函数,发送数据和 接受数据的
int send_cb(int fd,int events, void *arg);
int recv_cb(int fd,int events, void *arg);
//2、写event的 set,add,del
//将传入进来的qbevent指针对应的参数进行初始化
void qb_event_set(struct qbevent *ev,int fd,QBCALLBACK call,void * arg)
{
  ev->fd = fd;
  ev->arg = arg;
  ev->events = 0;
  ev->callback = call;
  ev->status = 0;
  ev->last_alive= time(NULL);
  return ;
}
//将用户的event 挂到 内核epoll的event上
int qb_event_add(int epollfd,int event,struct qbevent *ev)
{
  if(epollfd < 0) return -1;
  if(ev == NULL)  return -1;
  struct epoll_event epev={0,{0}};
  epev.data.ptr = ev;
  epev.events = ev->events =  event; //用户event 和 内核的event 分别赋值
  int opt;
  if(ev->status == 1)
  {
    opt = EPOLL_CTL_MOD;
  }
  else
  {
    opt = EPOLL_CTL_ADD;
    ev->status = 1;
  }
  if(epoll_ctl(epollfd,opt,ev->fd,&epev) < 0)
  {
    printf("event add failed [fd=%d], events[%d]\n", ev->fd,ev->events);
    return -1;
  }
  return 0;
}
//3、从内核epol l   event 中删除 客户端的fd,此处为用户的event
int qb_event_del(int epollfd,struct qbevent *ev)
{
  if(epollfd < 0) return  -1;
  if(ev == NULL)  return -1;
  if(ev->status != 1) return -1;
  struct epoll_event epev={0,{0}};;
  epev.data.ptr = ev;
  ev->status = 0;
  epoll_ctl(epollfd,EPOLL_CTL_DEL,ev->fd,&epev);
  return 0;
}
//回调函数 发送数据
int send_cb(int fd,int events, void *arg)
{
  struct qbreactor *reactor = (struct qbreactor *)arg;
  struct qbevent *ev = reactor->events +fd;
  int len = send(ev->fd,ev->buff,ev->length,0);
  if(len > 0){
  //说明数据发送ok
    printf("send[fd=%d], [%d]%s\n", fd, len, ev->buff);
    qb_event_del(reactor->epollfd,ev);
    qb_event_set(ev,ev->fd,recv_cb,reactor);
    qb_event_add(reactor->epollfd,EPOLLIN,ev);
  }
  else{
  //len < 0 可能写情况是,缓冲区里面需要send的数据是有的,
  //但是被别的事件提前send出去了,因此本次的send会返回-1
    close(ev->fd);
    qb_event_del(reactor->epollfd,ev);
    printf("send [fd=%d] error %s\n", fd, strerror(errno));
  }
  return 0;
}
//回调函数 接收数据
int recv_cb(int fd,int events, void *arg)
{
  struct qbreactor *reactor = (struct qbreactor *)arg;
  struct qbevent *ev = reactor->events +fd;
  int len = recv(ev->fd,ev->buff,BUFFER_LENGTH,0);
  qb_event_del(reactor->epollfd,ev);
  if(len >0){
    ev->buff[len] = '\0';
    ev->length = len;
  //可以进行判断包是否完整,黏包,拆包的事情在这里做
  //...
  printf("recv [%d]:%s\n", fd, ev->buff);
  //服务器读完了客户端的数据后,想客户端发送数据
  qb_event_set(ev,ev->fd, send_cb, reactor);
  qb_event_add(reactor->epollfd,EPOLLOUT,ev);
  }
  else if(len < 0){
    close(ev->fd);
    printf("recv[fd=%d] error[%d]:%s\n", fd, errno, strerror(errno));
  }
  else{
    close(ev->fd);
    printf("[fd=%d] pos[%ld], closed\n", fd, ev-reactor->events);
  }
  return 0;
}
//回调函数 监听新的连接
int accept_cb(int fd,int events, void *arg)
{
  struct qbreactor * reactor = (struct qbreactor *)arg;
  if (reactor == NULL) return -1;
//初始化客户端的地址结构
  struct sockaddr_in clientaddr;
  socklen_t len = sizeof(clientaddr);
  int clientfd;
  if ((clientfd = accept(fd, (struct sockaddr*)&clientaddr, &len)) == -1) {
    if (errno != EAGAIN && errno != EINTR) {
      //数据被别的事件读取了
    }
    printf("accept: %s\n", strerror(errno));
    return -1;
  }
  //开始 把clientfd 放到用户层事件里面
  int i = 0;
  do{
    //找到clientfd 需要加入的位置
    for(i = 0;i<MAX_EPOLL_EVENTS;i++){
      if(reactor->events[i].status == 0)
          break;
    }
    if(i == MAX_EPOLL_EVENTS){
      printf("%s: max connect limit[%d]\n", __func__, MAX_EPOLL_EVENTS);
      break;
    }
    int flag = 0;
    if ((flag = fcntl(clientfd, F_SETFL, O_NONBLOCK)) < 0) {
      printf("%s: fcntl nonblocking failed, %d\n", __func__, MAX_EPOLL_EVENTS);
      break;
    }
    //printf(" ii == %d -- client fd == %d , sockfd = %d\n",i,clientfd,fd);
    qb_event_set(&reactor->events[clientfd],clientfd,recv_cb,reactor);
    qb_event_add(reactor->epollfd,EPOLLIN,&reactor->events[clientfd]);
  }while(0);
  printf("new connect [%s:%d][time:%ld], pos[%d]\n", 
    inet_ntoa(clientaddr.sin_addr), ntohs(clientaddr.sin_port), reactor->events[i].last_alive, i);
  return 0;
}
//4、sock_init 初始化
int sock_init(int port)
{
  int sockfd = socket(AF_INET,SOCK_STREAM,0);
  if(sockfd < 0){
    perror("sockfd socket error\n");
    return -1;
  }
  //设置为非阻塞
  fcntl(sockfd,F_SETFL,O_NONBLOCK);
  //初始化属性
  struct sockaddr_in seraddr;
  memset(&seraddr,0,sizeof(seraddr));
  seraddr.sin_family = AF_INET;
  seraddr.sin_port = htons(port);
  seraddr.sin_addr.s_addr = INADDR_ANY;
  //绑定属性
  if(bind(sockfd,(struct sockaddr*)&seraddr,sizeof(seraddr)) < 0){
    perror("bind error");
    return -1;
  }
  //设置最大监听数量
  if(listen(sockfd,20)<0){
    perror("listen error");
    return -1;
  }
  return sockfd;
}
//5、qbreactor的初始化
int qbreactor_init(struct qbreactor *reactor)
{
  if(reactor == NULL) return -1;
  memset(reactor,0,sizeof(struct qbreactor));
  reactor->epollfd = epoll_create(1);
  if(reactor->epollfd <= 0 ){
    printf("create epfd in %s err %s\n", __func__, strerror(errno));
    return -1;
  }
  reactor->events = (struct qbevent *)malloc(sizeof(struct qbevent )* MAX_EPOLL_EVENTS);
  if(reactor->events == NULL ){
    printf("malloc events in %s err %s\n", __func__, strerror(errno));
    close(reactor->epollfd);
    return -1;
  }
  return 0;
}
//6、qbreactor destroy
void qbreactor_destroy(struct qbreactor *reactor)
{
  if(reactor == NULL) return ;
  if(reactor->events == NULL) return ;
  close(reactor->epollfd);
  free(reactor->events);
  return ;
}
//7、qbreactor的listen + 传参回调 //添加listen
int qbreactor_addlistener(struct qbreactor * reactor,int sockfd,QBCALLBACK accept)
{
  if(reactor == NULL) return -1;
  if(sockfd < 0) return -1;
  if(accept == NULL) return -1;
  qb_event_set(&reactor->events[sockfd],sockfd,accept,reactor);
  qb_event_add(reactor->epollfd,EPOLLIN,&reactor->events[sockfd]);//作为用来连接的,使用水平触发
  return 0;
}
//8、 run qbreactor
 void qbreactor_runner(struct qbreactor *reactor)
{
  if(reactor == NULL) return ;
  if(reactor->epollfd <= 0) return ;
  if(reactor->events == NULL) return ;
//阻塞等待事件发生
  struct epoll_event epev[MAX_EPOLL_EVENTS+1]={0};
  while(1){
    int nready = epoll_wait(reactor->epollfd,epev,MAX_EPOLL_EVENTS,-1);
    for(int i = 0;i<nready;i++){
      struct qbevent * ev = (struct qbevent *) epev[i].data.ptr;
      //只关心读事件和写事件
      //读事件 EPOLLIN
      if((ev->events & EPOLLIN) && (epev[i].events & EPOLLIN )){
        ev->callback(ev->fd,ev->events,ev->arg);
      }
      //写事件 EPOLLOUT
      if((ev->events & EPOLLOUT) && (epev[i].events & EPOLLOUT )){
        ev->callback(ev->fd,ev->events,ev->arg);
      }
    }
  }
  return ;
}
 int main(int argc,char * argv[])
{
//串口号从命令行传入
  if(argc < 2){
    perror("please input port number\n");
    return -1;
  }
//获取用于连接的套接字
  int sockfd = sock_init(atoi(argv[1]));
//初始化reactor
  struct qbreactor *reactor = (struct qbreactor *)malloc(sizeof(struct qbreactor));
  qbreactor_init(reactor);  
//注册监听
  qbreactor_addlistener(reactor,sockfd,accept_cb);
//开始运行
  qbreactor_runner(reactor);
//摧毁reactor
  qbreactor_destroy(reactor);
  close(sockfd);
  return 0;
}

 

 

 

相关文章
基于Reactor模型的高性能网络库之地址篇
这段代码定义了一个 InetAddress 类,是 C++ 网络编程中用于封装 IPv4 地址和端口的常见做法。该类的主要作用是方便地表示和操作一个网络地址(IP + 端口)
453 58
|
网络协议 算法 Java
基于Reactor模型的高性能网络库之Tcpserver组件-上层调度器
TcpServer 是一个用于管理 TCP 连接的类,包含成员变量如事件循环(EventLoop)、连接池(ConnectionMap)和回调函数等。其主要功能包括监听新连接、设置线程池、启动服务器及处理连接事件。通过 Acceptor 接收新连接,并使用轮询算法将连接分配给子事件循环(subloop)进行读写操作。调用链从 start() 开始,经由线程池启动和 Acceptor 监听,最终由 TcpConnection 管理具体连接的事件处理。
399 2
|
负载均衡 算法 安全
基于Reactor模式的高性能网络库之线程池组件设计篇
EventLoopThreadPool 是 Reactor 模式中实现“一个主线程 + 多个工作线程”的关键组件,用于高效管理多个 EventLoop 并在多核 CPU 上分担高并发 I/O 压力。通过封装 Thread 类和 EventLoopThread,实现线程创建、管理和事件循环的调度,形成线程池结构。每个 EventLoopThread 管理一个子线程与对应的 EventLoop(subloop),主线程(base loop)通过负载均衡算法将任务派发至各 subloop,从而提升系统性能与并发处理能力。
610 3
基于Reactor模型的高性能网络库之Tcpconnection组件
TcpConnection 由 subLoop 管理 connfd,负责处理具体连接。它封装了连接套接字,通过 Channel 监听可读、可写、关闭、错误等
331 1
基于Reactor模型的高性能网络库之Poller(EpollPoller)组件
封装底层 I/O 多路复用机制(如 epoll)的抽象类 Poller,提供统一接口支持多种实现。Poller 是一个抽象基类,定义了 Channel 管理、事件收集等核心功能,并与 EventLoop 绑定。其子类 EPollPoller 实现了基于 epoll 的具体操作,包括事件等待、Channel 更新和删除等。通过工厂方法可创建默认的 Poller 实例,实现多态调用。
519 60
基于Reactor模型的高性能网络库之Channel组件篇
Channel 是事件通道,它绑定某个文件描述符 fd,注册感兴趣的事件(如读/写),并在事件发生时分发给对应的回调函数。
556 60
|
安全 调度
基于Reactor模型的高性能网络库之核心调度器:EventLoop组件
它负责:监听事件(如 I/O 可读写、定时器)、分发事件、执行回调、管理事件源 Channel 等。
541 57
基于Reactor模型的高性能网络库之时间篇
是一个用于表示时间戳(精确到微秒)**的简单封装类
344 57
|
监控 应用服务中间件 Linux
掌握并发模型:深度揭露网络IO复用并发模型的原理。
总结,网络 I/O 复用并发模型通过实现非阻塞 I/O、引入 I/O 复用技术如 select、poll 和 epoll,以及采用 Reactor 模式等技巧,为多任务并发提供了有效的解决方案。这样的模型有效提高了系统资源利用率,以及保证了并发任务的高效执行。在现实中,这种模型在许多网络应用程序和分布式系统中都取得了很好的应用成果。
394 35