RabbitMQ 的核心概念,看了必懂!

作者:海向 出处:cnblogs.com/haixiang/p/10853467.html

RabbitMQ 特点

RabbitMQ 相较于其他消息队列,有一系列防止消息丢失的措施,拥有强悍的高可用性能,它的吞吐量可能没有其他消息队列大,但是其消息的保障性出类拔萃,被广泛用于金融类业务。

AMQP 协议

AMQP: Advanced Message Queuing Protocol 高级消息队列协议

AMQP定义:是具有现代特征的二进制协议。是一个提供统一消息服务的应用层标准高级消息队列协议,是应用层协议的一个开放标准,为面向消息的中间件设计。

Erlang语言最初在于交换机领域的架构模式,这样使得RabbitMQ在Broker之间进行数据交互的性能是非常优秀的,Erlang的优点: Erlang有着和原生Socket一样的延迟。

RabbitMQ是一个开源的消息代理和队列服务器,用来通过普通协议在完全不同的应用之间共享数据, RabbitMQ是使用Erlang语言来编写的,并且RabbitMQ是基于AMQP协议的。关注公众号Java技术栈获取系列RabbitMQ教程。

RabbitMQ 消息传递机制

生产者发送消息到指定的 Exchange,Exchange 依据自身的类型(direct、topic等),根据 routing key 将消息发送给 0 - n 个 队列,队列再将消息转发给了消费者。

Server: 又称Broker, 接受客户端的连接,实现AMQP实体服务,这里指RabbitMQ 服务器

Connection: 连接,应用程序与Broker的网络连接。

Channel: 网络信道,几乎所有的操作都在 Channel 中进行,Channel是进行消息读写的通道。客户端可建立多个Channel:,每个Channel代表一个会话任务。

Virtual host: 虚似地址,用于迸行逻辑隔离,是最上层的消息路由。一个 Virtual Host 里面可以有若干个 Exchange和 Queue ,同一个 VirtualHost 里面不能有相同名称的 Exchange 或 Queue。权限控制的最小粒度是Virtual Host。

Binding: Exchange 和 Queue 之间的虚拟连接,binding 中可以包含 routing key。

Routing key: 一 个路由规则,虚拟机可用它来确定如何路由一个特定消息,即交换机绑定到 Queue 的键。

Queue: 也称为Message Queue,消息队列,保存消息并将它们转发给消费者。

Message

消息,服务器和应用程序之间传送的数据,由 Properties 和 Body 组成。Properties 可以对消息进行修饰,比如消息的优先级、延迟等高级特性;,Body 则就 是消息体内容。

properties 中我们可以设置消息过期时间以及是否持久化等,也可以传入自定义的map属性,这些在消费端也都可以获取到。

生产者

