002 JavaClent操作RabbitMQ

Java Client操作RabbitMQ

文章目录

  • Java Client操作RabbitMQ
    • 1.pom依赖
    • 2.连接工具类
    • 3.简单模式
    • 4.工作队列模式(work)
      • 公平调度
      • 示例
    • 5.发布/订阅模式(fanout)
      • 交换机
      • 绑定
      • 示例代码
    • 6.路由模式(direct)
    • 7.Topic匹配模式

1.pom依赖

<!-- https://mvnrepository.com/artifact/com.rabbitmq/amqp-client -->
<dependency><groupId>com.rabbitmq</groupId><artifactId>amqp-client</artifactId><version>5.20.0</version>
</dependency>

2.连接工具类

/*** rabbitmq连接工具类* @author moshangshang*/
@Slf4j
public class RabbitMQUtil {private static final String HOST_ADDRESS="192.168.1.102";private static final Integer PORT=5672;private static final String VIRTUAL_HOST="my_vhost";private static final String USER_NAME="root";private static final String PASSWORD="root";public static Connection getConnection() throws Exception {com.rabbitmq.client.ConnectionFactory factory=new com.rabbitmq.client.ConnectionFactory();factory.setHost(HOST_ADDRESS);factory.setPort(PORT);factory.setVirtualHost(VIRTUAL_HOST);factory.setUsername(USER_NAME);factory.setPassword(PASSWORD);return factory.newConnection();}public static void main(String[] args) {Connection connection = null;try {connection = getConnection();} catch (Exception e) {log.error("get rabbitmq connection exception....",e);}finally {try {if(connection!=null){connection.close();}} catch (IOException e) {log.error("close rabbitmq connection exception....",e);}}}}

3.简单模式

生产者投递消费到队列进行消费

在这里插入图片描述

消息发送

public class Send {private final static String QUEUE_NAME = "hello";public static void main(String[] argv) throws Exception {Connection connection = RabbitMQUtil.getConnection();try (Channel channel = connection.createChannel()) {//声明队列/**  如果队列不存在,则会创建*  Rabbitmq不允许创建两个相同的队列名称,否则会报错。**  @params1: queue 队列的名称*  @params2: durable 队列是否持久化*  @params3: exclusive 是否排他,即是否私有的,如果为true,会对当前队列加锁,其他的通道不能访问,并且连接自动关闭*  @params4: autoDelete 是否自动删除,当最后一个消费者断开连接之后是否自动删除消息。*  @params5: arguments 可以设置队列附加参数,设置队列的有效期,消息的最大长度,队列的消息生命周期等等。* */channel.queueDeclare(QUEUE_NAME, false, false, false, null);String message = "Hello World!";channel.basicPublish("", QUEUE_NAME, null, message.getBytes());System.out.println(" [x] Sent '" + message + "'");}}
}

消息接收

因为希望在消费者异步监听消息到达时,当前程序能够继续执行,而不是退出。

因为提供了一个DeliverCallback回调,该回调将缓冲消息,直到准备使用它们。

public class Recv {private final static String QUEUE_NAME = "hello";public static void main(String[] argv) throws Exception {Connection connection = RabbitMQUtil.getConnection();Channel channel = connection.createChannel();channel.queueDeclare(QUEUE_NAME, false, false, false, null);DeliverCallback deliverCallback = (consumerTag, delivery) -> {String message = new String(delivery.getBody(), StandardCharsets.UTF_8);System.out.println(" [x] Received '" + message + "'");};channel.basicConsume(QUEUE_NAME, true, deliverCallback, consumerTag -> { });}
}

4.工作队列模式(work)

生产者直接投递消息到队列,存在多个消费者情况

