公平分发
你可能会注意到,消息的分发可能并没有如我们想要的那样公平分配。比如,对于两个工作者,当奇数个消息的任务比较重,但是偶数个消息任务比较轻时,奇数个工作者始终处于忙碌状态,而偶数个工作者始终处于空闲状态。但是RabbitMQ并不知道这些,它仍然会平均依次的分发消息。
为了改变这一状态,我们可以使用BasicQos方法,设置prefetchCount=1。这样就告诉RabbitMQ不要在同一时间给一个工作者发送多于1个的消息,或者换句话说,在一个工作者还在处理消息,并且没有响应消息之前,不要给他分发新的消息。相反,将这条新的消息发送给下一个不那么忙碌的工作者。
channel.BasicQos(0, 1, false);Code language: JavaScript (javascript)
消息确认
在一些场合,如转账、付费时每一条消息都必须保证成功的被处理。AMQP是金融级的消息队列协议,有很高的可靠性,这里介绍在使用RabbitMQ时怎么保证消息被成功处理的。消息确认可以分为两种:一种是生产者发送消息到Broke时,Broker给生产者发送确认回执,用于告诉生产者消息已被成功发送到Broker;一种是消费者接收到Broker发送的消息时,消费者给Broker发送确认回执,用于通知消息已成功被消费者接收。
下边分别介绍生产者端和消费者端的消息确认方法。准备条件:使用Web管理工具添加exchange、queue并绑定,bindingKey为“mykey”,如下所示:



生产者端消息确认-tx机制
生产者端的消息确认:当生产者将消息发送给Broker,Broker接收到消息后给生产者发送确认回执。生产者端的消息确认有两种方式:tx机制和Confirm模式。
tx机制可以叫做事务机制,RabbitMQ中有三个与tx机制的方法:txSelect()、txCommit()和txRollback()。 channel.txSelect() 用于将当前channel设置成transaction模式, channel.txCommit() 提交事务, channel.txRollback() 回滚事务。使用tx机制,我们首先要通过txSelect方法开启事务,然后发布消息给broker服务器,如果txCommit提交成功了,则说明消息成功被broker接收了;如果在txCommit执行之前broker异常崩溃或者由于其他原因抛出异常,这个时候我们可以捕获异常,通过txRollback回滚事务。看一个tx机制的简单实现:
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())
{
var properties = channel.CreateBasicProperties();
properties.DeliveryMode = 2;//1非持久化 2持久化
string message = "Hello:" + random.Next(2, 10); //传递的消息内容
var body = Encoding.UTF8.GetBytes(message);
try
{
//开启事务机制
channel.TxSelect();
//发送消息
channel.BasicPublish(exchange: "myexchange",
routingKey: "mykey",
basicProperties: null,
body: body);
//事务提交
channel.TxCommit();
Console.WriteLine($"【{message}】发送到Broker成功!");
}
catch (Exception)
{
Console.WriteLine($"【{message}】发送到Broker失败!");
//回滚事务
channel.TxRollback();
}
}
}
}Code language: JavaScript (javascript)
生产者端消息确认-Confirm模式
C#的RabbitMQ API中,有三个与Confirm相关的方法:ConfirmSelect()、WaitForConfirms()和WaitForConfirmsOrDie。 channel.ConfirmSelect() 表示开启Confirm模式; channel.WaitForConfirms() 等待所有消息确认,如果所有的消息都被服务端成功接收返回true,只要有一条没有被成功接收就返回false。 channel.WaitForConfirmsOrDie() 和WaitForConfirms作用类似,也是等待所有消息确认,区别在于该方法没有返回值(Void),如果有任意一条消息没有被成功接收,该方法会立即抛出一个OperationInterruptedException类型异常。看一个Confirm模式的简单实现:
//开启Confirm模式
channel.ConfirmSelect();
//发送消息
channel.BasicPublish(exchange: "myexchange",
routingKey: "mykey",
basicProperties: null,
body: body);
//WaitForConfirms确认消息(可以同时确认多条消息)是否发送成功
if (channel.WaitForConfirms())
{
Console.WriteLine($"【{message}】发送到Broker成功!");
}Code language: JavaScript (javascript)