LongLong's Blog

分享IT技术,分享生活感悟,热爱摄影,热爱航天。

用Redis做Gearman的持久化队列

Gearman很久之前就开始支持持久化队列,但用sqlite进行持久化性能实在是着急,之后支持了MySQL后才感觉靠谱一些,但对于高并发的情况下需要反复的写入和删除数据——消息到达需要写入消息,处理完成后删除消息,对MySQL的性能有着非常高的要求。而在Gearman的源代码中已经提供了将Redis作为持久化队列的代码,但相关的代码有较多的bug,感觉目前只是个隐藏功能,如果不对代码进行处理是没法正常使用的。

1. 修复代码中的bug

其中主要修复了默认只能够连接本地Reids、启动时不能够正确获取队列中的全部任务、不能正确解析获取到的任务标识和一些段错误,增加了断线重连的机制。感觉官方根本就没有认真的写这个模块,更没有好好的进行测试。补丁的代码如下,针对最新的1.1.12版本。希望官方之后能修复自己的一大堆bug。

diff -ruNa gearmand-1.1.12/libgearman-server/plugins/queue/redis/queue.cc gearmand-1.1.12.patch/libgearman-server/plugins/queue/redis/queue.cc
--- gearmand-1.1.12/libgearman-server/plugins/queue/redis/queue.cc      2014-02-12 08:05:28.000000000 +0800
+++ gearmand-1.1.12.patch/libgearman-server/plugins/queue/redis/queue.cc        2014-06-18 10:57:02.147575821 +0800
@@ -85,6 +85,7 @@
   ~Hiredis();

   gearmand_error_t initialize();
