如今不少知名的互联网公司都有用到RabbitMQ,其性能,可扩展性让不少大公司青睐于使用它,不过想要彻底使用好RabbitMQ须要掌握其核心的一些概念,这里就说说掌握RabbitMQ所需的必要知识java
生产者: 建立消息,而后发送到代理服务器(RabbitMQ)的程序web
消费者:链接到代理服务器,并订阅到队列上接收消息shell
AMQP协议规定,AMQP消息必须有三部分,交换机,队列和绑定。生产者把消息发送到交换机,交换机与队列的绑定关系决定了消息如何路由到特定的队列,最终被消费者接收。编程
Note: 消息是不能直接到达队列(Queue)的api
消息实际上投递到的是交换机,具体路由到那个队列由交换机根据路由键(routing key)完成。bash
交换机在队列与消息中间起到了中间层的做用,有了交换机咱们能够实现更灵活的功能,RabbitMQ中有三种经常使用的交换机类型:服务器
*
能够替换一个单词#
能够替换全部的单词能够理解,direct为1v1, fanout为1v全部,topic比较灵活,能够1v任意。并发
每个虚拟主机(vhost)至关于mini版的RabbitMQ服务器,拥有本身的队列,交换机和绑定,权限... 这使得一个RabbitMQ服务众多的应用程序,而不会互相冲突。ide
rabbitMQ默认的虚拟主机为: "/" ,通常咱们在建立Rabbit的用户时会再给用户分配一个虚拟主机。高并发
操做虚拟主机,除了命令行以外还有一个web管理页面
#建立虚拟主机
rabbitmqctl add vhost [vhost_name]
#删除虚拟主机
rabbitmqctl delete vhost [vhost_name]
#列出虚拟主机
rabbitmqctl list_vhosts
复制代码
默认状况下RabbitMQ的队列和交换机在RabbitMQ服务器重启以后会消失,缘由在于队列和交换机的durable属性,该属性默认状况下为false.
能从AMQP服务器崩溃中恢复的消息称为持久化消息,若是想要从崩溃中恢复那么消息必须
缺点:消息写入磁盘性能差不少。除非特别关键的消息会使用
以上都是概念性的内容,实际咱们仍是要经过编程来实现咱们的目的,RabbitMQ的客户端api提供了不少功能,经过看代码,来了解它的强大之处。
基本步骤以前的RabbitMQ快速入门已经提过了,Channel类是关键的部分:包含了不少咱们想要的功能
生成端能够添加监听事件:
channel.addConfirmListener(new ConfirmListener() {
@Override
public void handleNack(long deliveryTag, boolean multiple) throws IOException {
System.err.println("-------no ack!-----------");
}
@Override
public void handleAck(long deliveryTag, boolean multiple) throws IOException {
System.err.println("-------ack!-----------");
}
});
复制代码
消费端能够确认消息状态:
public class MyConsumer extends DefaultConsumer {
private Channel channel ;
public MyConsumer(Channel channel) {
super(channel);
this.channel = channel;
}
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
System.err.println("-----------consume message----------");
System.err.println("body: " + new String(body));
try {
Thread.sleep(2000);
} catch (InterruptedException e) {
e.printStackTrace();
}
if((Integer)properties.getHeaders().get("num") == 0) {
channel.basicNack(envelope.getDeliveryTag(), false, true);
} else {
channel.basicAck(envelope.getDeliveryTag(), false);
}
}
}
复制代码
channel.basicAck与basicNack最后一个参数指定消息是否重回队列。
咱们的消息生产者经过指定交换机和路由键来把消息送到队列中,但有时候指定的路由键不存在,或者交换机不存在,那么消息就会return,咱们能够经过添加return listener来实现:
channel.addReturnListener(new ReturnListener() {
@Override
public void handleReturn(int replyCode, String replyText, String exchange, String routingKey, BasicProperties properties, byte[] body) throws IOException {
System.err.println("---------handle return----------");
System.err.println("replyCode: " + replyCode);
System.err.println("replyText: " + replyText);
System.err.println("exchange: " + exchange);
System.err.println("routingKey: " + routingKey);
System.err.println("properties: " + properties);
System.err.println("body: " + new String(body));
}
});
channel.basicPublish(exchange, routingKeyError, true, null, msg.getBytes());
复制代码
在basicPublish中的Mandatory要设置为true才会生效,不然broker会删除该消息
假设MQ服务器上面囤积了成千上万条的消息的时候,这个时候忽然链接消费端,那么巨量的消息所有推过来,可是客户端没法一次性处理这么多的数据。
在高并发的时候,瞬间产生的流量很大,消息很大,而MQ有个重要的做用就是限流,限流则是消费端作的。
RabbitMQ提供了一种Qos(服务质量保证)功能,即在非自动确认消息的前提下,在必定数量的消息未被消费前,不进行消费新的消息。
// prefetchSize消息的限制大小,通常设置为0,在生产端限制
// prefetchCount 咱们一次最多消费多少条消息,通常设置为1
// global,通常设置为false,在消费端进行限制
channel.basicQos(int prefetchSize, int prefetchCount, boolean global)
// 使用
channel.basicQos(0, 1, false);
channel.basicConsume(queueName, false, new MyConsumer(channel));
复制代码
Note: autoAck设置为false, 必定要手工签收消息
当消息在队列中变成死信,没有消费者进行消费的时候,消息可能会被从新发布到另一个队列中,这个队列就是死信队列。
如下状况会致使消息进入死信队列:
basic.reject/basic.nack 而且 requeue为false(不重回队列)的时候,消息就是死信
消息TTL过时
队列达到最大的长度
死信队列也是正常的Exchange,和通常的Exchange没什么区别,不过要作一点操做。
设置死信队列包括:
// 这就是一个普通的交换机 和 队列 以及路由
String exchangeName = "test_dlx_exchange";
String routingKey = "dlx.#";
String queueName = "test_dlx_queue";
channel.exchangeDeclare(exchangeName, "topic", true, false, null);
Map<String, Object> agruments = new HashMap<String, Object>();
agruments.put("x-dead-letter-exchange", "dlx.exchange");
//这个agruments属性,要设置到声明队列上
channel.queueDeclare(queueName, true, false, false, agruments);
channel.queueBind(queueName, exchangeName, routingKey);
//要进行死信队列的声明:
channel.exchangeDeclare("dlx.exchange", "topic", true, false, null);
channel.queueDeclare("dlx.queue", true, false, false, null);
channel.queueBind("dlx.queue", "dlx.exchange", "#");
复制代码
这里主要讲了一些使用RabbitMQ中常常涉及到的概念,懂了概念,在进行应用的时候才不至于糊涂。而后列举了MQ的Java客户端重要的几个API。