3.rabbitmq

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也有多个消费者,但是只有一个消息队列,并且一个消息只会被某一个消费者消费。但是订阅者模式不一样,它有多个消息队列,也有多个消费者,而且一条消息可以被多个消费者消费,类似广播模式。下面通过实例代码看看这种模式是如何收发消息的。

1 1 package com.maozw.mq.pubsub; 2 2 3 3 import com.maozw.mq.config.RabbitConfig; 4 4 import com.rabbitmq.client.Channel; 5 5 import org.slf4j.Logger; 6 6 import org.slf4j.LoggerFactory; 7 7 import org.springframework.amqp.rabbit.connection.Connection; 8 8 import org.springframework.amqp.rabbit.connection.ConnectionFactory; 9 9 import org.springframework.beans.factory.annotation.Autowired; 1010 import org.springframework.web.bind.annotation.PathVariable; 1111 import org.springframework.web.bind.annotation.RequestMapping; 1212 import org.springframework.web.bind.annotation.RestController; 1313 1414 import java.io.IOException; 1515 import java.util.concurrent.TimeoutException; 1616 1717 import static org.apache.log4j.varia.ExternallyRolledFileAppender.OK; 1818 1919 /** 2020 * work 模式 2121 * 两种分发: 轮询分发 + 公平分发 2222 * 轮询分发:消费端:自动确认消息;boolean autoAck = true; 2323 * 公平分发: 消费端:手动确认消息 boolean autoAck = false; channel.basicAck(envelope.getDeliveryTag(),false); 2424 * 2525 * @author MAOZW 2626 * @Description: ${todo} 2727 * @date 2018/11/26 15:06 2828 */ 2929 @RestController 3030 @RequestMapping("/publish") 3131 public class PublishProducer { 3232 private static final Logger LOGGER = LoggerFactory.getLogger(PublishProducer.class); 3333 @Autowired 3434 RabbitConfig rabbitConfig; 3535 3636 3737 @RequestMapping("/send/{exchangeName}/{queueName}") 3838 public String send(@PathVariable String exchangeName, @PathVariable String queueName) throws IOException, TimeoutException { 3939 Connection connection = null; 4040 Channel channel= null; 4141 try { 4242 ConnectionFactory connectionFactory = rabbitConfig.connectionFactory(); 4343 connection = connectionFactory.createConnection(); 4444 channel = connection.createChannel(false); 4545 4646 /** 4747 * 申明交换机 4848 */ 4949 channel.exchangeDeclare(exchangeName,"fanout"); 5050 5151 /** 5252 * 发送消息 5353 * 每个消费者 发送确认消息之前,消息队列不会发送下一个消息给消费者,一次只处理一个消息 5454 * 自动模式无需设置下面设置 5555 */ 5656 int prefetchCount = 1; 5757 channel.basicQos(prefetchCount); 5858 5959 String Hello = ">>>> Hello Simple <<<<"; 6060 for (int i = 0; i < 5; i++) { 6161 String message = Hello + i; 6262 channel.basicPublish(RabbitConfig.EXCHANGE_AAAAA, "", null, message.getBytes()); 6363 LOGGER.info("生产消息: " + message); 6464 } 6565 return "OK"; 6666 }catch (Exception e) { 6767 6868 } finally { 6969 connection.close(); 7070 channel.close(); 7171 return OK; 7272 } 7373 } 7474 }

 订阅1

