当前位置:首页 > PHP教程 > PHP总结归纳

深入剖析redis事件驱动

概述 redis 内部有一个小型的事件驱动,它和 libevent 网络库的事件驱动一样,都是依托 i/o 多路复用技术支撑起来的。 利用 i/o 多路复用技术,监听感兴趣的文件 i/o 事件,例如读事件,写事件等,同时也要维护一个以文件描述符为主键,数据为某个预设函数的

概述

redis 内部有一个小型的事件驱动,它和 libevent 网络库的事件驱动一样,都是依托 i/o 多路复用技术支撑起来的。

利用 i/o 多路复用技术,监听感兴趣的文件 i/o 事件,例如读事件,写事件等,同时也要维护一个以文件描述符为主键,数据为某个预设函数的事件表,这里其实就是一个数组或者链表 。当事件触发时,比如某个文件描述符可读,系统会返回文件描述符值,用这个值在事件表中找到相应的数据项,从而实现回调。同样的,定时事件也是可以实现的,因为系统提供的 i/o 多路复用技术中的函数允许我们设定时间值。

上面一段话比较综合,可能需要一些 linux 系统编程和网络编程的基础,但你会看到多数事件驱动程序都是这么实现的(?)。

redis 事件驱动数据结构

redis 事件驱动内部有四个主要的数据结构,分别是:事件循环结构体,文件事件结构体,时间事件结构体和触发事件结构体。

// 文件事件结构体 /* file event structure */ typedef struct aefileevent { int mask; /* one of ae_(readable|writable) */ // 回调函数指针 aefileproc *rfileproc; aefileproc *wfileproc; // clientdata 参数一般是指向 redisclient 的指针 void *clientdata; } aefileevent; // 时间事件结构体 /* time event structure */ typedef struct aetimeevent { long long id; /* time event identifier. */ long when_sec; /* seconds */ long when_ms; /* milliseconds */ // 定时回调函数指针 aetimeproc *timeproc; // 定时事件清理函数,当删除定时事件的时候会被调用 aeeventfinalizerproc *finalizerproc; // clientdata 参数一般是指向 redisclient 的指针 void *clientdata; // 定时事件表采用链表来维护 struct aetimeevent *next; } aetimeevent; // 触发事件 /* a fired event */ typedef struct aefiredevent { int fd; int mask; } aefiredevent; // 事件循环结构体 /* state of an event based program */ typedef struct aeeventloop { int maxfd; /* highest file descriptor currently registered */ int setsize; /* max number of file descriptors tracked */ // 记录最大的定时事件 id + 1 long long timeeventnextid; // 用于系统时间的矫正 time_t lasttime; /* used to detect system clock skew */ // i/o 事件表 aefileevent *events; /* registered events */ // 被触发的事件 aefiredevent *fired; /* fired events */ // 定时事件表 aetimeevent *timeeventhead; // 事件循环结束标识 int stop; // 对于不同的 i/o 多路复用技术,有不同的数据,详见各自实现 void *apidata; /* this is used for polling api specific data */ // 新的循环前需要执行的操作 aebeforesleepproc *beforesleep; } aeeventloop;

上面的数据结构能给我们很好的提示:事件循环结构体维护 i/o 事件表,定时事件表和触发事件表。

事件循环中心

redis 的主函数中调用 initserver() 函数从而初始化事件循环中心(eventloop),它的主要工作是在 aecreateeventloop() 中完成的。

