柯南君:看大数据时代下的IT架构(6)消息队列之RabbitMQ--案例(Publish/Subscribe起航)

一、回顾

让我们回顾一下,在上几章里都讲了什么?总结如下:




二、Publish/Subscribe(发布/订阅)(using the Java Client)

	在前面的教程中,我们创建了一个work Queue(工作队列)。工作队列背后的假设是每个任务是交付给一个工作者(worker) 也就是均匀分给每个消费者。在本部分,我们将做一些完全不同的事情,我们将提供一个消息到多个消费者。这种模式被称为“发布/订阅”。
	为了说明这个模式,我们将构建一个简单的日志系统。它将包括两个项目:
  1. 第一个将发出日志消息
  2. 第二个将接收并打印它们。
	在我们的日志系统,每运行一次,接收器项目将得到消息的副本。这样我们能够运行一个接收机并且可以直接记录到磁盘,同时我们可以运行另一个接收器,看到屏幕上的日志。
	注:从本质上讲,发表日志消息广播给所有的接收者。
	下面让我们脑中带几个问题,让我们一步一步去解决:
  • 如果我把消息分配给所有的消费者,我们将怎么做呢?

三、Exchanges(交换机)

在前部分的教程中,我们从一个队列发送和接收消息。现在是时候让Rabbit推出完整的消息模型。

让我们快速复习我们前面的教程::

  • 生产者是一个用户发送消息的应用程序。
  • 一个队列是存储消息的缓冲区。
  • 消费者是一个用户应用程序接收消息。

RabbitMQ的消息模型的核心思想是,生产者从未直接向队列发送任何消息。实际上,经常生产者甚至不知道消息是否会被运送到任何队列。

相反,生产者只能发送Exchanges(消息交换区)。交换是一个非常简单的事情。一方面它从生产者那收到消息并推他们到另一边队列。交换区必须知道如何处理它收到一条消息:

  1. 它应该被加到一个特定的队列吗?
  2. 它应该被加到多队列?
  3. 或者它应该丢弃吗?

交换的规则定义的类型。

如上图所示:X表示Exchange(交换机);

有一些可用的交换类型direct, topic, headers and fanout。我们将专注于最后一个——fanout。让我们创建一个这种类型的交换,称之为日志:

channel.exchangeDeclare("logs", "fanout");

fanout交换非常简单。你大概可以猜到的名字,只是广播所有的消息接收队列它知道。而这正是我们需要为我们的记录器。

问题:

exchange list

列出所有 (交换机)列表

$ 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(交换机)

在前部分的教程中我们对exchange 一无所知,,但仍然能够将消息发送到队列。这是可能的,因为我们是使用一个默认的交换,我们确定的空字符串(" ")

记得之前我们发布一个消息:

    channel.basicPublish("", "hello", null, message.getBytes());

第一个参数是该交换区的名称;空字符串表示默认或无名的交换,:如果routingKey存在的话,消息路由到指定的队列的名称。

现在,我们可以发布我们的交换器:

    channel.basicPublish( "logs", "", null, message.getBytes());

四、Temporary queues(临时队列)

	你可能记得以前我们使用的队列都是指定名称的(还记得hello和task_queue吗?)。对我们来说命名一个队列是至关重要的,
	当你想在生产者和消费者中分享队列的时候,给一个队列的名称是必须的。	
    但是那些都不是日志记录系统所需要的,我们希望能够获得所有的日志信息,而不只是其中的一部分,而且我们只对当前正在传递的信息感兴趣,
    对旧的日志信息不感兴趣,要解决这些问题,我们需要分两个步骤:
  • 首先当我们链接到RabbitMQ服务器的时候,需要一个新的、空的队列,为了做到这点,可以创建一个随机名的队列,
或者更好的方法就是让服务器选择一个随机的队列名。
  • 其次,当断开与队列的连接时,消费者应该被自动删除掉
在Java客户端,我们通过一个无参数的queueDeclare()方法为我们创建一个非持久的、唯一的、能自动删除的队列与队列名称
 String queueName = channel.queueDeclare().getQueue();

