MQ总结 v1.0
- 什么时候用?
场景:
-
系统处理高耗时 需要异步操作
-
解决原来的线程池丢消息的情况
-
使用
2.1. 开始
-
安装mq
-
安装mq管理UI 初始账户密码为guest
-
导入依赖 编写demo 测试
-
创建工厂 建立连接
-
创建信道 用信道声明队列 声明交换机 绑定交换机等操作 生产者消费者在声明队列时需参数保持一致
2.2. 交换机 exchange
介绍: 如计网里面的交换机根据MAC地址的转发数据帧的效果类似 不同的交换机可以根据不同的规则转发消息
分类:
fanout;
direct;
topic;
2.2.1. fanout
介绍: '扇出' 交换机 名如其机 感觉一下子消息全部扇到了脸上一样
特点: 就像计网的广播帧 转发给所有信道 每个都可以收到消息
使用:https://www.rabbitmq.com/tutorials/tutorial-three-java
生产者信道在声明交换机时 就声明fanout 所有绑定该交换机的消费者信道都会受到消息
2.2.2. direct
介绍: 直连交换机 可以根据路由key 直接发送消息
单绑定:

多绑定:

使用:https://www.rabbitmq.com/tutorials/tutorial-four-java
在声明队列时 声明路由参数 那么创建的交换机在转发消息时就会根据路由key发送到相应的队列
可以应用的场景:
根据不同的路由 发送不同的日志消息到消费者

2.2.3. topic
介紹: topic 交换机是用根据规则来模糊匹配路由key的交换机