aeeventloop *aecreateeventloop(int setsize) { aeeventloop *eventloop; int i; // 分配空间 if ((eventloop = zmalloc(sizeof(*eventloop))) == null) goto err; // 分配文件事件结构体空间 eventloop->events = zmalloc(sizeof(aefileevent)*setsize); // 分配已触发事件结构体空间 eventloop->fired = zmalloc(sizeof(aefiredevent)*setsize); if (eventloop->events == null || eventloop->fired == null) goto err; eventloop->setsize = setsize; eventloop->lasttime = time(null); // 时间事件链表头 eventloop->timeeventhead = null; // 后续提到 eventloop->timeeventnextid = 0; eventloop->stop = 0; eventloop->maxfd = -1; // 进入事件循环前需要执行的操作,此项会在 redis main() 函数中设置 eventloop->beforesleep = null; // 在这里,aeapicreate() 函数对于每个 io 多路复用模型的实现都有不同,具体参见源代码,因为每种 io 多路复用模型的初始化都不同 if (aeapicreate(eventloop) == -1) goto err; /* events with mask == ae_none are not set. so let's initialize the * vector with it. */ // 初始化事件类型掩码为无事件状态 for (i = 0; i < setsize; i++) eventloop->events[i].mask = ae_none; return eventloop; err: if (eventloop) { zfree(eventloop->events); zfree(eventloop->fired); zfree(eventloop); } return null; }

有上面初始化工作只是完成了一个空空的事件中心而已。要想驱动事件循环,还需要下面的工作。

事件注册详解

文件 i/o 事件注册主要操作在 aecreatefileevent() 中完成。aecreatefileevent() 会根据文件描述符的数值大小在事件循环结构体的 i/o 事件表中取一个数据空间,利用系统提供的 i/o 多路复用技术监听感兴趣的 i/o 事件,并设置回调函数。

int aecreatefileevent(aeeventloop *eventloop, int fd, int mask, aefileproc *proc, void *clientdata) { if (fd >= eventloop->setsize) { errno = erange; return ae_err; } // 在 i/o 事件表中选择一个空间 aefileevent *fe = &eventloop->events[fd]; // aeapiaddevent() 只在此函数中调用,对于不同 io 多路复用实现,会有所不同 if (aeapiaddevent(eventloop, fd, mask) == -1) return ae_err; fe->mask |= mask; // 设置回调函数 if (mask & ae_readable) fe->rfileproc = proc; if (mask & ae_writable) fe->wfileproc = proc; fe->clientdata = clientdata; if (fd > eventloop->maxfd) eventloop->maxfd = fd; return ae_ok; }

对于不同版本的 i/o 多路复用,比如 epoll,select,kqueue 等,redis 有各自的版本,但接口统一,譬如 aeapiaddevent()。

之于定时事件,在事件循环结构体中用链表来维护。定时事件操作在 aecreatetimeevent() 中完成:分配定时事件结构体,设置触发时间和回调函数,插入到定时事件表中。

long long aecreatetimeevent(aeeventloop *eventloop, long long milliseconds, aetimeproc *proc, void *clientdata, aeeventfinalizerproc *finalizerproc) { /* 自增 timeeventnextid 会在处理执行定时事件时会用到,用于防止出现死循环。 如果超过了最大 id,则跳过这个定时事件,为的是避免死循环,即: 如果事件一执行的时候注册了事件二,事件一执行完毕后事件二得到执行,紧接着如果事件一有得到执行就会成为循环,因此维护了 timeeventnextid 。*/ long long id = eventloop->timeeventnextid++; aetimeevent *te; // 分配空间 te = zmalloc(sizeof(*te)); if (te == null) return ae_err; // 填充时间事件结构体 te->id = id; // 计算超时时间 aeaddmillisecondstonow(milliseconds,&te->when_sec,&te->when_ms); // proc == servercorn te->timeproc = proc; te->finalizerproc = finalizerproc; te->clientdata = clientdata; // 头插法 te->next = eventloop->timeeventhead; eventloop->timeeventhead = te; return id; } 准备监听工作

initserver() 中调用了 aecreateeventloop() 完成了事件中心的初始化,initserver() 还做了监听的准备。

/* open the tcp listening socket for the user commands. */ // listentoport() 中有调用 listen() if (server.port != 0 && listentoport(server.port,server.ipfd,&server.ipfd_count) == redis_err) exit(1); // unix 域套接字 /* open the listening unix domain socket. */ if (server.unixsocket != null) { unlink(server.unixsocket); /* don't care if this fails */ server.sofd = anetunixserver(server.neterr,server.unixsocket,server.unixsocketperm); if (server.sofd == anet_err) { redislog(redis_warning, "opening socket: %s", server.neterr); exit(1); } }

从上面可以看出,redis 提供了 tcp 和 unix 域套接字两种工作方式。以 tcp 工作方式为例,listenport() 创建绑定了套接字并启动了监听。

为监听套接字注册事件

在进入事件循环前还需要做一些准备工作。紧接着,initserver() 为所有的监听套接字注册了读事件,响应函数为 accepttcphandler() 或者 acceptunixhandler()。

// 创建接收 tcp 或者 unix 域套接字的事件处理 // tcp /* create an event handler for accepting new connections in tcp and unix * domain sockets. */ for (j = 0; j < server.ipfd_count; j++) { // accepttcphandler() tcp 连接接受处理函数 if (aecreatefileevent(server.el, server.ipfd[j], ae_readable, accepttcphandler,null) == ae_err) { redispanic( "unrecoverable error creating server.ipfd file event."); } } // unix 域套接字 if (server.sofd > 0 && aecreatefileevent(server.el,server.sofd,ae_readable, acceptunixhandler,null) == ae_err) redispanic("unrecoverable error creating server.sofd file event.");

来看看accepttcphandler() 做了什么:

// 用于 tcp 接收请求的处理函数 void accepttcphandler(aeeventloop *el, int fd, void *privdata, int mask) { int cport, cfd; char cip[redis_ip_str_len]; redis_notused(el); redis_notused(mask); redis_notused(privdata); // 接收客户端请求 cfd = anettcpaccept(server.neterr, fd, cip, sizeof(cip), &cport); // 出错 if (cfd == ae_err) { redislog(redis_warning,"accepting client connection: %s", server.neterr); return; } // 记录 redislog(redis_verbose,"accepted %s:%d", cip, cport); // 真正有意思的地方 acceptcommonhandler(cfd,0); }

接收套接字与客户端建立连接后,调用 acceptcommonhandler()。acceptcommonhandler() 主要工作就是:

  • 建立并保存服务端与客户端的连接信息,这些信息保存在一个 struct redisclient 结构体中;
  • 为与客户端连接的套接字注册读事件,相应的回调函数为 readqueryfromclient(),readqueryfromclient() 作用是从套接字读取数据,执行相应操作并回复客户端。
  • redis 事件循环

    以上做好了准备工作,可以进入事件循环。跳出 initserver() 回到 main() 中,main() 会调用 aemain()。进入事件循环发生在 aeprocessevents() 中:

  • 根据定时事件表计算需要等待的最短时间;
  • 调用 redis api aeapipoll() 进入监听轮询,如果没有事件发生就会进入睡眠状态,其实就是 i/o 多路复用 select() epoll() 等的调用;
  • 有事件发生会被唤醒,处理已触发的 i/o 事件和定时事件。
  • void aemain(aeeventloop *eventloop) { eventloop->stop = 0; while (!eventloop->stop) { // 进入事件循环可能会进入睡眠状态。在睡眠之前,执行预设置的函数 aesetbeforesleepproc()。 if (eventloop->beforesleep != null) eventloop->beforesleep(eventloop); // ae_all_events 表示处理所有的事件 aeprocessevents(eventloop, ae_all_events); } } // 先处理定时事件,然后处理套接字事件 int aeprocessevents(aeeventloop *eventloop, int flags) { int processed = 0, numevents; /* nothing to do? return asap */ if (!(flags & ae_time_events) && !(flags & ae_file_events)) return 0; /* note that we want call select() even if there are no * file events to process as long as we want to process time * events, in order to sleep until the next time event is ready * to fire. */ if (eventloop->maxfd != -1 || ((flags & ae_time_events) && !(flags & ae_dont_wait))) { int j; aetimeevent *shortest = null; // tvp 会在 io 多路复用的函数调用中用到,表示超时时间 struct timeval tv, *tvp; // 得到最短将来会发生的定时事件 if (flags & ae_time_events && !(flags & ae_dont_wait)) shortest = aesearchnearesttimer(eventloop); // 计算睡眠的最短时间 if (shortest) { // 存在定时事件 long now_sec, now_ms; /* calculate the time missing for the nearest * timer to fire. */ // 得到当前时间 aegettime(&now_sec, &now_ms); tvp = &tv; tvp->tv_sec = shortest->when_sec - now_sec; if (shortest->when_ms < now_ms) { // 需要借位 // 减法中的借位,毫秒向秒借位 tvp->tv_usec = ((shortest->when_ms+1000) - now_ms)*1000; tvp->tv_sec --; } else { // 不需要借位,直接减 tvp->tv_usec = (shortest->when_ms - now_ms)*1000; } // 当前系统时间已经超过定时事件设定的时间 if (tvp->tv_sec < 0) tvp->tv_sec = 0; if (tvp->tv_usec < 0) tvp->tv_usec = 0; } else { /* if we have to check for events but need to return * asap because of ae_dont_wait we need to set the timeout * to zero */ // 如果没有定时事件,见机行事 if (flags & ae_dont_wait) { tv.tv_sec = tv.tv_usec = 0; tvp = &tv; } else { /* otherwise we can block */ tvp = null; /* wait forever */ } } // 调用 io 多路复用函数阻塞监听 numevents = aeapipoll(eventloop, tvp); // 处理已经触发的事件 for (j = 0; j < numevents; j++) { // 找到 i/o 事件表中存储的数据 aefileevent *fe = &eventloop->events[eventloop->fired[j].fd]; int mask = eventloop->fired[j].mask; int fd = eventloop->fired[j].fd; int rfired = 0; /* note the fe->mask & mask & ... code: maybe an already processed * event removed an element that fired and we still didn't * processed, so we check if the event is still valid. */ // 读事件 if (fe->mask & mask & ae_readable) { rfired = 1; fe->rfileproc(eventloop,fd,fe->clientdata,mask); } // 写事件 if (fe->mask & mask & ae_writable) { if (!rfired || fe->wfileproc != fe->rfileproc) fe->wfileproc(eventloop,fd,fe->clientdata,mask); } processed++; } } // 处理定时事件 /* check time events */ if (flags & ae_time_events) processed += processtimeevents(eventloop); return processed; /* return the number of processed file/time events */ } 事件触发

    这里以 select 版本的 redis api 实现作为讲解,aeapipoll() 调用了 select() 进入了监听轮询。aeapipoll() 的 tvp 参数是最小等待时间,它会被预先计算出来,它主要完成:

  • 拷贝读写的 fdset。select() 的调用会破坏传入的 fdset,实际上有两份 fdset,一份作为备份,另一份用作调用。每次调用 select() 之前都从备份中直接拷贝一份;
  • 调用 select();
  • 被唤醒后,检查 fdset 中的每一个文件描述符,并将可读或者可写的描述符记录到触发表当中。
  • 接下来的操作便是执行相应的回调函数,代码在上一段中已经贴出:先处理 i/o 事件,再处理定时事件。

    static int aeapipoll(aeeventloop *eventloop, struct timeval *tvp) { aeapistate *state = eventloop->apidata; int retval, j, numevents = 0; /* 真有意思,在 aeapistate 结构中: typedef struct aeapistate { fd_set rfds, wfds; fd_set _rfds, _wfds; } aeapistate; 在调用 select() 的时候传入的是 _rfds 和 _wfds,所有监听的数据在 rfds 和 wfds 中。 在下次需要调用 selec() 的时候,会将 rfds 和 wfds 中的数据拷贝进 _rfds 和 _wfds 中。*/ memcpy(&state->_rfds,&state->rfds,sizeof(fd_set)); memcpy(&state->_wfds,&state->wfds,sizeof(fd_set)); retval = select(eventloop->maxfd+1, &state->_rfds,&state->_wfds,null,tvp); if (retval > 0) { // 轮询 for (j = 0; j maxfd; j++) { int mask = 0; aefileevent *fe = &eventloop->events[j]; if (fe->mask == ae_none) continue; if (fe->mask & ae_readable && fd_isset(j,&state->_rfds)) mask |= ae_readable; if (fe->mask & ae_writable && fd_isset(j,&state->_wfds)) mask |= ae_writable; // 添加到触发事件表中 eventloop->fired[numevents].fd = j; eventloop->fired[numevents].mask = mask; numevents++; } } return numevents; } 总结

    redis 的事件驱动总结如下:

  • 初始化事件循环结构体
  • 注册监听套接字的读事件
  • 注册定时事件
  • 进入事件循环
  • 如果监听套接字变为可读,会接收客户端请求,并为对应的套接字注册读事件
  • 如果与客户端连接的套接字变为可读,执行相应的操作
  • 后续分享更多内容。

    —-

    捣乱 2014-3-9

    http://daoluan.net



    【说明】本文章由站长整理发布,文章内容不代表本站观点,如文中有侵权行为,请与本站客服联系(QQ:)!