RabbitMQ Work Queue模式讲解

RabbitMQ Work Queue模式讲解

一.W/Q 模式的结构

有一个负责发送任务,多个执行单元接受任务后进行处理。该模式中,没个消费者接受的任务量是均衡的。不会根据具体任务的处理时间进行更精确的分配

发布者代码 var message = GetMessage(args);var body = Encoding.UTF8.GetBytes(message); var properties = channel.CreateBasicProperties(); properties.Persistent = true;

channel.BasicPublish(exchange: “”, routingKey: “task_queue”, basicProperties: properties, body: body);

接收者代码 var consumer = new EventingBasicConsumer(channel); consumer.Received += (model, ea) => { var body = ea.Body.ToArray(); var message = Encoding.UTF8.GetString(body); Console.WriteLine(" [x] Received {0}", message);

int dots = message.Split(.).Length - 1;
Thread.Sleep(dots * 1000);

Console.WriteLine(" [x] Done");

}; channel.BasicConsume(queue: “task_queue”, autoAck: true, consumer: consumer);

发送者的Console

shell 3cd NewTask

dotnet run “First message.” dotnet run “Second message…” dotnet run “Third message…” dotnet run “Fourth message…” dotnet run “Fifth message…”

Worker1 的Console

shell 1# => [*] Waiting for messages. To exit press CTRL+C# => [x] Received ‘First message.’# => [x] Received ‘Third message…’# => [x] Received ‘Fifth message…’

Work2的Console

shell 2# => [*] Waiting for messages. To exit press CTRL+C# => [x] Received ‘Second message…’# => [x] Received ‘Fourth message…’

以上的,在两个接受者收到消息后,消息队列就自己删除了,不会等待接受者对于任务处理的反馈。

二.消息的确认

int dots = message.Split(.).Length - 1;
Thread.Sleep(dots * 1000);

Console.WriteLine(" [x] Done");

// Note: it is possible to access the channel via
//       ((EventingBasicConsumer)sender).Model here
channel.BasicAck(deliveryTag: ea.DeliveryTag, multiple: false);

}; channel.BasicConsume(queue: “task_queue”, autoAck: false, consumer: consumer);

三.消息的持续

以上我们解决了worker崩溃造成消息的删除。那么 如果Server崩溃造成消息的删除,要根据下面的操作进行

channel.QueueDeclare(queue: “hello”, durable: true, exclusive: false, autoDelete: false, arguments: null);

var properties = channel.CreateBasicProperties(); properties.Persistent = true;

ior we can use the BasicQos method with the prefetchCount = 1 setting.

channel.BasicQos(prefetchSize: 0, prefetchCount: 1, global: false);

经验分享 程序员 微信小程序 职场和发展