RabbitMQ消息确认机制


RabbitMQ消息确认的本质也就是为了解决RabbitMQ消息丢失问题,因为哪怕我们做了,其实也并不能保证解决我们的消息丢失问题


RabbitMQ的消息确认有两种

  • 第一种是消息发送确认。这种是用来确认生产者将消息发送给交换器,交换器传递给队列的过程中,消息是否成功投递。发送确认分为两步,一是确认是否到达交换器,二是确认是否到达队列。
  • 第二种是消费接收确认。这种是确认消费者是否成功消费了队列中的消息。

1.消息发送确认(生产者)

当生产者发送消息给rabbitmq服务器时,消息是否真正的到达了服务器?为了保证生产者发送的消息能够可靠的发送到服务器(即消息落地),rabbitmq提供了两种方式:

  • 通过事务实现
  • 通过发送方确认机制(publisher confirm)实现

事务实现

  • channel.txSelect(): 将当前信道设置成事务模式
  • channel.txCommit(): 用于提交事务
  • channel.txRollback(): 用于回滚事务

通过事务实现机制,只有消息成功被rabbitmq服务器接收,事务才能提交成功,否则便可在捕获异常之后进行回滚,然后进行消息重发,但是事务非常影响rabbitmq的性能。还有就是事务机制是阻塞的过程,只有等待服务器回应之后才会处理下一条消息

 待更新。。。

2.消息接收确认(消费者)

消息接收确认机制,分为消息自动确认模式和消息手动确认模式,当消息确认后,我们队列中的消息将会移除

那这两种模式要如何选择呢?

  • 如果消息不太重要,丢失也没有影响,那么自动ACK会比较方便。好处就是可以提高吞吐量,缺点就是会丢失消息
  • 如果消息非常重要,不容丢失,则建议手动ACK,正常情况都是更建议使用手动ACK。虽然可以解决消息不会丢失的问题,但是可能会造成消费者过载

消息自动确认模式的实现

注:自动确认模式,消费者不会判断消费者是否成功接收到消息,也就是当我们程序代码有问题,我们的消息都会被自动确认,消息被自动确认了,我们队列就会移除该消息,这就会造成我们的消息丢失

/**
 * 消费者
 */
public class Recv {
    //设定队列名称(已存在的队列)
    private static final String QUEUE_NAME = "queue1";
    public static void main(String[] args) throws IOException, TimeoutException {
        //从mq工具类获取连接信息
        Connection connection = MqConnectionUtils.getConnection();
        //获取一个通道
        Channel channel = connection.createChannel();
        //监听该队列,true代表自动确认
        channel.basicConsume(QUEUE_NAME,true,new DefaultConsumer(channel){
            @Override
            public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties basicProperties, byte[] body) throws IOException{
                System.out.println("接收到的消息:"+ new String(body,"UTF-8"));
            }
        });
    }
}

实现效果,消费者会将我们队列中的消息全部接收然后确认,并移除队列

消息手动确认模式的实现

/**
 * 消费者
 */
public class Recv {
    //设定队列名称(已存在的队列)
    private static final String QUEUE_NAME = "queue1";
    public static void main(String[] args) throws IOException, TimeoutException {
        //从mq工具类获取连接信息
        Connection connection = MqConnectionUtils.getConnection();
        //获取一个通道
        Channel channel = connection.createChannel();
        //监听该队列,false代表手动确认
        channel.basicConsume(QUEUE_NAME,false,new DefaultConsumer(channel){
            @Override
            public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties basicProperties, byte[] body) throws IOException{
                System.out.println("接收到的消息:"+ new String(body,"UTF-8"));
            }
        });
    }
}

手动确认模式下,当我们消费者成功接收到消息后,在队列中消息会进入Unacked项,也就是待确认模式

所以我们还需要加上下列代码,来实现消息者在成功接收到消息后,手动确认

#添加红色字段

/**
 * 消费者
 */
public class Recv {
    //设定队列名称(已存在的队列)
    private static final String QUEUE_NAME = "queue1";
    public static void main(String[] args) throws IOException, TimeoutException {
        //从mq工具类获取连接信息
        Connection connection = MqConnectionUtils.getConnection();
        //获取一个通道
        Channel channel = connection.createChannel();
        //监听该队列
        channel.basicConsume(QUEUE_NAME,false,new DefaultConsumer(channel){
            @Override
            public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties basicProperties, byte[] body) throws IOException{
                System.out.println("接收到的消息:"+ new String(body,"UTF-8"));

                //获取消息的编号,我们需要根据消息的编号来确认消息
                long tag = envelope.getDeliveryTag();
                //获取当前内部类中的通道
                Channel c = this.getChannel();
                //手动确认消息,确认以后,则表示消息已经成功处理,消息就会从队列中移除
                c.basicAck(tag,true);            }
        });
    }
}

此时,我们的消息才会成功被确认,并移除队列