*是用来匹配一个词
#是用来匹配0个或者多个词
举例:
Q1: .orange.
Q2: ..rabbith 和lazy.#.
那么 一下方式
quick.orange.rabbit 匹配Q1 Q2
lazy.orange.elephant 匹配Q1 Q2
quick.orange.fox 匹配Q1
lazy.brown.fox 匹配Q2
lazy.pink.rabbit 匹配Q2 只匹配一次
使用
https://www.rabbitmq.com/tutorials/tutorial-five-java
2.3. 发布确认机制 Publish Confirm
2.3.1. 设置消息确认的到期时间
▼java复制代码while (thereAreMessagesToPublish()) { byte[] body = ...; BasicProperties properties = ...; channel.basicPublish(exchange, queue, properties, body); // uses a 5 second timeout channel.waitForConfirmsOrDie(5_000); }
在规定的到期时间内 未确认消息 或者该消息被拒绝 就会抛出一个异常
处理该异常通常是 重发或者是打印错误日志
优点: 简单
缺点: 当确认消息阻塞发布消息时 发布效率就会很低 不适合高吞吐量的系统
2.3.2. 批量发布
▼java复制代码// 设置批次数量大小 int batchSize = 100; int outstandingMessageCount = 0; while (thereAreMessagesToPublish()) { byte[] body = ...; BasicProperties properties = ...; channel.basicPublish(exchange, queue, properties, body); outstandingMessageCount++; // 发送数量没到就不设置过期时间 if (outstandingMessageCount == batchSize) { channel.waitForConfirmsOrDie(5_000); outstandingMessageCount = 0; } } // 发送完了 发现都还没到数量 就直接设置过期时间 if (outstandingMessageCount > 0) { channel.waitForConfirmsOrDie(5_000); }
优点: 相比于接受单独的确认 效率好很多
缺点: 出问题不知道是哪个出了错
2.3.3. 手动确认
2.4. 死信
死信:
-
被拒绝的消息
-
过期的消息
-
被删除的消息
当定义了死信交换机和死信队列 这些消息就会由该交换机转发到死信队列
2.4.1. 死信交换机
设置: https://www.rabbitmq.com/docs/dlx
▼java复制代码channel.exchangeDeclare("some.exchange.name", "direct"); // Important: prefer using policies over hardcoded x-arguments Map<String, Object> args = new HashMap<String, Object>(); //绑定个死信交换机 args.put("x-dead-letter-exchange", "some.exchange.name"); //给绑定的死信交换机 指明它所要转发死信到哪条死信队列上去 args.put("x-dead-letter-routing-key", "some-routing-key"); //此时 该队列就绑定了死信交换机 channel.queueDeclare("myqueue", false, false, false, args);
死信交换机就像direct交换机一样 正常创建就好
2.4.2. 示例

创建死信的direct交换机和队列
▼java复制代码public class DieLetterWorker { private final static String QUEUE_NAME1 = "m3"; private final static String QUEUE_NAME2 = "m4"; public static void main(String[] argv) throws Exception { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("localhost"); Connection connection = factory.newConnection(); Channel channel1 = connection.createChannel(); Channel channel2 = connection.createChannel(); //声明死信交换机 channel1.exchangeDeclare("e2", "direct"); //声明死信队列 channel1.queueDeclare(QUEUE_NAME1, true, false, false, null); //绑定队列到交换机 channel1.queueBind(QUEUE_NAME1, "e2", "worker1"); //声明交换机 同上 channel2.exchangeDeclare("e2", "direct"); channel2.queueDeclare(QUEUE_NAME2, true, false, false, null); channel2.queueBind(QUEUE_NAME2, "e2", "worker2"); System.out.println(" [*] Waiting for messages. To exit press CTRL+C"); DeliverCallback deliverCallback1 = (consumerTag, delivery) -> { try { String message = new String(delivery.getBody(), StandardCharsets.UTF_8); System.out.println(" [x]worker1 Received1 '" + message + "'"); //模拟消息处理 Thread.sleep(1000); //如果消息处理失败,将消息发送到死信队列 } catch (InterruptedException e) { throw new RuntimeException(e); } }; DeliverCallback deliverCallback2 = (consumerTag, delivery) -> { try { String message = new String(delivery.getBody(), StandardCharsets.UTF_8); System.out.println(" [x]worker2 Received1 '" + message + "'"); //模拟消息处理 Thread.sleep(1000); } catch (InterruptedException e) { throw new RuntimeException(e); } }; channel1.basicConsume(QUEUE_NAME1, true, deliverCallback1, consumerTag -> { }); channel2.basicConsume(QUEUE_NAME2, true, deliverCallback2, consumerTag -> { }); } }
生产者
▼java复制代码public class DieLetterProvider { private final static String QUEUE_NAME1 = "hello1"; public static void main(String[] argv) throws Exception { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("localhost"); try (Connection connection = factory.newConnection(); //创建两个消息队列 Channel channel = connection.createChannel(); ) { //声明交换机 channel.exchangeDeclare("e1", "direct"); //手动输入 Scanner scanner = new Scanner(System.in); while (scanner.hasNextLine()) { String message = scanner.nextLine(); String[] strs = message.split(" "); //将消息发送到交换机 不指定路由 channel.basicPublish("e1", strs[1], null, message.getBytes(StandardCharsets.UTF_8)); System.out.println(" [x] Sent '" + message + "'"); } } } }
消费者
▼java复制代码public class DieLetterConsumer { private final static String QUEUE_NAME1 = "m1"; private final static String QUEUE_NAME2 = "m2"; public static void main(String[] argv) throws Exception { ConnectionFactory factory = new ConnectionFactory(); factory.setHost("localhost"); Connection connection = factory.newConnection(); Channel channel1 = connection.createChannel(); Channel channel2 = connection.createChannel(); //声明死信交换机 channel1.exchangeDeclare("e2", "direct"); channel2.exchangeDeclare("e2", "direct"); //声明死信队列参数 HashMap<String, Object> args1 = new HashMap<>(); args1.put("x-dead-letter-exchange", "e2"); args1.put("x-dead-letter-routing-key", "worker1"); HashMap<String, Object> args2 = new HashMap<>(); args2.put("x-dead-letter-exchange", "e2"); args2.put("x-dead-letter-routing-key", "worker2"); //声明交换机 channel1.exchangeDeclare("e1", "direct"); //声明队列 channel1.queueDeclare(QUEUE_NAME1, true, false, false, args1); //绑定队列到交换机 channel1.queueBind(QUEUE_NAME1, "e1", "exc1"); //声明交换机 同上 channel2.exchangeDeclare("e1", "direct"); channel2.queueDeclare(QUEUE_NAME2, true, false, false, args2); channel2.queueBind(QUEUE_NAME2, "e1", "exc2"); System.out.println(" [*] Waiting for messages. To exit press CTRL+C"); DeliverCallback deliverCallback1 = (consumerTag, delivery) -> { try { String message = new String(delivery.getBody(), StandardCharsets.UTF_8); if (message.equals("hello exc1")) { throw new InterruptedException("hello exc1"); } System.out.println(" [x]exc1 Received1 '" + message + "'"); //模拟消息处理 Thread.sleep(1000); //如果消息处理失败,将消息发送到死信队列 //确认消息 channel1.basicAck(delivery.getEnvelope().getDeliveryTag(), false); } catch (InterruptedException e) { //如果消息处理失败,将消息发送到死信队列 channel1.basicNack(delivery.getEnvelope().getDeliveryTag(), false, false); } }; DeliverCallback deliverCallback2 = (consumerTag, delivery) -> { try { String message = new String(delivery.getBody(), StandardCharsets.UTF_8); if (message.equals("hello exc2")) { throw new InterruptedException("hello exc2"); } System.out.println(" [x]exc2 Received1 '" + message + "'"); //模拟消息处理 Thread.sleep(1000); //确认消息 channel2.basicAck(delivery.getEnvelope().getDeliveryTag(), false); } catch (InterruptedException e) { //如果消息处理失败,将消息发送到死信队列 channel2.basicNack(delivery.getEnvelope().getDeliveryTag(), false, false); } }; channel1.basicConsume(QUEUE_NAME1, false, deliverCallback1, consumerTag -> { }); channel2.basicConsume(QUEUE_NAME2, false, deliverCallback2, consumerTag -> { }); } }
