4.6.1 简介

前面几种模式的通信都是基于Producer发送消息到Consumer,然后Consumer进行消费,假设我们需要Consumer操作完毕之后返回给Producer一个回调呢?前面几种模式就行不通了;

例如我们要做一个远程调用加钱操作,客户端远程调用服务端进行加钱操作,操作完毕之后服务端将用户最新的余额返回给客户端;客户端进行后续操作,例如更新到数据库等;

  • RPC业务分析

在这里插入图片描述

在RPC模式中,客户端和服务器都是Producer也都是Consumer;

RPC模式官网介绍:https://www.rabbitmq.com/tutorials/tutorial-five-java.html

  • RPC调用图解:

在这里插入图片描述

4.6.2 客户端
package com.dfbz.rabbitmq;

import com.rabbitmq.client.\*;

import java.io.IOException;
import java.util.UUID;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.TimeoutException;

public class RPCClient implements AutoCloseable {

    public Connection connection;
    public Channel channel;
    public static final String RPC_QUEUE_NAME = "rpc\_queue";

    public RPCClient() throws IOException, TimeoutException {
        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost("192.168.40.132");
        factory.setPort(5672);
        factory.setUsername("lscl");
        factory.setPassword("admin");
        factory.setVirtualHost("/lscl");

        connection = factory.newConnection();
        channel = connection.createChannel();
    }

    public static void main(String[] argv) throws Exception {
        // 初始化信息
        RPCClient rpcClient = new RPCClient();
        // 发起远程调用
        Integer response = rpcClient.call(20);

        System.out.println(response);

        rpcClient.channel.close();
        rpcClient.connection.close();
    }

    public Integer call(Integer money) throws IOException, InterruptedException {

        // 随机生成一个correlationId(密钥)
        final String corrId = UUID.randomUUID().toString();

        // 后期服务端回调给客户端的队列名(随机生成的回调队列名)
        String replyQueueName = channel.queueDeclare().getQueue();
        
        // 设置发送消息的一些参数
        AMQP.BasicProperties props = new AMQP.BasicProperties
 .Builder()
                .correlationId(corrId)			// 密钥
                .replyTo(replyQueueName)		// 回调队列名
                .build();

		// 采用Simple模式发送给Server端
        channel.basicPublish("", RPC_QUEUE_NAME, props, (money + "").getBytes("UTF-8"));

        // 定义延迟队列
        final BlockingQueue<Integer> response = new ArrayBlockingQueue<>(1);

        channel.basicConsume(replyQueueName, true, new DefaultConsumer(channel) {

            // 回调方法,当收到消息之后,会自动执行该方法
            public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {

                if (properties.getCorrelationId().equals(corrId)) {
                    System.out.println("响应的消息:" + new String(body));

                    // 往延迟队列中添加信息(服务端响应的最新余额)
                    response.offer(Integer.parseInt(new String(body, "UTF-8")));
                }
            }
        });

        // 获取延迟队列中的信息(如果没有信息将一直阻塞)
        return response.take();
    }

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

4.6.2 服务端
package com.dfbz.rabbitmq;

import com.rabbitmq.client.\*;

import java.io.IOException;

public class RPCServer {

    private static final String RPC_QUEUE_NAME = "rpc\_queue";

    // 总金额
    private static Integer money = 0;

    /\*\*
 \* 加钱方法
 \* @param n
 \* @return
 \*/
    private static Integer addMoney(int n) {
        money += n;
        return money;
    }

    public static void main(String[] argv) throws Exception {
        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost("192.168.40.132");
        factory.setPort(5672);
        factory.setUsername("lscl");


## 学习路线:

这个方向初期比较容易入门一些,掌握一些基本技术,拿起各种现成的工具就可以开黑了。不过,要想从脚本小子变成黑客大神,这个方向越往后,需要学习和掌握的东西就会越来越多以下是网络渗透需要学习的内容:  
 ![在这里插入图片描述](https://img-blog.csdnimg.cn/7a04c5d629f1415a9e35662316578e07.png#pic_center)



Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