1 1 package com.maozw.mq.pubsub; 2 2 3 3 import com.maozw.mq.config.RabbitConfig; 4 4 import com.rabbitmq.client.AMQP; 5 5 import com.rabbitmq.client.Channel; 6 6 import com.rabbitmq.client.DefaultConsumer; 7 7 import com.rabbitmq.client.Envelope; 8 8 import org.slf4j.Logger; 9 9 import org.slf4j.LoggerFactory; 1010 import org.springframework.amqp.rabbit.connection.Connection; 1111 import org.springframework.amqp.rabbit.connection.ConnectionFactory; 1212 1313 import java.io.IOException; 1414 1515 /** 1616 * @author MAOZW 1717 * @Description: ${todo} 1818 * @date 2018/11/26 15:06 1919 */ 2020 2121 public class SubscribeConsumer { 2222 private static final Logger LOGGER = LoggerFactory.getLogger(SubscribeConsumer.class); 2323 2424 public static void main(String[] args) throws IOException { 2525 ConnectionFactory connectionFactory = RabbitConfig.getConnectionFactory(); 2626 Connection connection = connectionFactory.createConnection(); 2727 Channel channel = connection.createChannel(false); 2828 /** 2929 * 创建队列申明 3030 */ 3131 boolean durable = true; 3232 channel.queueDeclare(RabbitConfig.QUEUE_PUBSUB_FANOUT, durable, false, false, null); 3333 /** 3434 * 绑定队列到交换机 3535 */ 3636 channel.queueBind(RabbitConfig.QUEUE_PUBSUB_FANOUT, RabbitConfig.EXCHANGE_AAAAA,""); 3737 3838 /** 3939 * 改变分发规则 4040 */ 4141 channel.basicQos(1); 4242 DefaultConsumer consumer = new DefaultConsumer(channel) { 4343 @Override 4444 public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException { 4545 super.handleDelivery(consumerTag, envelope, properties, body); 4646 System.out.println("[2] 接口数据 : " + new String(body, "utf-8")); 4747 try { 4848 Thread.sleep(300); 4949 } catch (InterruptedException e) { 5050 e.printStackTrace(); 5151 } finally { 5252 System.out.println("[2] done!"); 5353 //消息应答:手动回执,手动确认消息 5454 channel.basicAck(envelope.getDeliveryTag(),false); 5555 } 5656 } 5757 }; 5858 //监听队列 5959 /** 6060 * autoAck 消息应答 6161 * 默认轮询分发打开:true :这种模式一旦rabbitmq将消息发送给消费者,就会从内存中删除该消息,不关心客户端是否消费正常。 6262 * 使用公平分发需要关闭autoAck:false 需要手动发送回执 6363 */ 6464 boolean autoAck = false; 6565 channel.basicConsume(RabbitConfig.QUEUE_PUBSUB_FANOUT,autoAck, consumer); 6666 } 6767 6868 } 69 70 1 package com.maozw.mq.pubsub; 71 2 72 3 import com.maozw.mq.config.RabbitConfig; 73 4 import com.rabbitmq.client.AMQP; 74 5 import com.rabbitmq.client.Channel; 75 6 import com.rabbitmq.client.DefaultConsumer; 76 7 import com.rabbitmq.client.Envelope; 77 8 import org.slf4j.Logger; 78 9 import org.slf4j.LoggerFactory; 7910 import org.springframework.amqp.rabbit.connection.Connection; 8011 import org.springframework.amqp.rabbit.connection.ConnectionFactory; 8112 8213 import java.io.IOException; 8314 8415 /** 8516 * @author MAOZW 8617 * @Description: ${todo} 8718 * @date 2018/11/26 15:06 8819 */ 8920 9021 public class SubscribeConsumer2 { 9122 private static final Logger LOGGER = LoggerFactory.getLogger(SubscribeConsumer2.class); 9223 9324 public static void main(String[] args) throws IOException { 9425 ConnectionFactory connectionFactory = RabbitConfig.getConnectionFactory(); 9526 Connection connection = connectionFactory.createConnection(); 9627 Channel channel = connection.createChannel(false); 9728 /** 9829 * 创建队列申明 9930 */ 10031 boolean durable = true; 10132 channel.queueDeclare(RabbitConfig.QUEUE_PUBSUB_FANOUT2, durable, false, false, null); 10233 /** 10334 * 绑定队列到交换机 10435 */ 10536 channel.queueBind(RabbitConfig.QUEUE_PUBSUB_FANOUT2, RabbitConfig.EXCHANGE_AAAAA,""); 10637 10738 /** 10839 * 改变分发规则 10940 */ 11041 channel.basicQos(1); 11142 DefaultConsumer consumer = new DefaultConsumer(channel) { 11243 @Override 11344 public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException { 11445 super.handleDelivery(consumerTag, envelope, properties, body); 11546 System.out.println("[2] 接口数据 : " + new String(body, "utf-8")); 11647 try { 11748 Thread.sleep(400); 11849 } catch (InterruptedException e) { 11950 e.printStackTrace(); 12051 } finally { 12152 System.out.println("[2] done!"); 12253 //消息应答:手动回执,手动确认消息 12354 channel.basicAck(envelope.getDeliveryTag(),false); 12455 } 12556 } 12657 }; 12758 //监听队列 12859 /** 12960 * autoAck 消息应答 13061 * 默认轮询分发打开:true :这种模式一旦rabbitmq将消息发送给消费者,就会从内存中删除该消息,不关心客户端是否消费正常。 13162 * 使用公平分发需要关闭autoAck:false 需要手动发送回执 13263 */ 13364 boolean autoAck = false; 13465 channel.basicConsume(RabbitConfig.QUEUE_PUBSUB_FANOUT2,autoAck, consumer); 13566 } 13667 }
点赞
收藏

评论区

加载中...

相关推荐

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(

手写Java HashMap源码

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

Kafka安装步骤

基本概念1.Producer:消息生产者,就是向kafkabroker发消息的客户端2.Consumer:消息消费者,向kafkabroker取消息的客户端3.ConsumerGroup(CG):消费者组,由多个consumer组成。消费者组内每个消费者负责消费不同分区的数据,一个分区只能由一

RabbitMQ操作

注意:在rabbitmq中,可以存在多个exchange,exchange只是负责接收消息,然后消息必须发送到给queue中,如果没有queue,消息就丢失了,exchange就相当于交换机,不负责存消息,queue是必须声明的,所以exchange负责转发,queue负责接收!(https://oscimg.oschina.net/oscnet/1

00:Java简单了解

浅谈Java之概述Java是SUN(StanfordUniversityNetwork),斯坦福大学网络公司)1995年推出的一门高级编程语言。Java是一种面向Internet的编程语言。随着Java技术在web方面的不断成熟,已经成为Web应用程序的首选开发语言。Java是简单易学,完全面向对象,安全可靠,与平台无关的编程语言。