rabbitmq-----发布订阅模式

 模型组成

一个消费者Producer,一个交换机Exchange,多个消息队列Queue,多个消费者Consumer

一个生产者,多个消费者,每一个消费者都有自己的一个队列,生产者没有将消息直接发送到队列,而是发送到了交换机,每个队列绑定交换机,生产者发送的消息经过交换机,到达队列,实现一个消息被多个消费者获取的目的。需要注意的是,如果将消息发送到一个没有队列绑定的exchange上面,那么该消息将会丢失,这是因为在rabbitMQ中exchange不具备存储消息的能力,只有队列具备存储消息的能力。

Exchange

相比较于前两种模型Hello World和Work,这里多一个一个Exchange。其实Exchange是RabbitMQ的标配组成部件之一,前两种没有提到Exchange是为了简化模型,即使模型中没有看到Exchange的声明,其实还是声明了一个默认的Exchange。

RabbitMQ中实际发送消息并不是直接将消息发送给消息队列,消息队列也没那么聪明知道这条消息从哪来要到哪去。RabbitMQ会先将消息发送个Exchange,Exchange会根据这条消息打上的标记知道该条消息从哪来到哪去。

Exchange凭什么知道消息的何去何从,因为Exchange有几种类型:direct,fanout,topic和headers。这里说的订阅者模式就可以认为是fanout模式了。

RabbitMQ中,所有生产者提交的消息都由Exchange来接受,然后Exchange按照特定的策略转发到Queue进行存储 
RabbitMQ提供了四种Exchangefanout,direct,topic,header .发布/订阅模式就是是基于fanout Exchange实现的。

    • fanout这种模式不需要指定队列名称,需要将Exchangequeue绑定,他们之间的关系是‘多对多’的关系 
      任何发送到fanout Exchange的消息都会被转发到与该Exchange绑定的queue上面。

订阅者模式有何不同
订阅者模式相对前面的Work模式有和不同?Work也有多个消费者,但是只有一个消息队列,并且一个消息只会被某一个消费者消费。但是订阅者模式不一样,它有多个消息队列,也有多个消费者,而且一条消息可以被多个消费者消费,类似广播模式。下面通过实例代码看看这种模式是如何收发消息的。

 package com.maozw.mq.pubsub;

 import com.maozw.mq.config.RabbitConfig;
import com.rabbitmq.client.Channel;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.amqp.rabbit.connection.Connection;
import org.springframework.amqp.rabbit.connection.ConnectionFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController; import java.io.IOException;
import java.util.concurrent.TimeoutException; import static org.apache.log4j.varia.ExternallyRolledFileAppender.OK; /**
* work 模式
* 两种分发: 轮询分发 + 公平分发
* 轮询分发:消费端:自动确认消息;boolean autoAck = true;
* 公平分发: 消费端:手动确认消息 boolean autoAck = false; channel.basicAck(envelope.getDeliveryTag(),false);
*
* @author MAOZW
* @Description: ${todo}
* @date 2018/11/26 15:06
*/
@RestController
@RequestMapping("/publish")
public class PublishProducer {
private static final Logger LOGGER = LoggerFactory.getLogger(PublishProducer.class);
@Autowired
RabbitConfig rabbitConfig; @RequestMapping("/send/{exchangeName}/{queueName}")
public String send(@PathVariable String exchangeName, @PathVariable String queueName) throws IOException, TimeoutException {
Connection connection = null;
Channel channel= null;
try {
ConnectionFactory connectionFactory = rabbitConfig.connectionFactory();
connection = connectionFactory.createConnection();
channel = connection.createChannel(false); /**
* 申明交换机
*/
channel.exchangeDeclare(exchangeName,"fanout"); /**
* 发送消息
* 每个消费者 发送确认消息之前,消息队列不会发送下一个消息给消费者,一次只处理一个消息
* 自动模式无需设置下面设置
*/
int prefetchCount = ;
channel.basicQos(prefetchCount); String Hello = ">>>> Hello Simple <<<<";
for (int i = ; i < ; i++) {
String message = Hello + i;
channel.basicPublish(RabbitConfig.EXCHANGE_AAAAA, "", null, message.getBytes());
LOGGER.info("生产消息: " + message);
}
return "OK";
}catch (Exception e) { } finally {
connection.close();
channel.close();
return OK;
}
}
}

订阅1

 package com.maozw.mq.pubsub;

 import com.maozw.mq.config.RabbitConfig;