1import com.rabbitmq.client.AMQP; 2import com.rabbitmq.client.Channel; 3import com.rabbitmq.client.Connection; 4import com.rabbitmq.client.ConnectionFactory; 5import java.util.HashMap; 6import java.util.Map; 7 8public class MessageProducer { 9 public static void main(String[] args) throws Exception { 10 //1. 创建一个 ConnectionFactory 并进行设置 11 ConnectionFactory factory = new ConnectionFactory(); 12 factory.setHost("localhost"); 13 factory.setVirtualHost("/"); 14 factory.setUsername("guest"); 15 factory.setPassword("guest"); 16 17 //2. 通过连接工厂来创建连接 18 Connection connection = factory.newConnection(); 19 20 //3. 通过 Connection 来创建 Channel 21 Channel channel = connection.createChannel(); 22 23 //4. 声明 使用默认交换机 以队列名作为 routing key 24 String queueName = "msg_queue"; 25 26 /** 27 * deliverMode 设置为 2 的时候代表持久化消息 28 * expiration 意思是设置消息的有效期,超过10秒没有被消费者接收后会被自动删除 29 * headers 自定义的一些属性 30 * */ 31 //5. 发送 32 Map<String, Object> headers = new HashMap<String, Object>(); 33 headers.put("myhead1", "111"); 34 headers.put("myhead2", "222"); 35 36 AMQP.BasicProperties properties = new AMQP.BasicProperties().builder() 37 .deliveryMode(2) 38 .contentEncoding("UTF-8") 39 .expiration("100000") 40 .headers(headers) 41 .build(); 42 String msg = "test message"; 43 channel.basicPublish("", queueName, properties, msg.getBytes()); 44 System.out.println("Send message : " + msg); 45 46 //6. 关闭连接 47 channel.close(); 48 connection.close(); 49 50 } 51}

消费者

1import com.rabbitmq.client.*; 2import java.io.IOException; 3import java.util.Map; 4 5public class MessageConsumer { 6 public static void main(String[] args) throws Exception{ 7 //1. 创建一个 ConnectionFactory 并进行设置 8 ConnectionFactory factory = new ConnectionFactory(); 9 factory.setHost("localhost"); 10 factory.setVirtualHost("/"); 11 factory.setUsername("guest"); 12 factory.setPassword("guest"); 13 factory.setAutomaticRecoveryEnabled(true); 14 factory.setNetworkRecoveryInterval(3000); 15 16 //2. 通过连接工厂来创建连接 17 Connection connection = factory.newConnection(); 18 19 //3. 通过 Connection 来创建 Channel 20 Channel channel = connection.createChannel(); 21 22 //4. 声明 23 String queueName = "msg_queue"; 24 channel.queueDeclare(queueName, false, false, false, null); 25 26 //5. 创建消费者并接收消息 27 Consumer consumer = new DefaultConsumer(channel) { 28 @Override 29 public void handleDelivery(String consumerTag, Envelope envelope, 30 AMQP.BasicProperties properties, byte[] body) 31 throws IOException { 32 String message = new String(body, "UTF-8"); 33 Map<String, Object> headers = properties.getHeaders(); 34 System.out.println("head: " + headers.get("myhead1")); 35 System.out.println(" [x] Received '" + message + "'"); 36 System.out.println("expiration : "+ properties.getExpiration()); 37 } 38 }; 39 40 //6. 设置 Channel 消费者绑定队列 41 channel.basicConsume(queueName, true, consumer); 42 } 43} 44 45 46Send message : test message 47 48head: 111 49 [x] Received 'test message' 50100000

Exchange

1. 简介

Exchange 就是交换机,接收消息,根据路由键转发消息到绑定的队列。有很多的 Message 进入到 Exchange 中,Exchange 根据 Routing key 将 Message 分发到不同的 Queue 中。

2. 类型

RabbitMQ中的 Exchange 有多种类型,类型不同,Message 的分发机制不同,如下:

  • fanout:广播模式。这种类型的 Exchange 会将 Message 分发到绑定到该 Exchange 的所有 Queue。

  • direct:这种类型的 Exchange 会根据 Routing key(精确匹配,将Message分发到指定的Queue。

  • Topic:这种类型的 Exchange 会根据 Routing key(模糊匹配,将Message分发到指定的Queue。

  • headers: 主题交换机有点相似,但是不同于主题交换机的路由是基于路由键,头交换机的路由值基于消息的header数据。主题交换机路由键只有是字符串,而头交换机可以是整型和哈希值 .

3. 属性

1/** 2* Declare an exchange, via an interface that allows the complete set of 3* arguments. 4* @see com.rabbitmq.client.AMQP.Exchange.Declare 5* @see com.rabbitmq.client.AMQP.Exchange.DeclareOk 6* @param exchange the name of the exchange 7* @param type the exchange type 8* @param durable true if we are declaring a durable exchange (the exchange will survive a server restart) 9* @param autoDelete true if the server should delete the exchange when it is no longer in use 10* @param internal true if the exchange is internal, i.e. can't be directly 11* published to by a client. 12* @param arguments other properties (construction arguments) for the exchange 13* @return a declaration-confirm method to indicate the exchange was successfully declared 14* @throws java.io.IOException if an error is encountered 15*/ 16Exchange.DeclareOk exchangeDeclare(String exchange, 17 String type,boolean durable, 18 boolean autoDelete,boolean internal, 19 Map<String, Object> arguments) throws IOException;
  • Name: 交换机名称

  • Type: 交换机类型direct、topic、 fanout、 headers

  • Durability: 是否需要持久化,true为持久化

  • Auto Delete: 当最后一个绑定到Exchange. 上的队列删除后,自动删除该Exchange

  • Internal: 当前Exchange是否用于RabbitMQ内部使用,默认为False

  • Arguments: 扩展参数,用于扩展AMQP协议自制定化使用

近期热文推荐:

1.终于靠开源项目弄到 IntelliJ IDEA 激活码了,真香!

2.我用 Java 8 写了一段逻辑,同事直呼看不懂,你试试看。。

3.吊打 Tomcat ,Undertow 性能很炸!!

4.国人开源了一款超好用的 Redis 客户端,真香!!

5.《Java开发手册(嵩山版)》最新发布,速速下载!

觉得不错,别忘了随手点赞+转发哦!

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

2020年前端实用代码段,为你的工作保驾护航

有空的时候,自己总结了几个代码段,在开发中也经常使用,谢谢。1、使用解构获取json数据let jsonData  id: 1,status: "OK",data: 'a', 'b';let  id, status, data: number   jsonData;console.log(id, status, number )