作者:裴潇艳 | 来源:互联网 | 2023-10-16 19:29
什么是消费端的ACK和重回队列?消费端的手工ACK和NACK消费端进行消费的时候,如果由于业务异常我们可以进行日志的记录,然后进行补偿如果由于服务器宕机等严重问题
什么是消费端的ACK和重回队列?
消费端的手工ACK和NACK
- 消费端进行消费的时候,如果由于业务异常我们可以进行日志的记录,然后进行补偿
- 如果由于服务器宕机等严重问题,那我们就需要手工进行ACK保障费端消费成功!
消费端重回队列
- 消费端重回队列是为了对没有处理成功的消息,把消息重新会递给Broker
- 一般我们在实际应用中,都会关闭重回队列,也就是设置为 False
生产端代码
package com.bfxy.rabbitmq.api.ack;import java.util.HashMap;
import java.util.Map;import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;public class Producer {public static void main(String[] args) throws Exception {ConnectionFactory connectionFactory = new ConnectionFactory();connectionFactory.setHost("localhost");
// connectionFactory.setHost("192.168.43.223");connectionFactory.setPort(5672);connectionFactory.setVirtualHost("/");connectionFactory.setUsername("guest");connectionFactory.setPassword("guest");Connection connection &#61; connectionFactory.newConnection();Channel channel &#61; connection.createChannel();String exchange &#61; "test_ack_exchange";String routingKey &#61; "ack.save";for(int i &#61;0; i<5; i &#43;&#43;){Map headers &#61; new HashMap();headers.put("num", i);AMQP.BasicProperties properties &#61; new AMQP.BasicProperties.Builder().deliveryMode(2).contentEncoding("UTF-8").headers(headers).build();String msg &#61; "Hello RabbitMQ ACK Message " &#43; i;channel.basicPublish(exchange, routingKey, true, properties, msg.getBytes());}}
}
消费者
package com.bfxy.rabbitmq.api.ack;import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.QueueingConsumer;
import com.rabbitmq.client.QueueingConsumer.Delivery;public class Consumer {public static void main(String[] args) throws Exception {ConnectionFactory connectionFactory &#61; new ConnectionFactory();connectionFactory.setHost("localhost");
// connectionFactory.setHost("192.168.43.223");connectionFactory.setPort(5672);connectionFactory.setVirtualHost("/");connectionFactory.setUsername("guest");connectionFactory.setPassword("guest");Connection connection &#61; connectionFactory.newConnection();Channel channel &#61; connection.createChannel();String exchangeName &#61; "test_ack_exchange";String queueName &#61; "test_ack_queue";String routingKey &#61; "ack.#";channel.exchangeDeclare(exchangeName, "topic", true, false, null);channel.queueDeclare(queueName, true, false, false, null);channel.queueBind(queueName, exchangeName, routingKey);// 手工签收 必须要关闭 autoAck &#61; falsechannel.basicConsume(queueName, false, new MyConsumer(channel));}
}
package com.bfxy.rabbitmq.api.ack;import java.io.IOException;import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.DefaultConsumer;
import com.rabbitmq.client.Envelope;public class MyConsumer extends DefaultConsumer {private Channel channel ;public MyConsumer(Channel channel) {super(channel);this.channel &#61; channel;}&#64;Overridepublic void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {System.err.println("-----------consume message----------");System.err.println("body: " &#43; new String(body));try {Thread.sleep(2000);} catch (InterruptedException e) {e.printStackTrace();}if((Integer)properties.getHeaders().get("num") &#61;&#61; 0) {channel.basicNack(envelope.getDeliveryTag(), false, true);} else {channel.basicAck(envelope.getDeliveryTag(), false);}}
}
运行结果&#xff1a;
重回队列&#xff0c;回到队列末尾
有一条消息没有被消费