在上一章中,我們創建了一個工作隊列,工作隊列模式的設想是每一條消息只會被轉發給一個消費者。本章將會講解完全不一樣的場景: 我們會把一個消息轉發給多個消費者,這種模式稱之為發佈-訂閱模式。 為了闡述這個模式,我們將會搭建一個簡單的日誌系統,它包含兩種程式:一種發送日誌消息,另一種接收並列印日誌消息。在 ...
在上一章中,我們創建了一個工作隊列,工作隊列模式的設想是每一條消息只會被轉發給一個消費者。本章將會講解完全不一樣的場景: 我們會把一個消息轉發給多個消費者,這種模式稱之為發佈-訂閱模式。
為了闡述這個模式,我們將會搭建一個簡單的日誌系統,它包含兩種程式:一種發送日誌消息,另一種接收並列印日誌消息。在這個日誌系統里,每一個運行的消費者都可以獲取到消息,在這種情況下,我們可以實現這種需求:一個消費者接收消息並寫入磁碟,另一個消費者接收消息並列印在電腦屏幕上。簡單來說,生產者發佈的消息將會以廣播的形式轉發到所有的消費者。
1、交換器(Exchange)
在前兩章節我們,我們往隊列中發佈消息或獲取消息,然而,前面的講解其實並不完整,接下來,是時候介紹完整的RabbitMq消息模型了。
回憶一下我們前兩章指南中包含的內容:
-
- 一個生產者用以發送消息;
- 一個隊列緩存消息;
- 一個消費者用以消費隊列中的消息。
RabbitMq消息模式的核心思想是:一個生產者並不會直接往一個隊列中發送消息,事實上,生產者根本不知道它發送的消息將被轉發到哪些隊列。
實際上,生產者只能把消息發送給一個exchange,exchange只做一件簡單的事情:一方面它們接收從生產者發送過來的消息,另一方面,它們把接收到的消息推送給隊列。一個exchage必須清楚地知道如何處理一條消息。
有四種類型的交換器,分別是:direct、topic、headers、fanout。本章主要講解最後一種:fanous(廣播模式)。下麵創建一個fanout類型的交換器,我們稱之為:logs:
1 channel.exchangeDeclare("logs", "fanout");
廣播模式交換器很簡單,從字面意思也能理解,它其實就是把接收到的消息推送給所有它知道的隊列。在我們的日誌系統中正好需要這種模式。
如果想查看當前系統中有多少個exchange,可以使用以下命令:
sudo rabbitmqctl list_exchanges
或者通過控制台查看:
可以看到有很多以amq.*開頭的交換器,以及(AMQP default)預設交換器,這些是預設創建的交換器。
在前面兩章的指南中,我們並不知道交換器的存在,但是依然可以將消息發送到隊列中,那其實並不是因為我們可以不使用交換器,實際上是我們使用了預設的交換器(我們通過指定交換器為字字元串:""),回顧一下我們之前是如何發送消息的:
1 channel.basicPublish("", "hello", null, message.getBytes());
第一個參數是交換器的名字,空字元串表示它是一個預設或無命名的交換器,消息將會由指定的路由鍵(第二個參數,routingKey,後面會講)轉發到隊列。
你可能會有疑問:既然exchange可以指定為空字元串(""),那麼可否指定為null?
答案是:不能!
通過跟蹤發佈消息的代碼,在AMQImpl類中的Publish()方面中,可以看到,不光是exchange不能為null,同時routingKey路由鍵也不能為null,否則會拋出異常:
接著上面的講解,我們創建一個命名的交換器:
1 channel.basicPublish( "logs", "", null, message.getBytes());
2、臨時隊列
在前兩章的例子中,我們使用的隊列都是有具體的隊列名,創建命名隊列是很必要的,因為我們需要將消費者指向同一名字的隊列。因此,要想在生產者和消費者中間共用隊列就必須要使用命名隊列。
但是,本章講解的日誌系統也可以使用非命名隊列(可以不手動命名),我們希望收到所有日誌消息,而不是部分。並且我們希望總是接收到新的日誌消息而不是舊的日誌消息。為瞭解決這個問題,需要分兩步走。
首先,無論何時我們的消費者連接到RabbitMq,我們都需要一個新的、空的隊列來接收日誌消息,因此,消費者在連接上RabbitMq之後需要創建一個任意名字的隊列,或者讓RabbitMq生成任意的隊列名字。
其次,一旦該消費者斷開了與RabbitMq的連接,隊列也被自動刪除。
通過JAVA客戶端的無參方法:queueDeclare()來創建一個非持久化、專有的、自動刪除的、名字隨機生成的隊列。
1 String queueName = channel.queueDeclare().getQueue();
3、綁定(Binding)
前面廣播模式的交換器和隊列已經創建好了,接下來就是要告訴交換器向隊列里發送消息。交換器與隊列之間的關係稱之為綁定關係。
1 channel.queueBind(queueName, "logs", "");
至此,交換器已經可以往隊列中發送消息了。
可以通過下列命令來查看隊列的綁定關係:
4、完整的代碼
EmitLog.java
1 import com.rabbitmq.client.BuiltinExchangeType; 2 import com.rabbitmq.client.Channel; 3 import com.rabbitmq.client.Connection; 4 import com.rabbitmq.client.ConnectionFactory; 5 6 public class EmitLog { 7 8 private static final String EXCHANGE_NAME = "logs"; 9 10 public static void main(String[] args) throws Exception { 11 12 ConnectionFactory factory = new ConnectionFactory(); 13 factory.setHost("192.168.92.130"); 14 15 try (Connection connection = factory.newConnection(); 16 Channel channel = connection.createChannel();) { 17 18 channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.FANOUT); 19 20 String message = "RabbitMq fanout。。。。。。"; 21 channel.basicPublish(EXCHANGE_NAME,"",null,message.getBytes("utf-8")); 22 23 System.out.println(" [x] Sent '" + message + "'"); 24 } 25 } 26 }
正好你所看到的,Connection創建完成之後,定義了exchange,這一步是必要的,因為如果沒有交換器將無法發送消息。
如此沒有隊列綁定到該交換器上,那麼,交換器收到的消息將會丟失,但是對我們本章的日誌系統來說沒問題的,當沒有消費者時,我們可以安全地放棄掉數據,我們只接收最新的日誌消息。
ReceiveLogs.java
1 public class ReceiveLogs { 2 3 private static final String EXCHANGE_NAME = "logs"; 4 5 public static void main(String[] args) throws Exception { 6 7 ConnectionFactory factory = new ConnectionFactory(); 8 factory.setHost("192.168.92.130"); 9 10 Connection connection = factory.newConnection(); 11 Channel channel = connection.createChannel(); 12 13 channel.exchangeDeclare(EXCHANGE_NAME, BuiltinExchangeType.FANOUT); 14 15 final String queue = channel.queueDeclare().getQueue(); 16 channel.queueBind(queue,EXCHANGE_NAME,""); 17 18 System.out.println(" [*] Waiting for messages. To exit press CTRL+C"); 19 20 DeliverCallback deliverCallback = (consumerTa,delivery) -> { 21 22 String message = new String(delivery.getBody(), "UTF-8"); 23 System.out.println(" [x] Received '" + message + "'"); 24 25 }; 26 27 channel.basicConsume(queue,true,deliverCallback,consumerTag -> {}); 28 } 29 }
這裡的autoAck設置為true,因為我們這裡是廣播模式,每個消費者都會收到一樣的消息,並且這裡給消費者生產的隨機名稱的隊列相當於是獨有的,所以在接收到消息之後立即發送確認回執是OK的。
但是這裡先提出一個疑問:在這種模式下,每個隊列收到的消息是否也會有Ready和Unacked狀態?
5、測試結果
一、首先啟動生產者,再啟動兩個消費者
可以看到,生產者啟動後發送的消息丟失了,兩個消費者並沒有消費到,此時再看控制台:
可見RabbitMq為我們創建了兩個隨機命名的隊列,其Exclusive是Owner,表示是專有的,Parameters為AD(auto delete),擁有該隊列的消費者一占斷開連接,隊列將會被自動刪除。
二、其次啟動生產者發送一次消息
兩個消費都都收到了消息。
三、關閉所有消費者,觀察控制台變化
兩個專有隨機隊列自動刪除了。
6、SpringBoot的實現
工程結構圖:
一、配置文件application.properties:
生產者:
#RabbitMq spring.rabbitmq.host=192.168.92.130 spring.rabbitmq.exchange=logs
消費者:
#RabbitMq spring.rabbitmq.host=192.168.92.130 spring.rabbitmq.exchange=logs ##隊列--我們可以自己指定隊列名稱,也可以由RabbitMq自動生成,這裡為了方便,我們自己命名(如果需要,我也可以寫一個自動生成名稱的方法) rqbbitmq.log.fanout.info=info rqbbitmq.log.fanout.error=error
server.port=8090
二、生產者代碼
這裡為了讓系統生產者啟動時就自動發送一條消息,我加了一個EmitLogRunner類。
EmitLog.java
1 import org.springframework.amqp.core.AmqpTemplate; 2 import org.springframework.beans.factory.annotation.Autowired; 3 import org.springframework.beans.factory.annotation.Value; 4 import org.springframework.stereotype.Component; 5 6 @Component 7 public class EmitLog { 8 9 @Value("${spring.rabbitmq.exchange}") 10 private String exchange; 11 12 @Autowired 13 private AmqpTemplate amqpTemplate; 14 15 public void send(String msg) { 16 amqpTemplate.convertAndSend(exchange,"",msg); 17 } 18 }
EmitLogRunner.java
1 import org.springframework.beans.factory.annotation.Autowired; 2 import org.springframework.boot.ApplicationArguments; 3 import org.springframework.boot.ApplicationRunner; 4 import org.springframework.stereotype.Component; 5 6 @Component 7 public class EmitLogRunner implements ApplicationRunner { 8 9 @Autowired 10 private EmitLog emitLog; 11 12 @Override 13 public void run(ApplicationArguments args) throws Exception { 14 System.out.println("生產者發佈消息:" + msg); 15 emitLog.send("RabbitMq fanout test message"); 16 } 17 }
二、消費者代碼
ReceiveInfoLogs.java
1 @Component 2 @RabbitListener( 3 bindings = @QueueBinding( 4 value = @Queue(value = "${rqbbitmq.log.fanout.info}",autoDelete = "true"), 5 exchange = @Exchange(value = "${spring.rabbitmq.exchange}",type = ExchangeTypes.FANOUT) 6 ) 7 ) 8 public class ReceiveInfoLogs { 9 10 @Autowired 11 private AmqpTemplate amqpTemplate; 12 13 @RabbitHandler 14 public void receiveInfoLog (Object message) { 15 16 System.out.println("接收到info級別的日誌:" + message); 17 } 18 }
ReceiveErrorLogs.java
1 import org.springframework.amqp.core.AmqpTemplate; 2 import org.springframework.amqp.core.ExchangeTypes; 3 import org.springframework.amqp.rabbit.annotation.*; 4 import org.springframework.beans.factory.annotation.Autowired; 5 import org.springframework.stereotype.Component; 6 7 @Component 8 @RabbitListener( 9 bindings = @QueueBinding( 10 value = @Queue(value = "${rqbbitmq.log.fanout.error}",autoDelete = "true"), 11 exchange = @Exchange(value = "${spring.rabbitmq.exchange}",type = ExchangeTypes.FANOUT) 12 ) 13 ) 14 public class ReceiveErrorLogs { 15 16 @Autowired 17 private AmqpTemplate amqpTemplate; 18 19 @RabbitHandler 20 public void receiveErrorLog(Object message) { 21 System.out.println("接收到的error級別日誌:" + message); 22 } 23 }
註意看一下註解方式bindings裡面都是以@開頭並加上對應的要綁定的項,琢磨一下應該都能理解。
三、驗證
啟動消費者和生產者,查看控制台:
至此,發佈訂閱模式講解完了,在下一章中將會講解Routing(路由)的概念。