===
信息发布与订阅
Rabbit的核心组件包含Queue(消息队列)和Exchanges两部分,Exchange的主要部分就是对信息进行路由,通过将消息队列绑定到Exchange上,则可以实现订阅形式的消息发布及Publish/Subscribe在这种模式下消息发布者只需要将信息发布到相应的Exchange中,而Exchange则自动将信息分发到不同的Queue当中。
这种模式下Exchange充当的角色
在命令行中可以使用
sudo rabbitmqctl list_exchanges
sudo rabbitmqctl list_bindings
分别查看当前系统种存在的Exchange和Exchange上绑定的Queue信息。
消息发布者EmitLog.java
1import com.rabbitmq.client.Channel; 2import com.rabbitmq.client.Connection; 3import com.rabbitmq.client.ConnectionFactory; 4 5public class EmitLog { 6 7 private static final String EXCHANGE_NAME="logs"; 8 9 public static void main(String[] args) throws java.io.IOException{ 10 11 //创建链接工厂 12 ConnectionFactory factory = new ConnectionFactory(); 13 factory.setHost("localhost"); 14 //创建链接 15 Connection connection = factory.newConnection(); 16 17 //创建信息管道 18 Channel channel = connection.createChannel(); 19 20 //生命Exchange 非持久化 21 channel.exchangeDeclare(EXCHANGE_NAME, "fanout"); 22 23 String message = "Message "+Math.random(); 24 25 //第一个参数是对应的Exchange名称,如果为空则使用默认Exchange 26 channel.basicPublish(EXCHANGE_NAME, "", null, message.getBytes()); 27 System.out.println("[x] Sent '"+message+"'"); 28 29 //关闭链接 30 channel.close(); 31 connection.close(); 32 33 } 34 35}
消息消费者ReceiveLogs.java
1import java.io.IOException; 2 3import com.rabbitmq.client.Channel; 4import com.rabbitmq.client.Connection; 5import com.rabbitmq.client.ConnectionFactory; 6import com.rabbitmq.client.ConsumerCancelledException; 7import com.rabbitmq.client.QueueingConsumer; 8import com.rabbitmq.client.ShutdownSignalException; 9 10public class ReceiveLogs { 11 12 private static final String EXCHANGE_NAME = "logs"; 13 14 public static void main(String[] args) throws IOException, ShutdownSignalException, ConsumerCancelledException, InterruptedException { 15 16 //创建链接工厂 17 ConnectionFactory factory = new ConnectionFactory(); 18 factory.setHost("localhost"); 19 //创建链接 20 Connection connection = factory.newConnection(); 21 22 //创建消息管道 23 Channel channel = connection.createChannel(); 24 25 //声明Exchange 26 channel.exchangeDeclare(EXCHANGE_NAME, "fanout"); 27 28 //利用系统自动声明一个非持久化的消息队列,并返回唯一的队列名称 29 String queueName = channel.queueDeclare().getQueue(); 30 31 //将消息队列绑定到Exchange 32 channel.queueBind(queueName, EXCHANGE_NAME, ""); 33 34 System.out.println(" [*] Waiting for messages. To exit press CTRL+C"); 35 36 //声明一个消费者 37 QueueingConsumer consumer = new QueueingConsumer(channel); 38 channel.basicConsume(queueName, true, consumer); 39 40 while (true) { 41 42 //循环获取信息 43 QueueingConsumer.Delivery delivery = consumer.nextDelivery(); 44 String message = new String(delivery.getBody()); 45 System.out.println(" [x] Received '" + message + "'"); 46 47 } 48 49 } 50 51}
运行时启动一个EmitLog.java多个ReceiveLogs.java则可以看到发布者每次发布信息,只要绑定到了相应Exchange的消费者都可以获取到信息。
RabbitMQ信息持久化技术
上面的例子中我们实现了Publisher/Subscribe的消息分发方式,但是其中存在一些问题。比如当我们运行一个ReceiveLog都对应了一个特定的消息队列,可以利用list_queues进行查看,同时这些消息队列是帮到到名为logs的Exchange中,这是发布消息每个消费者都可以接收到,可以当关闭ReceiveLog程序后这些消息队列就都会自动销毁,因为他们是非持久化的。同样对于EmitLog程序也一样,每次关闭后之前生命的Exchange也将自动销毁。
这就产生了一些问题。如果当ReceiveLog为运行时,此时就并没有一个消息队列是绑定到Exchange上的,在发布消息后再启动ReceiveLog程序是无法接受到之前发布的信息。这就是为什么要进行消息的持久化。
通过持久化技术,我们可以生命一个持久化的Exchange,以及持久化的Queue这样,在把Queue绑定到Exchange后,即使没有消费者程序运行,信息依然能保存在Queue当中,当下次启动消费者程序时依然能获取到发布的所有信息。就好比当一个消费者程序在执行消息序列中的任务时,如果突然出现了异常那么重新启动后,依然能从上一次发生错误的位置继续运行,对于某些需要一个有序性和连续性的操作,这点显的尤为重要。
下面还是给出一个例子,在持久化过程中,可以借助list_exchanges,list_bindings,list_queues来查看服务器中相关信息来帮组分析过程。
Publisher.java
1import com.rabbitmq.client.Channel; 2import com.rabbitmq.client.Connection; 3import com.rabbitmq.client.ConnectionFactory; 4import com.rabbitmq.client.MessageProperties; 5 6public class Publisher { 7 8 private static final String EXCHANGE_NAME="persi";//定义Exchange名称 9 private static final boolean durable = true;//消息队列持久化 10 11 public static void main(String[] args) throws java.io.IOException { 12 13 ConnectionFactory factory = new ConnectionFactory();//创建链接工厂 14 factory.setHost("localhost"); 15 Connection connection = factory.newConnection();//创建链接 16 Channel channel = connection.createChannel();//创建信息通道 17 18 channel.exchangeDeclare(EXCHANGE_NAME, "fanout", durable);//创建交换机并生命持久化 19 20 String message = "Hello Wrold "+Math.random(); 21 //消息的持久化 22 channel.basicPublish(EXCHANGE_NAME, "", MessageProperties.PERSISTENT_TEXT_PLAIN, message.getBytes()); 23 24 System.out.println("[x] Sent '" + message + "'"); 25 26 channel.close(); 27 connection.close(); 28 29 } 30 31}
Subscriber.java
1public class Subscriber { 2 3 4 //private static final String[] QUEUE_NAMES= {"que_001","que_002","que_003","que_004","que_005"}; 5 private static final String[] QUEUE_NAMES= {"que_006","que_007","que_008","que_009","que_0010"}; 6 7 public static void main(String[] args){ 8 9 for(int i=0;i<QUEUE_NAMES.length;i++){ 10 11 SubscriberThead sub = new SubscriberThead(QUEUE_NAMES[i]); 12 Thread t = new Thread(sub); 13 t.start(); 14 15 } 16 17 } 18}
SubscriberThead.java
1import com.rabbitmq.client.Channel; 2import com.rabbitmq.client.Connection; 3import com.rabbitmq.client.ConnectionFactory; 4import com.rabbitmq.client.QueueingConsumer; 5import com.rabbitmq.client.AMQP.Queue.DeclareOk; 6 7public class SubscriberThead implements Runnable { 8 9 private String queue_name = null; 10 private static final String EXCHANGE_NAME = "persi";// 定义交换机名称 11 private static final boolean durable = true;//消息队列持久化 12 13 public SubscriberThead(String queue_name) { 14 15 this.queue_name = queue_name; 16 17 } 18 19 @Override 20 public void run() { 21 22 try{ 23 24 ConnectionFactory factory = new ConnectionFactory(); 25 factory.setHost("localhost"); 26 Connection connection = factory.newConnection(); 27 Channel channel = connection.createChannel(); 28 29 channel.exchangeDeclare(EXCHANGE_NAME, "fanout", durable); 30 31 DeclareOk ok = channel.queueDeclare(queue_name, durable, false, 32 false, null); 33 String queueName = ok.getQueue(); 34 35 36 channel.queueBind(queueName, EXCHANGE_NAME, ""); 37 38 System.out.println(" ["+queue_name+"] Waiting for messages. To exit press CTRL+C"); 39 40 channel.basicQos(1);//消息分发处理 41 QueueingConsumer consumer = new QueueingConsumer(channel); 42 channel.basicConsume(queueName, false, consumer); 43 44 while (true) { 45 46 Thread.sleep(2000); 47 QueueingConsumer.Delivery delivery = consumer.nextDelivery(); 48 String message = new String(delivery.getBody()); 49 System.out.println(" ["+queue_name+"] Received '" + message + "'"); 50 channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false); 51 52 } 53 }catch(Exception e){ 54 55 e.printStackTrace(); 56 } 57 58 59 } 60 61}
通过持久化处理后rabbitMQ将保存Exchange信息以及Queue信息,甚至在rabbitMQ服务器关闭后信息依然能保存,这样就提供了消息传递的可靠性