RabbitMQ - RPC in Java

这次试着用RabbitMQ进行RPC。

其实用RabbitMQ搞RPC也没什么特别的。

只是我们需要在请求中再加入一个callback queue。

比如这样:

callbackQueueName = channel.queueDeclare().getQueue();

BasicProperties props = new BasicProperties
                            .Builder()
                            .replyTo(callbackQueueName)
                            .build();

channel.basicPublish("", "rpc_queue", props, message.getBytes());

剩下的工作就是等待对方处理完成再从callback队列中读取响应消息。

上面用到了BasicProperties。

(注意:是com.rabbitmq.client.AMQP.BasicProperties 不是 com.rabbitmq.client.BasicProperties)

关于Message properties,AMQP协议为消息预定义了14种属性。

        private String contentType;
        private String contentEncoding;
        private Map<String,Object> headers;
        private Integer deliveryMode;
        private Integer priority;
        private String correlationId;
        private String replyTo;
        private String expiration;
        private String messageId;
        private Date timestamp;
        private String type;
        private String userId;
        private String appId;
        private String clusterId;

通常我们只需要使用其中一小部分:

·deliveryMode: 将消息设置为持久或者临时,2为持久,其余为临时。

·contentType: 指定mime-type,比如要使用JSON就是application/json

·replyTo: 指定callback queue的名字

·correlationId: 用来关联RPC请求和响应的标识。

上面那段代码中就是用到了correlationId。

另外需要说明这个correlationId。

其实在上面的代码中我们为每一个RPC请求都创建了一个回调队列。

但这样明显不效率,我们可以为每一个客户端只创建一个回调队列。

但这样我们又需要考虑另一个问题:<当我们将收到的消息放到队列时,如何确定该消息是属于哪个请求?>

这时我们可以使用correlationId解决这个问题。

我们可以用它来为每一个请求加上标识,获取信息时对比这个标识,以对应请求和响应。

如果我们收到了无法识别的correlationId,即该响应不与任何请求匹配,那么这个消息将会废除。

好了,代码比较简单。

RPCServer:

class RPCServer{
    private static final String RPC_QUEUE_NAME = "rpc_queue";
    public static void main(String[] args) throws Exception {

        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost("localhost");

        Connection connection = factory.newConnection();
        Channel channel = connection.createChannel();

        channel.queueDeclare(RPC_QUEUE_NAME, false, false, false, null);

        channel.basicQos(1);

        QueueingConsumer consumer = new QueueingConsumer(channel);
        channel.basicConsume(RPC_QUEUE_NAME, false, consumer);

        System.out.println(" [x] Awaiting RPC requests");

        while (true) {
            QueueingConsumer.Delivery delivery = consumer.nextDelivery();

            BasicProperties props = delivery.getProperties();
            BasicProperties replyProps = new BasicProperties
                    .Builder()
                    .correlationId(props.getCorrelationId())
                    .build();

            String message = new String(delivery.getBody());
            int n = Integer.parseInt(message);

            System.out.println(" [.] fib(" + message + ")");
            String response = "" + fib(n);

            channel.basicPublish( "", props.getReplyTo(), replyProps, response.getBytes());

            channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
        }

    }

    private static int fib(int n) throws Exception {
        if (n == 0) return 0;
        if (n == 1) return 1;
        return fib(n-1) + fib(n-2);
    }
}

由于是共享队列,这里我们就不用exchange和routing了。

另外,有时我们可能需要运行多个服务,为了让多个服务端负载均衡,我们可以使用prefetchCount。

这个属性在之前任务队列的例子里也用过,也就是

workerChannel.basicQos(1);

即让多个worker一次获取一个任务。

用basicConsume方法进入队列后循环等待请求,发现有请求到达时根据队列和CorrelationId对相应请求作出响应。

另外需要注意的一点,server中basicConsume的第二个参数是false。

其意义为是否自动作出回应,即:

true if the server should consider messages acknowledged once delivered; false if the server should expect explicit acknowledgements

于是循环时需要显示调用basicAck进行回应。

RPCClient:

class RPCClient{

    private Connection connection;
    private Channel channel;
    private String requestQueueName = "rpc_queue";
    private String replyQueueName;
    private QueueingConsumer consumer;

    public RPCClient() throws Exception {
        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost("localhost");
        connection = factory.newConnection();
        channel = connection.createChannel();

        replyQueueName = channel.queueDeclare().getQueue();
        consumer = new QueueingConsumer(channel);
        channel.basicConsume(replyQueueName, true, consumer);
    }

    public String call(String message) throws Exception {
        String response = null;
        String corrId = java.util.UUID.randomUUID().toString();

        BasicProperties props = new BasicProperties
                .Builder()
                .correlationId(corrId)
                .replyTo(replyQueueName)
                .build();

        channel.basicPublish("", requestQueueName, props, message.getBytes());

        while (true) {
            QueueingConsumer.Delivery delivery = consumer.nextDelivery();
            if (delivery.getProperties().getCorrelationId().equals(corrId)) {
                response = new String(delivery.getBody());
                break;
            }
        }

        return response;
    }

    public void close() throws Exception {
        connection.close();
    }
}

callback队列只是一个匿名队列,但切记需要将其设置到BasicProperties中。

