设为首页 收藏本站
查看: 195|回复: 0

[经验分享] RabbitMQ实例教程:RPC远程调用消息队列

[复制链接]
累计签到:1 天
连续签到:1 天
发表于 2015-10-23 09:51:12 | 显示全部楼层 |阅读模式
 在工作队列一章中,我们学会了如何使用工作队列来处理多个工作进程间分发任务,但如果我们想要运行远程计算机上的函数来获得结果呢?这就是本章要处理的问题RPC。
  本节我们会使用RabbitMQ构建一个RPC系统:一个客户端和一个可扩展的RPC服务器。因为我们没有任何耗时的任务值得分发下去,我们构建一个虚拟的服务来返回斐波纳契数列。

  客户端接口

  我们创建一个客户端类来说明如何使用RPC服务,暴露一个call方法来发送RPC请求和数据获取结果。

1
2
3
FibonacciRpcClient fibonacciRpc = new FibonacciRpcClient();  
String result = fibonacciRpc.call("4");
System.out.println( "fib(4) is " + result);




  尽管RPC是编程中一种常见的模式,但其也常常饱受批评。因为程序员常常不知道调用的方法是本地方法还是一个RPC方法,这在调试中常常增加一些不必要的复杂性。我们应该简化代码,而不是滥用RPC导致代码变的臃肿。

  回调队列

  一般来说,通过RabbitMQ实现RPC非常简单,客户端发送一个请求消息,服务端响应消息就完成了。为了接收到响应内容,我们在请求中发送”callback“队列地址,也可以使用默认的队列。

1
2
3
callbackQueueName = channel.queueDeclare().getQueue();
BasicProperties props = new BasicProperties .Builder().replyTo(callbackQueueName) .build();
channel.basicPublish("", "rpc_queue", props, message.getBytes());




  AMQP协议中预定了14个消息属性,除了下面几个,其它的都很少使用:

  deliveryMode : 标识消息是持久化还是瞬态的。

  contentType : 描述 mime-type的编码类型,如JSON编码为”application/json“。

  replyTo : 通常在回调队列中使用。

  correlationId : 在请求中关联RPC响应时使用。

  关联Id(Correlation Id)

  在前面的方法中,要求在每个RPC请求创建回调队列,这可真是一件繁琐的事情,但幸运的是我们有个好方法-在每个客户端创建一个简单的回调队列。

  这样问题又来了,队列如何知道这些响应来自哪个请求呢?这时候correlationId就出场了。我们在每个请求中都设置一个唯一的值,这样我们在回调队列中接收消息的时候就能知道是哪个请求发送的。如果收到未知的correlationId,就废弃该消息,因为它不是我们发出的请求。

  你可能会问,为什么抛弃未知消息而不是抛出错误呢?这是由服务器竞争资源所导致的。尽管这不太可能,试想一下,如果RPC服务器在发送完响应后而在发送应答消息前死掉了,重启RPC服务器会重新发送请求。这就是我们在客户机上优雅地处理重复的反应,RPC应该是等同的。

wKioL1YoynDhkgSMAACbyFOQwfU789.jpg

  (1)客户端启动,创建一个匿名且唯一的回调队列。

  (2)对每个RPC请求,客户端发送一个包含replyTo和correlationId两个属性的消息。

  (3)请求发送到rpc_queue队列。

  (4)RPC服务在队列中等待请求,当请求出现时,根据replyTo字段使用队列将结果发送到客户端。

  (5)客户端在回调队列中等待数据。当消息出现时,它会检查correlationId属性,如果该值匹配的话,就会返回响应结果给应用。

  示例代码

  RPCServer.java
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
package com.favccxx.favrabbit;

import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.QueueingConsumer;
import com.rabbitmq.client.AMQP.BasicProperties;

public class RPCServer {

    private static final String RPC_QUEUE_NAME = "rpc_queue";

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