  • 创建一个工作队列,用于在多个工作人员之间分配耗时的任务。
  • 工作队列(又名:任务队列)背后的主要思想是避免立即执行资源密集型任务,并必须等待其完成。相反,我们把任务安排在以后完成。我们将任务封装为消息并将其发送到队列。在后台运行的工作进程将弹出任务并最终执行作业。当你运行多个worker时,任务将在它们之间共享。
  • 这个概念在web应用程序中特别有用,因为在短的HTTP请求窗口内无法处理复杂的任务。
  • 默认情况下,消费者会进行轮询调度
  • RabbitMQ支持消息确认。消费者发送回一个确认,告诉RabbitMQ已经收到、处理了一条特定的消息,RabbitMQ可以自由删除它。
  • 如果一个消费者在没有发送ack的情况下死亡(其通道关闭、连接关闭或TCP连接丢失),RabbitMQ将理解消息未完全处理,并将其重新排队。如果同时有其他消费者在线,它将迅速将其重新传递给另一个消费者。这样,即使worker偶尔挂掉,也可以确保没有信息丢失。
  • 消费者交付确认时强制执行超时(默认为30分钟)。这有助于检测一直没有确认的消费者。
  • 默认情况下,手动消息确认已打开。在前面的示例中,我们通过autoAck=true标志明确地关闭了它们。一旦我们完成了一项任务,是时候将此标志设置为false并从worker发送适当的确认了。

在这里插入图片描述

公平调度

由于默认轮询调度,有些任务执行时间长,有些短,所以会导致部分worker压力大

使用预取计数=1设置的basicQos方法。这条消息告诉RabbitMQ一次不要给一个worker发送多条消息。在处理并确认前一条消息之前,不要向worker发送新消息。相反,它会将其发送给下一个不忙的worker。

int prefetchCount = 1;
channel.basicQos(prefetchCount);

示例

public class WorkProvider {private static final String TASK_QUEUE_NAME = "work_queue";public static void main(String[] argv) throws Exception {Connection connection = RabbitMQUtil.getConnection();try (Channel channel = connection.createChannel()) {channel.queueDeclare(TASK_QUEUE_NAME, true, false, false, null);for (int i = 0; i < 8; i++) {String message = String.valueOf(i);channel.basicPublish("", TASK_QUEUE_NAME,MessageProperties.PERSISTENT_TEXT_PLAIN,message.getBytes(StandardCharsets.UTF_8));System.out.println(" 消息发送 :'" + i + "'");}}}
}
public class WorkerConsumer1 {private static final String TASK_QUEUE_NAME = "work_queue";public static void main(String[] argv) throws Exception {final Connection connection = RabbitMQUtil.getConnection();final Channel channel = connection.createChannel();channel.queueDeclare(TASK_QUEUE_NAME, true, false, false, null);System.out.println(" 消息监听中。。。。。。");//控制ack流速,表示每次进行ack确认前只会处理一条消息//channel.basicQos(1);DeliverCallback deliverCallback = (consumerTag, delivery) -> {//获取消息String message = new String(delivery.getBody(), StandardCharsets.UTF_8);System.out.println("worker1 消息消费:'" + message + "'");try {doWork(message);} finally {System.out.println(" 执行结束。。");//消息确认,根据消息序号(false只确认当前一个消息收到,true确认所有比当前序号小的消息(成功消费,消息从队列中删除 ))channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);}};//设置自动应答channel.basicConsume(TASK_QUEUE_NAME, false, deliverCallback, consumerTag -> { });}}

消费者1的方法处理

  private static void doWork(String task) {try {Thread.sleep(1000);} catch (InterruptedException ignored) {Thread.currentThread().interrupt();}}

消费者2的方法处理

  private static void doWork(String task) {return;}

启动两个worker消费者,执行结果如下(轮询):

在这里插入图片描述
在这里插入图片描述

若设置公平调度

channel.basicQos(1);

测试结果:

在这里插入图片描述
在这里插入图片描述

5.发布/订阅模式(fanout)

  • 不同于工作队列,同一消息在所有消费者共享,但只能有一个消费者消费,而发布订阅则会将同一消息发送给多个消费者,则将消息广播给所有订阅者

  • 在之前的模式中,都是直接将消息发送给队列,然后从队列消费,事实上之前使用了一个默认的交换机,即“”空字符串的

  • RabbitMQ消息传递模型的核心思想是生产者从不直接向队列发送任何消息。实际上,很多时候,生产者甚至根本不知道消息是否会被传递到任何队列。

