【发布时间】:2015-06-15 17:02:42
【问题描述】:
我有以下设置: 有一个客户端、多个工作人员和一个接收器。 工作人员通过 ZeroMQ 消息接收来自客户端的作业请求。他们处理输入,并将答案发送到另一个进程(接收器)。处理一条消息大约需要 1 毫秒,我们需要处理大约 50,000 条消息/秒 - 这意味着我们需要 50 多个工作人员来处理负载。
我尝试了一个简单的设置,其中客户端创建一个 ZeroMQ PUSH 套接字,所有工作人员都连接到该套接字(通过 PULL)。类似地,sink 会创建一个 PULL 套接字,所有工作人员通过 PUSH 套接字连接到该套接字。
IIUC,ZeroMQ 使用“循环”将消息发送给工作人员 - 每次另一个工作人员获得工作。这种设置似乎可以在大约 10 个工作人员(和适当的负载)的情况下足够有效地工作。但是,当进一步增加工作人员的数量和负载时,这会很快中断并且系统开始累积延迟。
我知道有几种模式可以解决负载平衡问题,但是它们面向多个客户端并且需要在其间使用路由器,这意味着额外的代码 + cpu 周期。问题是:
1) 在单个客户端、多个工作器、单个接收器的情况下,最好的模式是什么?
2) 是否可以在客户端和工作人员之间没有路由器的情况下通过在客户端进行路由来执行此操作?
3) 应该使用什么样的 ZeroMQ 套接字?
谢谢!
编辑: 添加代码。
客户:
void *context = zmq_ctx_new ();
// Socket to send messages on
void *sender = zmq_socket (context, ZMQ_PUSH);
zmq_bind (sender, "tcp://*:5557");
// Socket to send start of batch message on
void *sink = zmq_socket (context, ZMQ_PUSH);
zmq_connect (sink, "tcp://localhost:5558");
printf ("Press Enter when the workers are ready: ");
getchar ();
printf ("Sending tasks to workers\n");
// The first message is "0" and signals start of batch
s_send (sink, "0");
unsigned long i;
const int nmsgs = atoi(argv[1]);
const int nmsgs_sec = atoi(argv[2]);
const int buff_size = 1024; // 1KB msgs
unsigned long t, t_start;
t_start = timestamp();
for (i = 0; i < nmsgs; i++) {
t = timestamp();
// Pace the sending according to nmsgs_sec
while( i * 1000000 / (t+1-t_start) > nmsgs_sec) {
// busy wait
t = timestamp();
}
char buffer [buff_size];
// Write current timestamp in the packet beginning
sprintf (buffer, "%lu", t);
zmq_send (sender, buffer, buff_size, 0);
}
printf("Total time: %lu ms Planned time: %d ms\n", (timestamp() - t_start)/1000, nmsgs * 1000 / nmsgs_sec);
zmq_close (sink);
zmq_close (sender);
zmq_ctx_destroy (context);
工人:
// Socket to receive messages on
void *context = zmq_ctx_new ();
void *receiver = zmq_socket (context, ZMQ_PULL);
zmq_connect (receiver, receiver_addr);
// Socket to send messages to
void *sender = zmq_socket (context, ZMQ_PUSH);
zmq_connect (sender, sender_addr);
// Process tasks forever
const int buff_size = 1024;
char buffer[buff_size];
while (1) {
zmq_recv (receiver, buffer, buff_size, 0);
s_send (sender, buffer);
}
zmq_close (receiver);
zmq_close (sender);
zmq_ctx_destroy (context);
水槽:
// Prepare our context and socket
void *context = zmq_ctx_new ();
void *receiver = zmq_socket (context, ZMQ_PULL);
zmq_bind (receiver, "tcp://*:5558");
// Wait for start of batch
char *string = s_recv (receiver);
free (string);
unsigned long t1;
unsigned long maxdt = 0;
unsigned long sumdt = 0;
int task_nbr;
int nmsgs = atoi(argv[1]);
printf("nmsgs = %d\n", nmsgs);
for (task_nbr = 0; task_nbr < nmsgs; task_nbr++) {
char *string = s_recv (receiver);
t1 = timestamp();
unsigned long t0 = atoll(string);
free (string);
unsigned long dt = t1-t0;
maxdt = (maxdt > dt ? maxdt : dt);
sumdt += dt;
if(task_nbr % 10000 == 0) {
printf("%d %lu\n", task_nbr, dt);
}
}
printf("Average time: %lu usec\tMax time: %lu usec\n", sumdt/nmsgs, maxdt);
zmq_close (receiver);
zmq_ctx_destroy (context);
【问题讨论】:
标签: sockets client-server load-balancing zeromq