在这点上,queueName包含了一个随机队列名称。例如它可能看起来像amq.gen-JzTY20BRgKO-HjmUJj0wLg。

五、Bindings(绑定)

我们已经创建了一个fanout exchange和一个队列,现在我们需要告诉exchange去发送消息到队列中,exchange和队列之间的关系被称为一个绑定(binding)

	channel.queueBind(queueName, "logs", "");

注意:从现在开始我们从logs exchange将被添加消息到队列中,使用rabbitmqctl list_bingdins能列出所有的绑定。

六、Putting it all together(发布者/订阅者 实现)

生产者代码和之前的发送消息的代码并没有太大的区别,最重要的变化是,我们现在要将发布的消息传递给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 + "'");
        }
    }
}

像以前一样,我们开始做编译

$ javac -cp rabbitmq-client.jar EmitLog.java ReceiveLogs.java

如果你想将日志保存到一个文件,打开一个控制台并运行

$ java -cp .:commons-io-1.2.jar:commons-cli-1.1.jar:rabbitmq-client.jar ReceiveLogs > logs_from_rabbit.log

如果你想看到日志在你的屏幕上,产生一个新的终端并运行:

$ java -cp .:commons-io-1.2.jar:commons-cli-1.1.jar:rabbitmq-client.jar ReceiveLogs

发布日志类型:

$ java -cp .:commons-io-1.2.jar:commons-cli-1.1.jar:rabbitmq-client.jar EmitLog

使用rabbitmqctl list_bindings实际上您可以验证绑定和队列的代码是否是我们想要的? 有两个ReceiveLogs。

$ sudo rabbitmqctl list_bindings
Listing bindings ...
logs    exchange        amq.gen-JzTY20BRgKO-HjmUJj0wLg  queue           []
logs    exchange        amq.gen-vso0PVvyiRIL2WoV3i48Yg  queue           []
...done.

时间: 2024-08-07 08:37:36

柯南君:看大数据时代下的IT架构(6)消息队列之RabbitMQ--案例(Publish/Subscribe起航)的相关文章

柯南君:看大数据时代下的IT架构(5)消息队列之RabbitMQ--案例(Work Queues起航)