+  bool init_redis();

   redisContext* redis()
   {
@@ -113,10 +114,15 @@
 {
 }

+bool Hiredis::init_redis() {
+  int service_port= atoi(service.c_str());
+  _redis = redisConnect(server.c_str(), service_port);
+  return _redis != NULL;
+}
+
 gearmand_error_t Hiredis::initialize()
 {
-  int service_port= atoi(service.c_str());
-  if ((_redis= redisConnect("127.0.0.1", service_port)) == NULL)
+  if (!init_redis())
   {
     return gearmand_gerror("Could not connect to redis server", GEARMAND_QUEUE_ERROR);
   }
@@ -148,7 +154,7 @@
                         const char *function_name,
                         size_t function_name_size)
 {
-  key.resize(function_name_size +unique_size +GEARMAND_QUEUE_GEARMAND_DEFAULT_PREFIX_SIZE +4);
+  key.resize(function_name_size +unique_size +GEARMAND_QUEUE_GEARMAND_DEFAULT_PREFIX_SIZE +2);
   int key_size= snprintf(&key[0], key.size(), GEARMAND_KEY_LITERAL,
                          GEARMAND_QUEUE_GEARMAND_DEFAULT_PREFIX,
                          (int)function_name_size, function_name,
@@ -202,13 +208,19 @@
   build_key(key, unique, unique_size, function_name, function_name_size);
   gearmand_log_debug(GEARMAN_DEFAULT_LOG_PARAM, "hires key: %u", (uint32_t)key.size());

-  redisReply *reply= (redisReply*)redisCommand(queue->redis(), "SET %b %b", &key[0], key.size(), data, data_size);
+  redisReply *reply= (redisReply*)redisCommand(queue->redis(), "SET %s %b", &key[0], data, data_size);
   gearmand_log_debug(GEARMAN_DEFAULT_LOG_PARAM, "got reply");
   if (reply == NULL)
   {
-    return gearmand_log_gerror(GEARMAN_DEFAULT_LOG_PARAM, GEARMAND_QUEUE_ERROR, "failed to insert '%.*s' into redis", key.size(), &key[0]);
+    if (!queue->init_redis())
+    {
+      return gearmand_log_gerror(GEARMAN_DEFAULT_LOG_PARAM, GEARMAND_QUEUE_ERROR, "failed to insert '%.*s' into redis", key.size(), &key[0]);
+    }
+  } 
+  else 
+  { 
+    freeReplyObject(reply);
   }
-  freeReplyObject(reply);

   return GEARMAND_SUCCESS;
 }
@@ -231,12 +243,16 @@
   std::vector<char> key;
   build_key(key, unique, unique_size, function_name, function_name_size);

-  redisReply *reply= (redisReply*)redisCommand(queue->redis(), "DEL %b", &key[0], key.size());
+  redisReply *reply= (redisReply*)redisCommand(queue->redis(), "DEL %s", &key[0]);
   if (reply == NULL)
   {
-    return GEARMAND_QUEUE_ERROR;
+    if (!queue->init_redis()) {
+      return GEARMAND_QUEUE_ERROR;
+    }
+  }
+  else {
+    freeReplyObject(reply);
   }
-  freeReplyObject(reply);

   return GEARMAND_SUCCESS;
 }
@@ -252,7 +268,7 @@
   
   gearmand_info("hiredis replay start");

-  redisReply *reply= (redisReply*)redisCommand(queue->redis(), "KEYS %s", GEARMAND_QUEUE_GEARMAND_DEFAULT_PREFIX);
+  redisReply *reply= (redisReply*)redisCommand(queue->redis(), "KEYS %s*", GEARMAND_QUEUE_GEARMAND_DEFAULT_PREFIX);
   if (reply == NULL)
   {
     return gearmand_gerror("Failed to call KEYS during QUEUE replay", GEARMAND_QUEUE_ERROR);
@@ -265,9 +281,7 @@
     char unique[GEARMAN_MAX_UNIQUE_SIZE];

     char fmt_str[100] = "";    
-    int fmt_str_length= snprintf(fmt_str, sizeof(fmt_str), "%%%ds-%%%ds-%%%ds",
-                                 int(GEARMAND_QUEUE_GEARMAND_DEFAULT_PREFIX_SIZE),
-                                 int(GEARMAN_FUNCTION_MAX_SIZE),
+    int fmt_str_length= snprintf(fmt_str, sizeof(fmt_str), "%%[^-]-%%[^-]-%%%ds",
                                  int(GEARMAN_MAX_UNIQUE_SIZE));
     if (fmt_str_length <= 0 or size_t(fmt_str_length) >= sizeof(fmt_str))
     {
@@ -293,7 +307,7 @@
     (void)(add_fn)(server, add_context,
                    unique, strlen(unique),
                    function_name, strlen(function_name),
-                   get_reply->str, get_reply->len,
+                   strndup(get_reply->str, get_reply->len), get_reply->len,
                    GEARMAN_JOB_PRIORITY_NORMAL, 0);
     freeReplyObject(get_reply);
   }

2.模块的安装

如果安装了hiredis库,则在安装Gearman时就会自动加载Redis持久化队列模块,hireids的代码在Redis代码的deps/hiredis目录下,直接编译安装即可。正确的安装了Redis持久化队列模块后,运行gearmand有以下的提示,则证明模块已经正常被安装和加载。

redis:
--redis-server arg    Redis server
--redis-port arg      Redis server port/service

3.运行的效果

启动一个gearmand并连接相应的Redis实例,如下启动一个监听4730端口的gearmand并连接14730端口的Redis,这里端口的设计方法是方便启动多个gearmand时批量进行管理,默认对应Redis的端口是gearmand端口加10000。gearmand在启动时就会和Redis进行连接,如果连接失败则启动会失败。

/usr/local/sbin/gearmand -d -p 4730 -u nobody -P /var/run/gearmand/gearmand.4730.pid -t 4 -j 3 -l /var/log/gearmand/gearmand.4730 -q redis --redis-server 127.0.0.1 --redis-port 14730

通过gearman命令行工具向test队列写入一个异步消息,消息的内容为123,然后连接Redis进行查看,最后再使用gearadmin查看队列的状况,重启gearmand再查看队列的状况,最后启动一个Worker获取消息

#写入消息
echo "123" | gearman -f test -b
#连接Redis
/opt/redis/bin/redis-cli  -p 14730
#查看Redis中的内容,发现有一条记录
127.0.0.1:14730> keys *
1) "_gear_-test-205d3b6e-c479-11e4-9ad2-00237d29f08a"
#查看数据,发现是发送过来的123
127.0.0.1:14730> get _gear_-test-205d3b6e-c479-11e4-9ad2-00237d29f08a
"123\n"
127.0.0.1:14730>
#查看队列的状况,发现有一个消息
gearadmin --status
test	1	0	0
.
#重启gearmand
/etc/init.d/gearmand.4730 restart
#再查看队列的状况,发现消息仍然存在
gearadmin --status
test	1	0	0
.
#启动一个worker,发现正常的获得到了发送的消息
gearman -w -f test
123

如此就完成了对持久化正确性的验证,由于Redis是内存的数据库因此相比MySQL性能会有一定的保证,之前测试的结果是相比不开启持久化队列性能大概要下降一半左右。不过为了保证可靠性必然要付出一定的性能代价。

Memcached事件模型分析——线程池与连接的建立

之前一直想看一下Memcached的事件模型的实现方法,最近终于有时间来研究一下。看到网上相关的文章也比较多,这里只写一下我自己的理解,如果有错误还请指出。代码为了表达主要的思想都进行了大幅度的简化。

1. 总体结构

Memcached总体的事件模型代码在thread.c和memcached.c两个文件中,其中thread.c中主要实现了线程池,memcached.c中主要实现事件响应和连接的处理。事件驱动部分使用了比较常见的libevent。

主线程主要负责响应外部的连接,而工作线程池则和主线程之间通过管道进行通信——每一个工作线程注册一个对管道的事件响应,主线程在接收到外部连接后,轮询工作线程池找到下一个工作线程,并通过向其对应的管道中写入一个字节的方式触发事件响应并完成处理操作。主线程和工作线程之间通过每个工作线程自身的一个队列进行数据的传递——主线程接收到连接后,将连接相关的信息封装为一个对象推入队列,并通过管道激活对应的工作线程,工作线程从自身的队列中取出相应的数据进行处理。

2. 线程池的实现

如第一节中所述,每个工作线程都包含两个管道的描述符,一个事件结构和一个队列。工作线程的结构如下

class conn_queue {
private:
    std::queue<int>* queue;
    pthread_mutex_t lock;
public:
    conn_queue();
    int pop();
    void push(int fd);
    int size();
    ~conn_queue();
};

typedef struct {
	pthread_t thread_id;       
	struct event_base *base;   //libevent句柄
	struct event notify_event; //通知事件结构
	int notify_receive_fd;     //触发工作线程的描述符
	int notify_send_fd;     
	conn_queue* queue; 	   	   //连接队列,使用一个封装了自动加锁的队列
} LIBEVENT_THREAD;

每个工作线程的初始化过程如下

//工作线程池指针
static LIBEVENT_THREAD *threads;

//初始化线程的事件和连接队列
static void setup_thread(LIBEVENT_THREAD *me) {
	//初始化事件
    me->base = event_init();

	//增加对notify_receive_fd描述符的事件响应
    event_set(&me->notify_event, me->notify_receive_fd,
              EV_READ | EV_PERSIST, thread_libevent_process, me);
    event_base_set(me->base, &me->notify_event);

    if (event_add(&me->notify_event, 0) == -1) {
        fprintf(stderr, "Can't monitor libevent notify pipe\n");
        exit(1);
    }

	//初始化连接队列
	me->queue = new conn_queue();
}

void memcached_thread_init(int nthreads, struct event_base *main_base) {
	pthread_mutex_init(&init_lock, NULL);
    pthread_cond_init(&init_cond, NULL);
	
	threads = calloc(nthreads, sizeof(LIBEVENT_THREAD));
	//初始化每个工作线程
    for (int i = 0; i < nthreads; i++) {
		//开启管道
        int fds[2];
		pipe(fds);

        threads[i].notify_receive_fd = fds[0];
        threads[i].notify_send_fd = fds[1];

		//初始化工作线程的事件结构和连接队列
        setup_thread(&threads[i]);
    }

	//启动线程 事件循环就在线程的回调函数worker_libevent中启动
    for (int i = 0; i < nthreads; i++) {
        create_worker(worker_libevent, &threads[i]);
    }

	//等待全部工作线程启动完毕
	pthread_mutex_lock(&init_lock);
    wait_for_thread_registration(nthreads);
    pthread_mutex_unlock(&init_lock);
}

3. 外部连接的处理和线程调度

对于外部连接的处理相关的逻辑在主线程中,大致和libevent的单线程用法类似——响应监听描述符的事件,这部分的逻辑在memcached的main函数中,原始的代码中对监听描述符和管道描述符的事件响应使用了同一个处理函数event_handler,而是通过连接的不同状态在drive_machine中进行区分处理

static struct event_base *main_base;
static int last_thread = -1;

void base_event_handler(int sock, short event, void* arg) {
	//接收连接
    struct sockaddr_in cli_addr;
    int newfd;
    socklen_t sin_size;
    sin_size = sizeof(struct sockaddr_in);
    newfd = accept(sock, (struct sockaddr*)&cli_addr, &sin_size);

	//选择线程,使用轮询的方式进行选择
    int tid = (last_thread + 1) % THREAD_NUM;
    LIBEVENT_THREAD* thread = threads + tid;
    last_thread = tid;

	//将连接描述符推入队列
	thread->queue->push(newfd);

	//向管道中写入一个空字符激活触发工作线程的事件响应 
    write(thread->notify_send_fd, " ", 1); 
}

int main() {
    main_base = event_init();
	//初始化线程池
    thread_init(THREAD_NUM, main_base);

	//开启对地址和端口的监听
    int sfd = server_socket("127.0.0.1", 11212);
	//增加对应监听描述符的事件响应
    struct event listen_ev;
    event_set(&listen_ev, sfd, EV_READ | EV_PERSIST, base_event_handler, NULL);
    event_base_set(main_base, &listen_ev);
    event_add(&listen_ev, NULL);

	//启动事件循环    
    event_base_loop(main_base, 0); 
    return 0;
}

4. 总结

Memcached的事件模型是使用libevent实现的多线程TCP类服务器的比较经典的实例,在开发一些轻量级的服务组件时非常具有参考意义,以上只是简单的分析了Memcached的线程池和连接处理的机制,后续的数据的读写和连接的保持和关闭部分还未涉及,之后会对剩下的部分加以分析和研究,最终目的是为了得到一个编写多线程TCP类服务组件的代码框架。