工作队列-模拟消费者处理时间

工作队列-模拟消费者处理时间

前面的例子展示了如何在指定的消息队列发送和接收消息。

现在我们创建一个工作队列(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)

我们进行如下的测试:

  1. 先打开生产者,创建消息
  2. 打开一个消费者,再打开第二个消费者
  3. 关闭生产者和消费者
  4. 先打开两个消费者
  5. 再打开一个生产者
  6. 关闭消费者、生产者
  7. 先打开两个消费者
  8. 打开生产者
  9. 关闭其中一个消费者,观察消息执行情况

发表回复

您的邮箱地址不会被公开。 必填项已用 * 标注