三、Exchanges(交换机)
让我们快速复习我们前面的教程::
一个队列是存储消息的缓冲区。
RabbitMQ的消息模型的核心思想是,生产者从未直接向队列发送任何消息。实际上,经常生产者甚至不知道消息是否会被运送到任何队列。
相反,生产者只能发送Exchanges(消息交换区)。交换是一个非常简单的事情。一方面它另一边队列。交换区必须知道如何处理它收到一条消息:
-
它应该被加到一个特定的队列吗?
-
它应该被加到多队列?
-
或者它应该丢弃吗?
X表示Exchange(交换机);
有一些可用的交换类型direct, topic, headers and fanout。我们将专注于最后一个——fanout。让我们创建一个这种类型的交换,称之为日志:
channel.exchangeDeclare("logs", "fanout");
问题:
列出所有 (交换机)列表
$ sudo rabbitmqctl list_exchanges
Listing exchanges ...
direct
amq.direct direct
amq.fanout fanout
amq.headers headers
amq.match headers
amq.rabbitmq.log topic
amq.rabbitmq.trace topic
amq.topic topic
logs fanout
...done.
在此列表中有一些amq* 交换器 与默认(匿名)交换。这些都是默认创建的,但可能你不需要使用它们。
② 缺省名字的 exchange(交换机)
channel.basicPublish("", "hello", null, message.getBytes());
routingKey存在的话,消息路由到指定的队列的名称。
channel.basicPublish( "logs", "", null, message.getBytes());
四、Temporary queues(临时队列)
当你想在生产者和消费者中分享队列的时候,给一个队列的名称是必须的。
但是那些都不是日志记录系统所需要的,我们希望能够获得所有的日志信息,而不只是其中的一部分,而且我们只对当前正在传递的信息感兴趣,
对旧的日志信息不感兴趣,要解决这些问题,我们需要分两个步骤:
- 首先当我们链接到RabbitMQ服务器的时候,需要一个新的、空的队列,为了做到这点,可以创建一个随机名的队列,
或者更好的方法就是让服务器选择一个随机的队列名。
- 其次,当断开与队列的连接时,消费者应该被自动删除掉。
在Java客户端,我们通过一个无参数的queueDeclare()方法为我们创建一个非持久的、唯一的、能自动删除的队列与队列名称
String queueName = channel.queueDeclare().getQueue(); 在这点上,queueName包含了一个随机队列名称。例如它可能看起来像amq.gen-JzTY20BRgKO-HjmUJj0wLg。
五、Bindings(绑定)
channel.queueBind(queueName, "logs", "");
生产者代码和之前的发送消息的代码并没有太大的区别,最重要的变化是,我们现在要将发布的消息传递给logs exchange来代替无名的exchange(之前的是""),
在发送消息时需要提供一个routingKey,它对于fanout exchange是非常重要的,不能被忽视的,这里的EmitLog.java代码如下
- </pre><pre name="code" class="java">import java.io.IOException;
- import com.rabbitmq.client.ConnectionFactory;
- import com.rabbitmq.client.Connection;
- import com.rabbitmq.client.Channel;
-
- public class EmitLog {
-
- private static final String EXCHANGE_NAME = "logs";
-
- public static void main(String[] argv)
- throws java.io.IOException {
-
- ConnectionFactory factory = new ConnectionFactory();
- factory.setHost("localhost");
- Connection connection = factory.newConnection();
- Channel channel = connection.createChannel();
-
- channel.exchangeDeclare(EXCHANGE_NAME, "fanout");
-
- String message = getMessage(argv);
-
- channel.basicPublish(EXCHANGE_NAME, "", null, message.getBytes());
- System.out.println(" [x] Sent ‘" + message + "‘");
-
- channel.close();
- connection.close();
- }
-
- }
接收端:
- import java.io.IOException;
- import com.rabbitmq.client.ConnectionFactory;
- import com.rabbitmq.client.Connection;
- import com.rabbitmq.client.Channel;
- import com.rabbitmq.client.QueueingConsumer;
-
- public class ReceiveLogs {
-
- private static final String EXCHANGE_NAME = "logs";
-
- public static void main(String[] argv)
- throws java.io.IOException,
- java.lang.InterruptedException {
-
- ConnectionFactory factory = new ConnectionFactory();
- factory.setHost("localhost");
- Connection connection = factory.newConnection();
- Channel channel = connection.createChannel();
-
- channel.exchangeDeclare(EXCHANGE_NAME, "fanout");
- String queueName = channel.queueDeclare().getQueue();
- channel.queueBind(queueName, EXCHANGE_NAME, "");
-
- System.out.println(" [*] Waiting for messages. To exit press CTRL+C");
-
- QueueingConsumer consumer = new QueueingConsumer(channel);
- channel.basicConsume(queueName, true, consumer);
-
- while (true) {
- QueueingConsumer.Delivery delivery = consumer.nextDelivery();
- String message = new String(delivery.getBody());
-
- System.out.println(" [x] Received ‘" + message + "‘");
- }
- }
- }
像以前一样,我们开始做编译