MQ总结 v1.0

  1. 什么时候用?

场景:

  1. 系统处理高耗时 需要异步操作

  2. 解决原来的线程池丢消息的情况

  3. 使用


2.1. 开始

  1. 安装mq

  2. 安装mq管理UI 初始账户密码为guest

  3. 导入依赖 编写demo 测试

  4. 创建工厂 建立连接

  5. 创建信道 用信道声明队列 声明交换机 绑定交换机等操作 生产者消费者在声明队列时需参数保持一致

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. 死信

死信:

  1. 被拒绝的消息

  2. 过期的消息

  3. 被删除的消息

当定义了死信交换机和死信队列 这些消息就会由该交换机转发到死信队列

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 -> { }); } }
0个评论
点击登录,快来和大家讨论吧~
表情
图片
暂无评论
羊子
作者分享
计算机网络思维导图 进程章
2
微服务学习入门
6
今天差点入职 java 实习了 还是一家大公司 结果走入职通道被告知上面领导说专升本不要 硬性条件 服了 原来我只写了两年本科 没有写专科 心情起伏太大了 专升本 gg 了吗🥹 我应该注意什么
7
周五和周一面了两家 今天又遭拷打 周五: 1.自我介绍 2.讲讲你的项目 3.ArrayList 和 linkedList 的区别(真的好爱问) 4.创建线程的方式 5.redis 有哪些数据结构 6.mysql 的 sql 优化的思路 7.索引失效的原因 8.模糊查询索引失效有什么解决办法 9.ThreadLocal 讲一下 反问 找人的倾向 实习生工作 上班时间 周一(今天): 1.自我介绍 2.为什么没有想去考研和考公 前端我不知道为啥问了我很多问题? 3.项目用到了 react ? 4.你怎么调用的 AI 哪家 AI 名字是什么 5.项目中你觉得难点在哪里 6.你的生成图表的项目临时自己更换生成的结果吗 结果是流式输出的吗 知道流式输出吗 后端怎么实现的(麻了) 7.vue 的生命周期 8.你用的 ts 那么 ts 和 js 的区别是什么 ts 有哪些类型 9.你用的 elementUI 你自定义过样式吗 怎么做的 10. 你是怎么调教 AI 提示词的 调过 openAI 吗 后端 1.ArrayList 和 linkedList 的区别 2.用过 mangodb 吗 用过 python 吗 3.Spring 和 springboot 的区别 4.自己搭建过 mysql 吗 遇到什么问题了吗(这些个问题什么意思) 自己搭建过 redis 吗 5.redis 怎么清除所有的 key 6.你知道哪些数据库 你知道 sql 和 nosql 的区别吗 7.你有深入了解过 Es 吗 讲一下 Es 8.创建线程的方式 9.你知道数据库的范式吗 现在让你设计一个课程表 你会怎么设计 有哪些范式要注意(真没见过) 10.你的绩点 3.5 平均分多少 你的课上完了吗 下学期有哪些课你知道吗 学分修了多少了 反问 实习生的业务 其他没问了 我感觉回答的很不好 很多都半懂 估计 g 了 太难啦!
7
面试总结
5
下载 APP