    public static void main(String[] argv) {
        Connection connection = null;
        Channel channel = null;
        try {
            ConnectionFactory factory = new ConnectionFactory();
            factory.setHost("localhost");

            connection = factory.newConnection();
            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) {
                String response = null;

                QueueingConsumer.Delivery delivery = consumer.nextDelivery();

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

                try {
                    String message = new String(delivery.getBody(), "UTF-8");
                    int n = Integer.parseInt(message);

                    System.out.println(" [.] fib(" + message + ")");
                    response = "" + fib(n);
                } catch (Exception e) {
                    System.out.println(" [.] " + e.toString());
                    response = "";
                } finally {
                    channel.basicPublish("", props.getReplyTo(), replyProps, response.getBytes("UTF-8"));

                    channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
                }
            }
        } catch (Exception e) {
            e.printStackTrace();
        } finally {
            if (connection != null) {
                try {
                    connection.close();
                } catch (Exception ignore) {
                }
            }
        }
    }
}



  RPCClient.java      

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
package com.favccxx.favrabbit;

import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.QueueingConsumer;
import com.rabbitmq.client.AMQP.BasicProperties;
import java.util.UUID;

public 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 = UUID.randomUUID().toString();

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

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

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

        return response;
    }

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

    public static void main(String[] argv) {
        RPCClient fibonacciRpc = null;
        String response = null;
        try {
            fibonacciRpc = new RPCClient();

            System.out.println(" [x] Requesting fib(30)");
            response = fibonacciRpc.call("30");
            System.out.println(" [.] Got '" + response + "'");
        } catch (Exception e) {
            e.printStackTrace();
        } finally {
            if (fibonacciRpc != null) {
                try {
                    fibonacciRpc.close();
                } catch (Exception ignore) {
                }
            }
        }
    }
}




  先启动RPCServer,然后运行RPCClient,控制台输出如下内容
RPCClient[x] Requesting fib(30)
RPCClient[.] Got '832040'

RPCServer[x] Awaiting RPC requests
RPCServer[.] fib(30)

   


运维网声明 1、欢迎大家加入本站运维交流群:群②:261659950 群⑤:202807635 群⑦870801961 群⑧679858003
2、本站所有主题由该帖子作者发表,该帖子作者与运维网享有帖子相关版权
3、所有作品的著作权均归原作者享有,请您和我们一样尊重他人的著作权等合法权益。如果您对作品感到满意,请购买正版
4、禁止制作、复制、发布和传播具有反动、淫秽、色情、暴力、凶杀等内容的信息,一经发现立即删除。若您因此触犯法律,一切后果自负,我们对此不承担任何责任
5、所有资源均系网友上传或者通过网络收集,我们仅提供一个展示、介绍、观摩学习的平台,我们不对其内容的准确性、可靠性、正当性、安全性、合法性等负责,亦不承担任何法律责任
6、所有作品仅供您个人学习、研究或欣赏,不得用于商业或者其他用途,否则,一切后果均由您自己承担,我们对此不承担任何法律责任
7、如涉及侵犯版权等问题,请您及时通知我们,我们将立即采取措施予以解决
8、联系人Email:admin@iyunv.com 网址:www.yunweiku.com

所有资源均系网友上传或者通过网络收集,我们仅提供一个展示、介绍、观摩学习的平台,我们不对其承担任何法律责任,如涉及侵犯版权等问题,请您及时通知我们,我们将立即处理,联系人Email:kefu@iyunv.com,QQ:1061981298 本贴地址:https://www.yunweiku.com/thread-129845-1-1.html 上篇帖子: fedora下命令yum 下篇帖子: RedHat6.0 ES 本地yum服务器搭建配置
您需要登录后才可以回帖 登录 | 立即注册

本版积分规则

扫码加入运维网微信交流群X

扫码加入运维网微信交流群

扫描二维码加入运维网微信交流群,最新一手资源尽在官方微信交流群!快快加入我们吧...

扫描微信二维码查看详情

客服E-mail:kefu@iyunv.com 客服QQ:1061981298


QQ群⑦:运维网交流群⑦ QQ群⑧:运维网交流群⑧ k8s群:运维网kubernetes交流群


提醒:禁止发布任何违反国家法律、法规的言论与图片等内容;本站内容均来自个人观点与网络等信息,非本站认同之观点.


本站大部分资源是网友从网上搜集分享而来,其版权均归原作者及其网站所有,我们尊重他人的合法权益,如有内容侵犯您的合法权益,请及时与我们联系进行核实删除!



合作伙伴: 青云cloud

快速回复 返回顶部 返回列表