工作队列-模拟消费者处理时间
前面的例子展示了如何在指定的消息队列发送和接收消息。
现在我们创建一个工作队列(work queue)来将一些耗时的任务分发给多个工作者(workers):

工作队列(work queues,又称任务队列Task Queues)的主要思想是为了避免立即执行并等待一些占用大量资源、时间的操作完成。而是把任务(Task)当作消息发送到队列中,稍后处理。一个运行在后台的工作者(worker)进程就会取出任务然后处理。当运行多个工作者(workers)时,任务会在它们之间共享。
这个在网络应用中非常有用,它可以在短暂的HTTP请求中处理一些复杂的任务。在一些实时性要求不太高的地方,我们可以处理完主要操作之后,以消息的方式来处理其他不紧要的操作,比如写日志等等。
使用工作队列的一个好处就是它能够并行地处理队列。如果堆积了很多任务,我们只需要添加更多的工作者(workers)就可以了,扩展很简单。
现在,我们先启动两个接收端,等待接收消息,然后启动一个发送端开始发送消息。
为了模拟耗时情况,我们给每条发送的文本增加一个数字,这个数字表示耗时时间。我们使用sleep来睡眠。
比如:
hello1:3
hello2:5
hello3:2
分别表示耗时3秒、5秒、2秒。然后我们在接收程序里面通过Thread.Sleep来实现。
默认,RabbitMQ会将每个消息按照顺序依次分发给下一个消费者。所以每个消费者接收到的消息个数大致是平均的。这种消息分发的方式称之为轮询(round-robin)。
发送者
class Program
{
private static Random random = new Random(100);
static void Main(string[] args)
{
var factory = new ConnectionFactory();
factory.HostName = "localhost";//RabbitMQ服务在本地运行
factory.UserName = "guest";//用户名
factory.Password = "guest";//密码
using (var connection = factory.CreateConnection())
{
using (var channel = connection.CreateModel())
{
channel.QueueDeclare("hello", false, false, false, null);//创建一个名称为hello的消息队列
for (int i = 0; i < 8; i++)
{
var properties = channel.CreateBasicProperties();
properties.DeliveryMode = 2;//1非持久化 2持久化
string message = "Hello:" + random.Next(2, 10); //传递的消息内容
var body = Encoding.UTF8.GetBytes(message);
channel.BasicPublish("", "hello", properties, body);
Console.WriteLine("已发送: {0}", message);
}
Console.ReadLine();
}
}
}
}Code language: JavaScript (javascript)
消费者
class Program
{
static void Main(string[] args)
{
var factory = new ConnectionFactory();
factory.HostName = "localhost";
factory.UserName = "guest";
factory.Password = "guest";
using (var connection = factory.CreateConnection())
{
using (var channel = connection.CreateModel())
{
channel.QueueDeclare("hello", false, false, false, null);
var consumer = new QueueingBasicConsumer(channel);
channel.BasicConsume("hello", true, consumer);
while (true)
{
var ea = (BasicDeliverEventArgs)consumer.Queue.Dequeue();
var body = ea.Body;
var message = Encoding.UTF8.GetString(body);
var time = Convert.ToInt32(message.Split(':')[1]);
Thread.Sleep(time * 1000);
Console.WriteLine("Received {0}", message);
Console.WriteLine("Done");
}
}
}
}
}Code language: JavaScript (javascript)
我们进行如下的测试:
- 先打开生产者,创建消息
- 打开一个消费者,再打开第二个消费者
- 关闭生产者和消费者
- 先打开两个消费者
- 再打开一个生产者
- 关闭消费者、生产者
- 先打开两个消费者
- 打开生产者
- 关闭其中一个消费者,观察消息执行情况
Previous: 第一个例子-消费者接收消息
Next: 消息响应-避免消息丢失