import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.DefaultConsumer;
import com.rabbitmq.client.Envelope;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.amqp.rabbit.connection.Connection;
import org.springframework.amqp.rabbit.connection.ConnectionFactory; import java.io.IOException; /**
* @author MAOZW
* @Description: ${todo}
* @date 2018/11/26 15:06
*/ public class SubscribeConsumer {
private static final Logger LOGGER = LoggerFactory.getLogger(SubscribeConsumer.class); public static void main(String[] args) throws IOException {
ConnectionFactory connectionFactory = RabbitConfig.getConnectionFactory();
Connection connection = connectionFactory.createConnection();
Channel channel = connection.createChannel(false);
/**
* 创建队列申明
*/
boolean durable = true;
channel.queueDeclare(RabbitConfig.QUEUE_PUBSUB_FANOUT, durable, false, false, null);
/**
* 绑定队列到交换机
*/
channel.queueBind(RabbitConfig.QUEUE_PUBSUB_FANOUT, RabbitConfig.EXCHANGE_AAAAA,""); /**
* 改变分发规则
*/
channel.basicQos(1);
DefaultConsumer consumer = new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
super.handleDelivery(consumerTag, envelope, properties, body);
System.out.println("[2] 接口数据 : " + new String(body, "utf-8"));
try {
Thread.sleep(300);
} catch (InterruptedException e) {
e.printStackTrace();
} finally {
System.out.println("[2] done!");
//消息应答:手动回执,手动确认消息
channel.basicAck(envelope.getDeliveryTag(),false);
}
}
};
//监听队列
/**
* autoAck 消息应答
* 默认轮询分发打开:true :这种模式一旦rabbitmq将消息发送给消费者,就会从内存中删除该消息,不关心客户端是否消费正常。
* 使用公平分发需要关闭autoAck:false 需要手动发送回执
*/
boolean autoAck = false;
channel.basicConsume(RabbitConfig.QUEUE_PUBSUB_FANOUT,autoAck, consumer);
} }
 package com.maozw.mq.pubsub;

 import com.maozw.mq.config.RabbitConfig;
import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.DefaultConsumer;
import com.rabbitmq.client.Envelope;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.amqp.rabbit.connection.Connection;
import org.springframework.amqp.rabbit.connection.ConnectionFactory; import java.io.IOException; /**
* @author MAOZW
* @Description: ${todo}
* @date 2018/11/26 15:06
*/ public class SubscribeConsumer2 {
private static final Logger LOGGER = LoggerFactory.getLogger(SubscribeConsumer2.class); public static void main(String[] args) throws IOException {
ConnectionFactory connectionFactory = RabbitConfig.getConnectionFactory();
Connection connection = connectionFactory.createConnection();
Channel channel = connection.createChannel(false);
/**
* 创建队列申明
*/
boolean durable = true;
channel.queueDeclare(RabbitConfig.QUEUE_PUBSUB_FANOUT2, durable, false, false, null);
/**
* 绑定队列到交换机
*/
channel.queueBind(RabbitConfig.QUEUE_PUBSUB_FANOUT2, RabbitConfig.EXCHANGE_AAAAA,""); /**
* 改变分发规则
*/
channel.basicQos(1);
DefaultConsumer consumer = new DefaultConsumer(channel) {
@Override
public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
super.handleDelivery(consumerTag, envelope, properties, body);
System.out.println("[2] 接口数据 : " + new String(body, "utf-8"));
try {
Thread.sleep(400);
} catch (InterruptedException e) {
e.printStackTrace();
} finally {
System.out.println("[2] done!");
//消息应答:手动回执,手动确认消息
channel.basicAck(envelope.getDeliveryTag(),false);
}
}
};
//监听队列
/**
* autoAck 消息应答
* 默认轮询分发打开:true :这种模式一旦rabbitmq将消息发送给消费者,就会从内存中删除该消息,不关心客户端是否消费正常。
* 使用公平分发需要关闭autoAck:false 需要手动发送回执
*/
boolean autoAck = false;
channel.basicConsume(RabbitConfig.QUEUE_PUBSUB_FANOUT2,autoAck, consumer);
}
}

最新文章

  1. Reporting Services 错误案例一则
  2. ASP.NET页面的字符编码设置
  3. JQuery制作简单的网页导航特效
  4. Unity3d5.0 新UI之2048
  5. MVC中html转义问题(直接输出html的方法)
  6. Servlet里写验证码
  7. maxscript, 数组和字符串下标是从1开始的
  8. CSS3动画制作的简单示例
  9. Laravel 5.1使用命令行模式(artisan)运行php脚本
  10. 很近没读书了,读书笔记之&lt;&lt;大道至简&gt;&gt;
  11. gradle使用国内源
  12. CentOS6.5 安装mysql5.6.30
  13. javaMail邮件发送的简单实现
  14. how tomcat works 读书笔记 十一 StandWrapper 下
  15. 以语音评测的PC端demo代码为例,讲解口语评测如何实现
  16. Rich feature hierarchies for accurate object detection and semantic segmentation(理解)
  17. idea: Unable to parse template &quot;class&quot;
  18. Js高级 部分内容 面向对象
  19. ubuntu16.04与mysql的运维注意事项
  20. linux 添加交换分区

热门文章

  1. hdu 2609 字符串最小表示法 虽然不是很懂 还是先贴上来吧。/,。/
  2. C# WebForm 屏蔽输入框的验证
  3. JS的一些简单基础运算题
  4. npm install 报错踩坑路
  5. “最不合格”的SAP应聘者: 从大学生到SAP成都研究院开发工程师
  6. sqlserver 2008修改数据库表的时候错误提示“阻止保存要求重新创建表的更改”
  7. 使用Fiddler工具在夜神模拟器或手机上抓包
  8. RobHess的SIFT代码解析之kd树
  9. 【2017-06-02】Linq高级查询,实现分页组合查询。
  10. 自定义控件之Canvas图形绘制基础练习-青春痘笑脸^_^