RabbitMQ实例教程:RPC远程调用消息队列

简介:

 在工作队列一章中,我们学会了如何使用工作队列来处理多个工作进程间分发任务,但如果我们想要运行远程计算机上的函数来获得结果呢?这就是本章要处理的问题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)






本文转自 genuinecx 51CTO博客,原文链接:http://blog.51cto.com/favccxx/1705357,如需转载请自行联系原作者
相关实践学习
RocketMQ一站式入门使用
从源码编译、部署broker、部署namesrv,使用java客户端首发消息等一站式入门RocketMQ。
消息队列 MNS 入门课程
1、消息队列MNS简介 本节课介绍消息队列的MNS的基础概念 2、消息队列MNS特性 本节课介绍消息队列的MNS的主要特性 3、MNS的最佳实践及场景应用 本节课介绍消息队列的MNS的最佳实践及场景应用案例 4、手把手系列:消息队列MNS实操讲 本节课介绍消息队列的MNS的实际操作演示 5、动手实验:基于MNS,0基础轻松构建 Web Client 本节课带您一起基于MNS,0基础轻松构建 Web Client
目录
相关文章
|
2月前
|
物联网
MQTT常见问题之用单片机接入阿里MQTT实例失败如何解决
MQTT(Message Queuing Telemetry Transport)是一个轻量级的、基于发布/订阅模式的消息协议,广泛用于物联网(IoT)中设备间的通信。以下是MQTT使用过程中可能遇到的一些常见问题及其答案的汇总:
|
2月前
|
消息中间件 Java RocketMQ
RocketMQ实战教程之RocketMQ安装
这是一篇关于RocketMQ安装的实战教程,主要介绍了在CentOS系统上使用传统安装和Docker两种方式安装RocketMQ。首先,系统需要是64位,并且已经安装了JDK 1.8。传统安装包括下载安装包,解压并启动NameServer和Broker。Docker安装则涉及安装docker和docker-compose,然后通过docker-compose.yaml文件配置并启动服务。教程还提供了启动命令和解决问题的提示。
|
2月前
|
消息中间件 前端开发 数据库
RocketMQ实战教程之MQ简介与应用场景
RocketMQ实战教程介绍了MQ的基本概念和应用场景。MQ(消息队列)是生产者和消费者模型,用于异步传输数据,实现系统解耦。消息中间件在生产者发送消息和消费者接收消息之间起到邮箱作用,简化通信。主要应用场景包括:1)应用解耦,如订单系统与库存系统的非直接交互;2)异步处理,如用户注册后的邮件和短信发送延迟处理,提高响应速度;3)流量削峰,如秒杀活动限制并发流量,防止系统崩溃。
|
2月前
|
消息中间件 存储 Apache
RocketMQ实战教程之常见概念和模型
Apache RocketMQ 实战教程介绍了其核心概念和模型。消息是基本的数据传输单元,主题是消息的分类容器,支持字节、数字和短划线命名,最长64个字符。消息类型包括普通、顺序、事务和定时/延时消息。消息队列是实际存储和传输消息的容器,是主题的分区。消费者分组是一组行为一致的消费者的逻辑集合,也有命名限制。此外,文档还提到了一些使用约束和建议,如主题和消费者组名的命名规则,消息大小限制,请求超时时间等。RocketMQ 提供了多种消息模型,包括发布/订阅模型,有助于理解和优化消息处理。
|
2月前
|
消息中间件 存储 Java
RocketMQ实战教程之NameServer与BrokerServer
这是一个关于RocketMQ实战教程的概要,主要讨论NameServer和BrokerServer的角色。NameServer负责管理所有BrokerServer,而BrokerServer存储和传输消息。生产者和消费者通过NameServer找到合适的Broker进行交互,不需要直接知道Broker的具体信息。工作流程包括生产者向NameServer查询后发送消息到Broker,以及消费者同样通过NameServer获取消息进行消费。这种设计类似于服务注册中心的概念,便于系统扩展和集群管理。
|
28天前
|
消息中间件 Java Spring
最新spingboot整合rabbitmq详细教程
最新spingboot整合rabbitmq详细教程
|
2月前
|
消息中间件 Cloud Native 自动驾驶
RocketMQ实战教程之MQ简介
Apache RocketMQ 是一个云原生的消息流平台,支持消息、事件和流处理,适用于云边端一体化场景。官网提供详细文档和下载资源:[RocketMQ官网](https://rocketmq.apache.org/zh/)。示例中提到了RocketMQ在物联网(如小米台灯)和自动驾驶等领域的应用。要开始使用,可从[下载页面](https://rocketmq.apache.org/zh/download)获取软件。
|
2月前
|
消息中间件 中间件 Java
RocketMQ实战教程之几种MQ优缺点以及选型
该文介绍了几种主流消息中间件,包括ActiveMQ、RabbitMQ、RocketMQ和Kafka。ActiveMQ和RabbitMQ是较老牌的选择,前者在中小企业中常见,后者因强大的并发能力和活跃社区而流行。RocketMQ是阿里巴巴的开源产品,适用于大规模分布式系统,尤其在数据可靠性方面进行了优化。Kafka最初设计用于大数据日志处理,强调高吞吐量。在选择MQ时,考虑因素包括性能、功能、开发语言、社区支持、学习难度、稳定性和集群功能。小型公司推荐使用RabbitMQ,而大型公司则可在RocketMQ和Kafka之间根据具体需求抉择。
|
1月前
|
消息中间件 存储 Java
RocketMQ下载安装、集群搭建保姆级教程
RocketMQ下载安装、集群搭建保姆级教程
43 0
|
2月前
|
消息中间件 人工智能 Java
Spring Boot+RocketMQ 实现多实例分布式环境下的事件驱动
Spring Boot+RocketMQ 实现多实例分布式环境下的事件驱动
56 1

热门文章

最新文章

相关产品

  • 云消息队列 MQ