corrId的生成方法有很多种,在这里使用UUID。

call方法中通过调用basicPublish进行RPC请求,参数中带着BasicProperties。

RabbitMQ - RPC in Java,布布扣,bubuko.com

时间: 2024-08-09 22:25:17

RabbitMQ - RPC in Java的相关文章

RabbitMQ指南(Java)

原文地址:http://www.rabbitmq.com/getstarted.html 翻译得不好,欢迎指出. 一.Hello World 1.基本概念介绍 RabbitMQ是一个消息代理(或者说消息队列),它的主要意图很明显,就是接收和转发消息.你可以把它想象成一个邮局:当你把一封邮件放入邮箱,邮递员会帮你把邮件送到收件人的手上.在这里,RabbitMQ就好比一个邮箱.邮局或者邮递员. RabbitMQ和邮局的主要区别在于,RabbitMQ不是处理邮件,而是接收.存储和将消息以二进制的方式转

RabbitMQ实例教程:Hello RabbitMQ World之Java实现

RabbitMQ要实现Hello World,其实也很简单.只需一个服务器来发送消息,另外有个客户端接收消息即可. 整体的设计流程如下: 消息生产者发送Hello到消息队列,消息消费者从队列中接收消息. 下载依赖Jar包 RabbitMQ要用Java实现发送消息,就必须使用Java客户端库.目前RabbizMQ的Java客户端库最新版为为 3.5.5 .可以从Maven仓库下载,也可以直接去官网下载. <dependency>    <groupId>com.rabbitmq<

Openstack中RabbitMQ RPC代码分析

在Openstack中,RPC调用是通过RabbitMQ进行的. 任何一个RPC调用,都有Client/Server两部分,分别在rpcapi.py和manager.py中实现. 这里以nova-scheduler调用nova-compute为例子. nova/compute/rpcapi.py中有ComputeAPI nova/compute/manager.py中有ComputeManager 两个类有名字相同的方法,nova-scheduler调用ComputeAPI中的方法,通过底层的R

使用rabbitmq rpc 模式

服务器端 安装 ubuntu 16.04 server 安装 rabbitmq-server 设置 apt 源 curl -s https://packagecloud.io/install/repositories/rabbitmq/rabbitmq-server/script.python.sh | bash 使用 apt-get install rabbitmq-server 安装 rabbitmq 服务器 按键Y或者 y 确认安装 rabbitmq-server 简单管理 rabbitm

RabbitMQ安装以及java使用(一)

最近闲来无事,整理下基础知识,本次安装 1.RabbitMQ版本是3.6.10 2.操作系统是centOS 7 64位  虚拟机IP:192.168.149.133 1.安装更新系统环境依赖 yum install gcc glibc-devel make ncurses-devel openssl-devel xmlto 2.安装配置erlang语言环境 因为RabbitMQ是使用erlang语言开发的,所以还需要配置以下erlang语言环境 下载安装包,地址http://www.erlang

module05-1-基于RabbitMQ rpc实现的主机管理

需求 题目:rpc命令端 需求: 可以异步的执行多个命令 对多台机器 >>:run "df -h" --hosts 192.168.3.55 10.4.3.4task id: 45334>>: check_task 45334>>: 实现需求 1. 实现全部需求 2.会缓存已建立过的连接,减少短时间内连接相同主机时再次建立连接的开销 3.定时清理缓存的连接 目录结构 rabbitmq_server ├ bin # 执行文件目录 | └ rabbitm

RabbitMQ 概念与Java例子

RabbitMQ简介目前RabbitMQ是AMQP 0-9-1(高级消息队列协议)的一个实现,使用Erlang语言编写,利用了Erlang的分布式特性.概念介绍:Broker:简单来说就是消息队列服务器实体.Exchange:消息交换机,它指定消息按什么规则,路由到哪个队列.Queue:消息队列载体,每个消息都会被投入到一个或多个队列.Binding:绑定,它的作用就是把exchange和queue按照路由规则绑定起来.Routing Key:路由关键字,exchange根据这个关键字进行消息投

python--基于RabbitMQ rpc实现的主机管理

要求: 可以异步的执行多个命令对多台机器>>:run "df -h" --hosts 192.168.3.55 10.4.3.4task id: 45334>>: check_task 45334>>: 思考:1.分解其中需要实现的功能(1)命令是发到远程主机上执行的,命令放在队列里,再发到主机处理,主机执行完结果放在队列里,提交命令的人自取.就需要2个进程,一个client,提交命令,取结果,一个server,处理命令,放结果(2)发送命令的时候,

RabbitMQ (消息队列)专题学习07 RPC

(使用Java客户端) 一.概述 在Work Queue的章节中我们学习了如何使用Work Queue分配耗时的任务给多个工作者,但是如果我们需要运行一个函数在远程计算机上,这是一个完全不同的情景,这种模式通常被称之为RPC. 在本章节的学习中,我们将使用RabbitMQ来构建一个RPC系统:一个远程客户端和一个可扩展的RPC服务器,我们没有任何费时的任务进行分配,我们将创建一个虚拟的RPC服务返回Fibonacci数. 1.1.客户端接口(Client Interface) 为了说明一个RPC