  • 相反,生产者只能向exchange发送消息。exchange是一件非常简单的事情。它一方面接收来自生产者的消息,另一方面将它们推送到队列。exchange必须确切地知道如何处理它收到的消息。

在这里插入图片描述

将消息生产投递到exchange,由交换机去投递消息到队列

交换机

有几种exchange类型可供选择:directtopicheadersfanout。我们将专注于最后一个fanout

//创建交换机名称为logs
channel.exchangeDeclare("logs", "fanout");/*** exchange:交换机的名称* type:交换机的类型* durable 队列是否持久化* autoDelete:是否自动删除,(当该交换机上绑定的最后一个队列解除绑定后,该交换机自动删除)* internal:是否是内置的,true表示内置交换器。(则无法直接发消息给内置交换机,只能通过其他交换机路由到该交换机)* argument:其他一些参数*/channel.exchangeDeclare(EXCHANGE_NAME,"fanout",false,false,false,null);

第一个参数是exchange的名称。空字符串表示默认或未命名的交换:消息将被路由到routingKey指定名称的队列(如果存在)。

channel.basicPublish( "logs", "", null, message.getBytes());

绑定

我们已经创建了一个fanout交换机和一个队列。现在我们需要告诉exchange向我们的队列发送消息。交换和队列之间的关系称为绑定。

交换机会向绑定的队列通过路由key将消息路由到指定的队列中,fanout分发不需要路由key

在这里插入图片描述

//其中,第一个参数为绑定的队列,第二个参数为绑定的交换机,第三个参数为路由key
channel.queueBind(queueName, "logs", "");
#列出所有得绑定
rabbitmqctl list_bindings

示例代码

public class FanoutProvider {//声明交换机public static final String EXCHANGE_NAME="fanoutTest";//声明队列public static final String QUEUE_NAME1="queue_name1";public static final String QUEUE_NAME2="queue_name2";public static void main(String[] argv) throws Exception {Connection connection = RabbitMQUtil.getConnection();try (Channel channel = connection.createChannel()) {//声明交换机channel.exchangeDeclare(EXCHANGE_NAME,"fanout",false,false,false,null);//声明队列channel.queueDeclare(QUEUE_NAME1,false,false,false,null);channel.queueDeclare(QUEUE_NAME2,false,false,false,null);//进行队列绑定channel.queueBind(QUEUE_NAME1,EXCHANGE_NAME,"");channel.queueBind(QUEUE_NAME2,EXCHANGE_NAME,"");String message = "fanout模式消息推送。。。。。";//消息推送//参数说明:交换机,路由key/队列,消息属性,消息体channel.basicPublish(EXCHANGE_NAME, "", MessageProperties.PERSISTENT_TEXT_PLAIN, message.getBytes(StandardCharsets.UTF_8));System.out.println(" 消息发送 :'" +message + "'");}}}
public class FanoutConsumer1 {public static final String EXCHANGE_NAME="fanoutTest";public static final String QUEUE_NAME1="queue_name1";public static void main(String[] argv) throws Exception {Connection connection = RabbitMQUtil.getConnection();//创建信道Channel channel = connection.createChannel();//声明交换机channel.exchangeDeclare(EXCHANGE_NAME, "fanout");//绑定队列channel.queueBind(QUEUE_NAME1, EXCHANGE_NAME, "");DeliverCallback deliverCallback = (consumerTag, delivery) -> {String message = new String(delivery.getBody(), StandardCharsets.UTF_8);System.out.println(" fanout 消费者1:'" + message + "'");};channel.basicConsume(QUEUE_NAME1, true, deliverCallback, consumerTag -> {});}}
public class FanoutConsumer2 {public static final String EXCHANGE_NAME="fanoutTest";public static final String QUEUE_NAME2="queue_name2";public static void main(String[] argv) throws Exception {Connection connection = RabbitMQUtil.getConnection();//创建信道Channel channel = connection.createChannel();//声明交换机channel.exchangeDeclare(EXCHANGE_NAME, "fanout");//绑定队列channel.queueBind(QUEUE_NAME2, EXCHANGE_NAME, "");DeliverCallback deliverCallback = (consumerTag, delivery) -> {String message = new String(delivery.getBody(), StandardCharsets.UTF_8);System.out.println(" fanout 消费者2:'" + message + "'");};channel.basicConsume(QUEUE_NAME2, true, deliverCallback, consumerTag -> { });}}

6.路由模式(direct)

由交换机通过路由key绑定key进行消息推送,也可以将同一个路由key绑定到多个队列或所有队列,此时相当于fanout

如果推送消息的路由key不存在,则该消息会丢弃

public class DirectProvider {public static final String EXCHANGE_NAME="direct-exchange";public static final String QUEUE_NAME1="direct-queue";public static final String ROUTING_KEY="change:direct";public static void main(String[] argv) throws Exception {Connection connection = RabbitMQUtil.getConnection();try (Channel channel = connection.createChannel()) {//声明交换机channel.exchangeDeclare(EXCHANGE_NAME,"direct",false,false,false,null);//声明队列channel.queueDeclare(QUEUE_NAME1,false,false,false,null);String message = "direct模式消息推送。。。。。";//参数说明:交换机,路由key/队列,消息属性,消息体channel.basicPublish(EXCHANGE_NAME, ROUTING_KEY, MessageProperties.PERSISTENT_TEXT_PLAIN, message.getBytes(StandardCharsets.UTF_8));System.out.println(" 消息发送 :'" +message + "'");}}}
public class DirectConsumer {public static final String EXCHANGE_NAME="direct-exchange";public static final String QUEUE_NAME1="direct-queue";public static final String BINDING_KEY="change:direct";public static void main(String[] argv) throws Exception {Connection connection = RabbitMQUtil.getConnection();//创建信道Channel channel = connection.createChannel();//声明交换机channel.exchangeDeclare(EXCHANGE_NAME, "direct");//绑定队列channel.queueBind(QUEUE_NAME1, EXCHANGE_NAME, BINDING_KEY);DeliverCallback deliverCallback = (consumerTag, delivery) -> {String message = new String(delivery.getBody(), StandardCharsets.UTF_8);System.out.println(" direct 消费者1:'" + message + "'");};channel.basicConsume(QUEUE_NAME1, true, deliverCallback, consumerTag -> { });}}

7.Topic匹配模式