一.回顾 让我们回顾一下,在上几章里都讲了什么?总结如下: <柯南君:看大数据时代下的IT架构(1)业界消息队列对比> <柯南君:看大数据时代下的IT架构(2)消息队列之RabbitMQ-基础概念详细介绍> <柯南君:看大数据时代下的IT架构(3)消息队列之RabbitMQ-安装.配置与监控> <柯南君:看大数据时代下的IT架构(4)消息队列之RabbitMQ--案例(Helloword起航)> 二.Work Queues(using the Java Cl

柯南君:看大数据时代下的IT架构(4)消息队列之RabbitMQ--案例(Helloword起航)

一.回顾 让我们回顾一下,在上几章里都讲了什么?总结如下: <柯南君:看大数据时代下的IT架构(1)业界消息队列对比> <柯南君:看大数据时代下的IT架构(2)消息队列之RabbitMQ-基础概念详细介绍> <柯南君:看大数据时代下的IT架构(3)消息队列之RabbitMQ-安装.配置与监控> 二.起航 本章节,柯南君将从几个层面,用官网例子讲解一下RabbitMQ的实操经典程序案例,让大家重新回到经典"Hello world!"(The simpl

柯南君:看大数据时代下的IT架构(3)消息队列之RabbitMQ-安装、配置与监控

柯南君上一章<看大数据时代下的IT架构(2)消息队列之RabbitMQ-基础概念详细介绍>中,粗略的讲了一下,目前消息队列的几种常见产品的优劣对比,接下来的几章节会分别详细阐述,本章介绍RabbitMQ,好吧,废话少说,正式开始: 一.安装 1.安装Erlang 1)系统编译环境(这里采用linux/unix 环境) ① 安装环境 虚拟机:VMware? Workstation 10.0.1 build Linux系统:CentOS6.5 rabbitMQ官网下载:http://www.rab

柯南君:看大数据时代下的IT架构(2)消息队列之RabbitMQ-基础概念详细介绍

柯南君上一章<柯南君:看大数据时代下的IT架构(1)业界消息队列对比 >中,粗略的讲了一下,目前消息队列的几种常见产品的优劣对比,接下来的几章节会分别详细阐述,本章介绍RabbitMQ,好吧,废话少说,正式开始: 一.基础概念详细介绍 1.引言 你是否遇到过两个(多个)系统间需要通过定时任务来同步某些数据?你是否在为异构系统的不同进程间相互调用.通讯的问题而苦恼.挣扎?如果是,那么恭喜你,消息服务让你可以很轻松地解决这些问题. 消息服务擅长于解决多系统.异构系统间的数据交换(消息通知/通讯)问

柯南君:看大数据时代下的IT架构(1)业界消息队列对比

一.MQ(Message Queue) 即消息队列,一般用于应用系统解耦.消息异步分发,能够提高系统吞吐量.MQ的产品有很多,有开源的,也有闭源,比如ZeroMQ.RabbitMQ.ActiveMQ.Kafka/Jafka.Kestrel.Beanstalkd.HornetQ.Apache Qpid.Sparrow.Starling.Amazon SQS.MSMQ等,甚至Redis也可以用来构造消息队列.至于如何取舍,取决于你的需求. 由于工作需要和兴趣爱好,曾经写过关于RabbitMQ的系列博

看大数据时代下的IT架构(1)业界消息队列对比

一.MQ(Message Queue) 即消息队列,一般用于应用系统解耦.消息异步分发,能够提高系统吞吐量.MQ的产品有很多,有开源的,也有闭源,比如ZeroMQ.RabbitMQ.ActiveMQ.Kafka/Jafka.Kestrel.Beanstalkd.HornetQ.Apache Qpid.Sparrow.Starling.Amazon SQS.MSMQ等,甚至Redis也可以用来构造消息队列.至于如何取舍,取决于你的需求. 由于工作需要和兴趣爱好,曾经写过关于RabbitMQ的系列博

看大数据时代下的IT架构(1)图片服务器之演进史

        柯南君的公司最近产品即将上线,由于产品业务对图片的需求与日俱增,花样百出,与此同时,在大数据时代,大流量的冲击下,对图片服务器的压力可想而知,那么今天,柯南君结合互联网的相关热文,加上自己的一点实践经验,与君探讨,与君共勉! 一.图片服务器的重要性 当前,不管哪一家网站(包括 电商行业.O2O行业.互联网行业等),不管哪一种渠道 (包括 web端,APP端甚至一些SNS应用),在大数据时代下,在内容为王的前提下,对图片的需求量越来越大,柯南君的公司是一家O2O公司,也不例外,图片

柯南君:看大数据时代下的IT架构(9)消息队列之RabbitMQ--案例(RPC起航)

二.Remote procedure call (RPC)(using the Java client) 三.Client interface(客户端接口) 为了展示一个RPC服务是如何使用的,我们将创建一段很简单的客户端class. 它将会向外提供名字为call的函数,这个call会发送RPC请求并且阻塞,直到收到RPC运算的结果.代码如下: fibonacci_rpc = FibonacciRpcClient() result = fibonacci_rpc.call(4) print "f

柯南君:看大数据时代下的IT架构(7)消息队列之RabbitMQ--案例(routing 起航)

二.Routing(路由) (using the Java client) 在前面的学习中,构建了一个简单的日志记录系统,能够广播所有的日志给多个接收者,在该部分学习中,将添加一个新的特点,就是可以只订阅一个特定的消息源,也就是说能够直接把关键的错误日志消息发送到日志文件保存起来,不重要的日志信息文件不保存在磁盘中,但是仍然能够在控制台输出,那么这便是我们这部分要学习的消息的路由分发机制. 三.Bindings(绑定) 在前面的学习中已经创建了绑定(bindings),代码如下: channel