  • 发送到主题交换的消息不能有任意的路由key,它必须是一个由点分隔的单词列表。单词可以是任何东西,但通常它们指定了与消息相关的一些特征。一些有效的路由key示例:stock.usd.nyse、nyse.vmw、quick.orange.rabbit。路由密钥中可以有任意多的单词,最多255个字节。
  • 绑定key也必须采用相同的形式。topic交换机背后的逻辑类似于direct交换机,使用特定路由key发送的消息将被传递到所有使用绑定key绑定的所有队列。但是,绑定密钥有两个重要的特殊情况:
  • *(星号)只能代替一个单词
  • #(hash)可以替代零个或多个单词
public class TopicProvider {public static final String EXCHANGE_NAME="topic-exchange";public static final String QUEUE_NAME1="topic-queue";public static final String ROUTING_KEY="com.orange.test";public static final String ROUTING_KEY2="com.orange.test.aaa";public static void main(String[] argv) throws Exception {Connection connection = RabbitMQUtil.getConnection();try (Channel channel = connection.createChannel()) {//声明交换机channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.TOPIC,false,false,false,null);//声明队列channel.queueDeclare(QUEUE_NAME1,false,false,false,null);String message1 = "topic test模式消息推送。。。。。";String message2 = "topic test.aaa模式消息推送。。。。。";//参数说明:交换机,路由key/队列,消息属性,消息体channel.basicPublish(EXCHANGE_NAME, ROUTING_KEY, MessageProperties.PERSISTENT_TEXT_PLAIN, message1.getBytes(StandardCharsets.UTF_8));channel.basicPublish(EXCHANGE_NAME, ROUTING_KEY2, MessageProperties.PERSISTENT_TEXT_PLAIN, message2.getBytes(StandardCharsets.UTF_8));System.out.println(" 消息发送 :'" +message1 + "'");System.out.println(" 消息发送 :'" +message2 + "'");}}}
public class TopicConsumer {public static final String EXCHANGE_NAME="topic-exchange";public static final String QUEUE_NAME1="topic-queue";public static final String BINDING_KEY="*.orange.#";public static void main(String[] argv) throws Exception {Connection connection = RabbitMQUtil.getConnection();//创建信道Channel channel = connection.createChannel();//声明交换机channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.TOPIC);//绑定队列channel.queueBind(QUEUE_NAME1, EXCHANGE_NAME, BINDING_KEY);DeliverCallback deliverCallback = (consumerTag, delivery) -> {String message = new String(delivery.getBody(), StandardCharsets.UTF_8);System.out.println(" topic 消费者1:'" + message + "'");};channel.basicConsume(QUEUE_NAME1, true, deliverCallback, consumerTag -> { });}}

本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若转载,请注明出处:http://www.mzph.cn/bicheng/53849.shtml

如若内容造成侵权/违法违规/事实不符,请联系多彩编程网进行投诉反馈email:809451989@qq.com,一经查实,立即删除!

相关文章

LVS--负载均衡调度器

文章目录 集群和分布式集群分布式 LVS介绍LVS特点LVS工作原理LVS集群架构 LVS集群中的术语CIPVIPRSDIPRIP LVS集群的工作模式NAT模式DR模式DR的工作原理DR的特点:DR的网络配置1.配置负载均衡器2.配置后端服务器lo接口的作用 3.测试连接&#xff1a; DR的典型应用场景 TUN模式 L…

《深度学习》【项目】 OpenCV 身份证号识别

目录 一、项目实施 1、自定义函数 2、定位模版图像中的数字 1&#xff09;模版图二值化处理 运行结果&#xff1a; 2&#xff09;展示所有数字 运行结果&#xff1a; 3、识别身份证号 1&#xff09;灰度图、二值化图展示 运行结果 2&#xff09;定位身份证号每一个数…

❤Node08-Express-jwt身份认证

❤Node08-Express-jwt身份认证 1、token基本概念​ Session认证的局限性​ Session 认证机制需要配合Cookie才能实现。由于 Cookie 默认不支持跨域访问&#xff0c;所以&#xff0c;当涉及到前端跨域请求后端接口的时候&#xff0c;需要做很多额外的配置&#xff0c;才能实现…

【JVM】JVM栈帧中的动态链接 与 Java的面向对象特性--多态

栈帧 每一次方法调用都会有一个对应的栈帧被压入栈&#xff08;虚拟机栈&#xff09;中&#xff0c;每一个方法调用结束后&#xff0c;都会有一个栈帧被弹出。 每个栈帧中包括&#xff1a;局部变量表、操作数栈、动态链接、方法返回地址。 JavaGuide&#xff1a;Java内存区域…

Debian项目实战——环境搭建篇

Debian系统安装 准备工作 1、系统镜像&#xff1a;根据自己的需要选择合适的版本格式&#xff1a;x86 / arm 架构 | 最好下载离线安装版本 | 清华镜像源 2、制作工具&#xff1a;balenaEtcher 3、系统媒介&#xff1a;16G以上U盘最佳 烧录镜像 打开balenaEtcher进行烧录&am…

改变事件

窗口的某些属性的状态发生改变时就会触发该事件 对应的事件类型包括 QEvent::ToolBarChange, QEvent::ActivationChange, QEvent::EnabledChange, QEvent::FontChange,QEvent::StyleChange, QEvent::PaletteChange, QEvent::WindowTitleChange, QEvent::IconTextChange, QEve…

【练习9】大数加法

链接&#xff1a;大数加法__牛客网 (nowcoder.com) 分析&#xff1a; 当作竖式计算 import java.util.*;public class Solution {public String solve (String s, String t) {StringBuffer ret new StringBuffer();//i是字符串s的最后一个字符的索引int i s.length() - 1;//j…

新能源汽车安全问题如何解决?细看“保护罩”连接器的守护使命

「当前市场上绝大部分电池的安全系数远远不够」。 在一场世界动力电池大会上&#xff0c;宁德时代的董事长曾毓群这样犀利直言。 从汽车开始向电动化转型升级那天起&#xff0c;动力电池的安全隐患就一直是个老生常谈的话题了。曾毓群的这句话&#xff0c;直接点明了行业的发展…

参数传了报错没传参数识别不到参数传丢

【记一次参数传值了但报错未传值的问题解决历程】 问题描述&#xff1a;同一个接口&#xff0c;用测试类调可以成功&#xff0c;用postman调用一直报错少参数&#xff0c;后又尝试了用idea自带的http调用&#xff0c;同样报错参数未传值。 如图&#xff0c;传值了报错未传值。…

Java并发:互斥锁,读写锁,Condition,StampedLock

3&#xff0c;Lock与Condition 3.1&#xff0c;互斥锁 3.1.1&#xff0c;可重入锁 锁的可重入性&#xff08;Reentrant Locking&#xff09;是指在同一个线程中&#xff0c;已经获取锁的线程可以再次获取该锁而不会导致死锁。这种特性允许线程在持有锁的情况下&#xff0c;可…

AI网盘搜索 1.2.6 智能文件搜索助手,一键搜索所有资源

对于经常需要处理大量文件的人来说&#xff0c;AI网盘检索简直是救星。它提供了智能对话式搜索功能&#xff0c;只需用自然语言描述就能找到需要的文件。此外&#xff0c;它还广泛支持各种文件类型&#xff0c;从文档到图片&#xff0c;全面覆盖。精准定位功能让您能够快速找到…

DSC+主备+异步备库搭建

DSC主备异步备库搭建 本次在DSC的基础上进行主备集群异步备库的搭建&#xff0c;实现DSC主备异步备库的集合。 这里DMDSC集群是看做一个数据库服务&#xff08;即DSC集群内的都叫主库&#xff09;&#xff0c;备库是一个单机实例 环境配置 服务器配置 端口配置 实例名PORT…

C#获取计算机信息

目录 效果 项目 代码 下载 效果 项目 代码 using System; using System.Collections.Generic; using System.ComponentModel; using System.Data; using System.Drawing; using System.Linq; using System.Text; using System.Windows.Forms; using System.Management; n…

Vulnhub:bassamCTF

靶机下载地址 信息收集 主机发现 扫描攻击机同网段存活主机。 nmap 192.168.31.0/24 -Pn -T4 靶机ip&#xff1a;192.168.31.165 端口扫描 nmap 192.168.31.165 -A -p- -T4 开放端口22&#xff0c;80。 网站信息收集 访问80端口的http服务。首页是空白页面&#xff0c;…

关于打不开SOAMANAGER如何解决

参考文章&#xff1a;https://blog.csdn.net/yannickdann/article/details/115396035 打开SE93

Django创建模型

1、根据创建好应用模块 python manage.py startapp tests 2、在models文件里创建模型 from django.db import modelsfrom book.models import User# Create your models here. class Tests(models.Model):STATUS_CHOICES ((0, 启用),(1, 停用),# 更多状态...)add_time mode…

人工智能(AI)领域各方向顶会和顶刊

在人工智能&#xff08;AI&#xff09;这个快速发展的领域&#xff0c;研究人员和从业者需要紧跟最新的研究动态和技术进展。顶级的会议和期刊是获取最新科研成果和交流思想的重要平台。以下是人工智能领域内不同方向的顶级会议和期刊概览。 顶级会议 人工智能基础与综合 A…

ROADM(可重构光分插复用器)-介绍

1. 引用 https://zhuanlan.zhihu.com/p/163369296 https://zhuanlan.zhihu.com/p/521352954 https://zhuanlan.zhihu.com/p/91103069 https://zhuanlan.zhihu.com/p/50610236 术语&#xff1a; 英文缩写描述灰光模块彩光模块CWDM&#xff1a;Coarse Wave-Length Division …

IT前端好用的工具集

在线抠图网站 https://www.remove.bg/ 将iconfont转成css显示 https://transfonter.org/ 免费的在线图片压缩 https://tinypng.com/ JSON在线格式化工具 https://www.sojson.com/ 国内人工智能kimi.moonshot工具 https://kimi.moonshot.cn/chat/crft7a6sdv14grouufs0 自动…

Android生成Java AIDL

AIDL:Android Interface Definition Language AIDL是为了实现进程间通信而设计的Android接口语言 Android进程间通信有多种方式&#xff0c;Binder机制是其中最常见的一种 AIDL的本质就是基于对Binder的运用从而实现进程间通信 这篇博文从实战出发&#xff0c;用一个尽可能…