RabbitMQ筆記二十二 :異步RPC之一(java client實(shí)現(xiàn)RPC功能)

異步RPC(Remote procedure call)

模型圖

Server:提供服務(wù)的服務(wù)述么,即RPC模型中的Server蝌数。
Client:調(diào)用服務(wù)的服務(wù),即RPC模型中的client度秘。

Client發(fā)送消息給服務(wù)端顶伞,消息中包含二個(gè)屬性,一個(gè)是reply_to屬性(value是一個(gè)隊(duì)列名剑梳,客戶端服務(wù)一直監(jiān)聽這個(gè)隊(duì)列)唆貌,一個(gè)是correkation_id(本次請(qǐng)求的唯一標(biāo)識(shí)),發(fā)送消息調(diào)用Server服務(wù)之后垢乙,服務(wù)端將調(diào)用結(jié)果發(fā)送到reply_to指定的隊(duì)列中锨咙。當(dāng)前也會(huì)根據(jù)客戶端得到correlation_id,并指明這次結(jié)果與correlation_id對(duì)應(yīng)追逮。

java client實(shí)現(xiàn)RPC功能

服務(wù)端代碼,監(jiān)聽sms隊(duì)列(客戶端請(qǐng)求消息發(fā)送到的隊(duì)列)

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;

import java.util.concurrent.TimeUnit;

public class Consume {
    public static void main(String[] args) throws Exception{
        ConnectionFactory connectionFactory = new ConnectionFactory();
        connectionFactory.setUri("amqp://zhihao.miao:123456@192.168.1.131:5672");

        Connection connection = connectionFactory.newConnection();

        Channel channel = connection.createChannel();

        System.out.println(channel.queueDeclare().getQueue());

        channel.basicConsume("sms",true,new SimpleConsumer(channel));
        System.out.println("短信服務(wù)已經(jīng)啟動(dòng)");
        TimeUnit.SECONDS.sleep(60);

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

接收到客戶端的消息之后酪刀,調(diào)用服務(wù)接口sendSMS方法粹舵,然后得到請(qǐng)求消息中的reply_tocorrelation_id屬性,將調(diào)用接口的結(jié)果發(fā)送到reply_to屬性的隊(duì)列中骂倘,指定correlation_id屬性眼滤。

import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.DefaultConsumer;
import com.rabbitmq.client.Envelope;
import com.zhihao.miao.rpc.server.SendSMSTool;

import java.io.IOException;

public class SimpleConsumer extends DefaultConsumer {

    public SimpleConsumer(Channel channel){
        super(channel);
    }

    @Override
    public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
        System.out.println("進(jìn)入RPC方法調(diào)用");
        String phone =properties.getHeaders().get("phone").toString();
        String content = new String(body);
        //調(diào)用服務(wù)
        boolean result = SendSMSTool.sendSMS(phone,content);
        System.out.println("消息處理成功");

        String reply = properties.getReplyTo();
        String id = properties.getCorrelationId();

        AMQP.BasicProperties props = new AMQP.BasicProperties.Builder().correlationId(id).build();
        this.getChannel().basicPublish("",reply,props,(result+"").getBytes());
        System.out.println("消息回復(fù)成功");
    }
}

服務(wù)接口:

public class SendSMSTool {

    public static boolean sendSMS(String phone,String content){
        System.out.println("發(fā)送短信內(nèi)容:【"+content+"】到手機(jī)號(hào):"+phone);
        return phone.length() > 6;
    }
}

總結(jié)
RPC Server步驟
1.創(chuàng)建服務(wù)
2.監(jiān)聽一個(gè)隊(duì)列(sms),監(jiān)聽客戶端發(fā)送的消息
3.收到消息之后历涝,調(diào)用服務(wù)诅需,得到調(diào)用結(jié)果
4.從消息屬性中,獲取reply_to, correlation_id屬性荧库,把調(diào)用結(jié)果發(fā)送給reply_to指定的隊(duì)列中堰塌,發(fā)送的消息屬性要帶上reply_to。
5.一次調(diào)用處理成功

Exchange綁定sms隊(duì)列

客戶端代碼
消息端先監(jiān)聽一個(gè)隊(duì)列sms.reply电爹,這個(gè)隊(duì)列是客戶端返回結(jié)果的隊(duì)列蔫仙,然后發(fā)送消息指定請(qǐng)求參數(shù)(參數(shù)頭和參數(shù)內(nèi)容)料睛,指定correlationIdreplyTo屬性丐箩。

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

import java.util.HashMap;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.TimeUnit;

public class Send {
    public static void main(String[] args) throws Exception{
        ConnectionFactory connectionFactory = new ConnectionFactory();

        connectionFactory.setUri("amqp://zhihao.miao:123456@192.168.1.131:5672");

        Connection connection = connectionFactory.newConnection();

        Channel channel = connection.createChannel();

        String correlationId = UUID.randomUUID().toString();
        String replyTo = "sms.reply";
        //自動(dòng)刪除屬性為true,當(dāng)connection關(guān)閉后會(huì)自動(dòng)刪除該隊(duì)列
        //channel.queueDeclare(replyTo,true,true,true,new HashMap<>());

        channel.basicConsume(replyTo,true,new SimpleConsumer(channel));

        Map<String,Object> headers = new HashMap<>();
        headers.put("phone","15790934342");

        AMQP.BasicProperties properties = new AMQP.BasicProperties.Builder().headers(headers).replyTo(replyTo).correlationId(correlationId).deliveryMode(2).
                contentEncoding("UTF-8").build();

        channel.basicPublish("send","sms",true,properties,"周年慶6折大促銷恤煞,只剩三天屎勘。詳情登錄****,去了解詳情居扒。".getBytes());

        TimeUnit.SECONDS.sleep(20);

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

接收到rpc調(diào)用返回接口的處理邏輯概漱,

import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.DefaultConsumer;
import com.rabbitmq.client.Envelope;

import java.io.IOException;


public class SimpleConsumer extends DefaultConsumer {

    public SimpleConsumer(Channel channel){
        super(channel);
    }

    @Override
    public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
        System.out.println("====收到RPC調(diào)用 回復(fù)了=====");
        System.out.println(properties.getCorrelationId()+",短信發(fā)送結(jié)果:"+new String(body));


    }
}

總結(jié)
RPC Client步驟:
1.監(jiān)聽reply_to對(duì)應(yīng)的隊(duì)列(RPC調(diào)用結(jié)果發(fā)送指定的隊(duì)列)
2.發(fā)送消息,消息屬性需要帶上reply_to喜喂,correlation_id屬性
3.服務(wù)端處理完成之后瓤摧,reply_to對(duì)應(yīng)的隊(duì)列就會(huì)收到異步處理結(jié)果消息
4.收到消息之后,進(jìn)行處理玉吁,根據(jù)消息屬性的correlation_id找到對(duì)應(yīng)的請(qǐng)求
5.一次客戶端調(diào)用就完成了照弥。

說一下邏輯:
就是客戶端往一個(gè)隊(duì)列sms發(fā)送消息,服務(wù)端監(jiān)聽這個(gè)隊(duì)列进副,接收到消息之后服務(wù)端在消息處理方法中調(diào)用服務(wù)这揣,并把調(diào)用的服務(wù)結(jié)果返回給客戶端,怎么返回呢(通過往一個(gè)隊(duì)列里發(fā)送(sms.reply),并把這次調(diào)用的id也返回給客戶端影斑,這個(gè)隊(duì)列名是reply_to屬性给赞,當(dāng)前的返回id是correlation_id,都是在客戶端往sms隊(duì)列發(fā)送消息的時(shí)候帶到服務(wù)端的),服務(wù)端將返回的結(jié)果發(fā)送到指定的隊(duì)列矫户,而客戶端一直在監(jiān)聽這個(gè)隊(duì)列拿到返回值片迅,即完成了一次異步的rpc調(diào)用。

Direct reply-to

Improve performance and simplicity of RPC clients by sending replies direct to a waiting channel.
通過直接將RPC的回復(fù)發(fā)送到等待的Channel中而提供RPC客戶端的性能皆辽。

動(dòng)機(jī)(原因)
RPC is a popular pattern to implement with a messaging broker like RabbitMQ. The typical way to do this is for RPC clients to send requests to a long lived server queue. The RPC server(s) read requests from this queue and then send replies to each client using the queue named by the client in the reply-to header.

But where does the client's queue come from? The client can declare a single-use queue for each request-response pair. But this is inefficient; even a transient unmirrored queue can be expensive to create and then delete (compared with the cost of sending a message). This is especially true in a cluster as all cluster nodes need to agree that the queue has been created, even if it is unmirrored.

So alternatively the client can create a long-lived queue for its replies. But this can be fiddly to manage, especially if the client itself is not long-lived.

The direct reply-to feature allows RPC clients to receive replies directly from their RPC server, without going through a reply queue. ("Directly" here still means going through AMQP and the RabbitMQ server; there is no separate network connection between RPC client and RPC server.)

大意就是使用RabbitMQ實(shí)現(xiàn)異步RPC柑蛇,通常的做法就是客戶端將請(qǐng)求發(fā)送到長期存在的服務(wù)器隊(duì)列罐旗。RPC服務(wù)器讀取接收此隊(duì)列的請(qǐng)求,然后將RPC調(diào)用結(jié)果返回到客戶端在發(fā)送請(qǐng)求時(shí)請(qǐng)求頭中的reply_to屬性中的隊(duì)列唯蝶。完成一次異步的RPC調(diào)用九秀。
那么請(qǐng)求頭中的隊(duì)列哪里來呢?客戶端可以聲明一個(gè)針對(duì)每一對(duì)請(qǐng)求的臨時(shí)隊(duì)列粘我。但是開銷太大(相比于發(fā)送消息本身)鼓蜒,在集群環(huán)境中更是如此。
如果創(chuàng)建一個(gè)長期的隊(duì)列呢征字?如果客戶端本身并不是長期的那么這個(gè)長期的隊(duì)列就很難管理都弹。

Direct reply-to就是解決這個(gè)問題的,不需要?jiǎng)?chuàng)建這樣一個(gè)回復(fù)的隊(duì)列匙姜。

示列

服務(wù)端
服務(wù)啟動(dòng)類畅厢,監(jiān)聽sms隊(duì)列

import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;

import java.util.concurrent.TimeUnit;

public class Consume {
    public static void main(String[] args) throws Exception{
        ConnectionFactory connectionFactory = new ConnectionFactory();
        connectionFactory.setUri("amqp://zhihao.miao:123456@192.168.1.131:5672");

        Connection connection = connectionFactory.newConnection();

        Channel channel = connection.createChannel();

        System.out.println(channel.queueDeclare().getQueue());

        channel.basicConsume("sms",true,new SimpleConsumer(channel));
        System.out.println("短信服務(wù)已經(jīng)啟動(dòng)");
        TimeUnit.SECONDS.sleep(60);

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

服務(wù)端監(jiān)聽消息處理器,接收到請(qǐng)求參數(shù)氮昧,調(diào)用服務(wù)接口框杜,返回參數(shù)

import com.rabbitmq.client.AMQP.BasicProperties;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.DefaultConsumer;
import com.rabbitmq.client.Envelope;

import java.io.IOException;

public class SimpleConsumer extends DefaultConsumer {

    public SimpleConsumer(Channel channel){
        super(channel);
    }

    @Override
    public void handleDelivery(String consumerTag, Envelope envelope, BasicProperties properties, byte[] body) throws IOException {
        System.out.println("進(jìn)入RPC方法調(diào)用");
        String phone =properties.getHeaders().get("phone").toString();
        String content = new String(body);
        //調(diào)用服務(wù)
        boolean result = SendSMSTool.sendSMS(phone,content);
        System.out.println("消息處理成功");

        String reply = properties.getReplyTo();
        String id = properties.getCorrelationId();

        BasicProperties props = new BasicProperties.Builder().correlationId(id).build();
        this.getChannel().basicPublish("",reply,props,(result+"").getBytes());
        System.out.println("消息回復(fù)成功");
    }
}

服務(wù)接口

public class SendSMSTool {

    public static boolean sendSMS(String phone,String content){
        System.out.println("發(fā)送短信內(nèi)容:【"+content+"】到手機(jī)號(hào):"+phone);
        return phone.length() > 6;
    }
}

客戶端
改變的是客戶端代碼,直接定義replyToamq.rabbitmq.reply-to

import com.rabbitmq.client.*;

import java.util.HashMap;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.TimeUnit;

public class Send {
    public static void main(String[] args) throws Exception{
        ConnectionFactory connectionFactory = new ConnectionFactory();

        connectionFactory.setUri("amqp://zhihao.miao:123456@192.168.1.131:5672");

        Connection connection = connectionFactory.newConnection();

        Channel channel = connection.createChannel();

        String correlationId = UUID.randomUUID().toString();

        String replyTo ="amq.rabbitmq.reply-to";

        channel.basicConsume(replyTo,true,new SimpleConsumer(channel));

        Map<String,Object> headers = new HashMap<>();
        headers.put("phone","15790934342");

        AMQP.BasicProperties properties = new AMQP.BasicProperties.Builder().headers(headers).replyTo(replyTo).correlationId(correlationId).deliveryMode(2).
                contentEncoding("UTF-8").build();

        channel.basicPublish("send","sms",true,properties,"周年慶6折大促銷袖肥,只剩三天筒主。詳情登錄****懊亡,去了解詳情知牌。".getBytes());

        TimeUnit.SECONDS.sleep(20);

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

定義消息監(jiān)聽消費(fèi)者

import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.DefaultConsumer;
import com.rabbitmq.client.Envelope;

import java.io.IOException;

public class SimpleConsumer extends DefaultConsumer {

    public SimpleConsumer(Channel channel){
        super(channel);
    }

    @Override
    public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
        System.out.println("====收到RPC調(diào)用恢復(fù)了=====");
        System.out.println(properties.getCorrelationId()+",短信發(fā)送結(jié)果:"+new String(body));


    }
}

管控臺(tái)中也不會(huì)創(chuàng)建amq.rabbitmq.reply-to這個(gè)隊(duì)列赶盔,而是直接將返回的消息發(fā)到等到的Channel中,提升異步rpc調(diào)用的性能寸癌,amq.rabbitmq.reply-to也叫做偽隊(duì)列专筷。

參考資料

Remote procedure call (RPC)
Direct reply-to

最后編輯于
?著作權(quán)歸作者所有,轉(zhuǎn)載或內(nèi)容合作請(qǐng)聯(lián)系作者
  • 序言:七十年代末,一起剝皮案震驚了整個(gè)濱河市蒸苇,隨后出現(xiàn)的幾起案子磷蛹,更是在濱河造成了極大的恐慌,老刑警劉巖填渠,帶你破解...
    沈念sama閱讀 216,919評(píng)論 6 502
  • 序言:濱河連續(xù)發(fā)生了三起死亡事件弦聂,死亡現(xiàn)場離奇詭異,居然都是意外死亡氛什,警方通過查閱死者的電腦和手機(jī)莺葫,發(fā)現(xiàn)死者居然都...
    沈念sama閱讀 92,567評(píng)論 3 392
  • 文/潘曉璐 我一進(jìn)店門,熙熙樓的掌柜王于貴愁眉苦臉地迎上來枪眉,“玉大人捺檬,你說我怎么就攤上這事∶惩” “怎么了堡纬?”我有些...
    開封第一講書人閱讀 163,316評(píng)論 0 353
  • 文/不壞的土叔 我叫張陵聂受,是天一觀的道長。 經(jīng)常有香客問我烤镐,道長蛋济,這世上最難降的妖魔是什么? 我笑而不...
    開封第一講書人閱讀 58,294評(píng)論 1 292
  • 正文 為了忘掉前任炮叶,我火速辦了婚禮碗旅,結(jié)果婚禮上,老公的妹妹穿的比我還像新娘镜悉。我一直安慰自己祟辟,他們只是感情好,可當(dāng)我...
    茶點(diǎn)故事閱讀 67,318評(píng)論 6 390
  • 文/花漫 我一把揭開白布侣肄。 她就那樣靜靜地躺著旧困,像睡著了一般。 火紅的嫁衣襯著肌膚如雪稼锅。 梳的紋絲不亂的頭發(fā)上吼具,一...
    開封第一講書人閱讀 51,245評(píng)論 1 299
  • 那天,我揣著相機(jī)與錄音缰贝,去河邊找鬼馍悟。 笑死畔濒,一個(gè)胖子當(dāng)著我的面吹牛剩晴,可吹牛的內(nèi)容都是我干的。 我是一名探鬼主播侵状,決...
    沈念sama閱讀 40,120評(píng)論 3 418
  • 文/蒼蘭香墨 我猛地睜開眼赞弥,長吁一口氣:“原來是場噩夢(mèng)啊……” “哼!你這毒婦竟也來了趣兄?” 一聲冷哼從身側(cè)響起绽左,我...
    開封第一講書人閱讀 38,964評(píng)論 0 275
  • 序言:老撾萬榮一對(duì)情侶失蹤,失蹤者是張志新(化名)和其女友劉穎艇潭,沒想到半個(gè)月后拼窥,有當(dāng)?shù)厝嗽跇淞掷锇l(fā)現(xiàn)了一具尸體,經(jīng)...
    沈念sama閱讀 45,376評(píng)論 1 313
  • 正文 獨(dú)居荒郊野嶺守林人離奇死亡蹋凝,尸身上長有42處帶血的膿包…… 初始之章·張勛 以下內(nèi)容為張勛視角 年9月15日...
    茶點(diǎn)故事閱讀 37,592評(píng)論 2 333
  • 正文 我和宋清朗相戀三年鲁纠,在試婚紗的時(shí)候發(fā)現(xiàn)自己被綠了。 大學(xué)時(shí)的朋友給我發(fā)了我未婚夫和他白月光在一起吃飯的照片鳍寂。...
    茶點(diǎn)故事閱讀 39,764評(píng)論 1 348
  • 序言:一個(gè)原本活蹦亂跳的男人離奇死亡改含,死狀恐怖,靈堂內(nèi)的尸體忽然破棺而出迄汛,到底是詐尸還是另有隱情捍壤,我是刑警寧澤骤视,帶...
    沈念sama閱讀 35,460評(píng)論 5 344
  • 正文 年R本政府宣布,位于F島的核電站鹃觉,受9級(jí)特大地震影響专酗,放射性物質(zhì)發(fā)生泄漏。R本人自食惡果不足惜盗扇,卻給世界環(huán)境...
    茶點(diǎn)故事閱讀 41,070評(píng)論 3 327
  • 文/蒙蒙 一笼裳、第九天 我趴在偏房一處隱蔽的房頂上張望。 院中可真熱鬧粱玲,春花似錦躬柬、人聲如沸。這莊子的主人今日做“春日...
    開封第一講書人閱讀 31,697評(píng)論 0 22
  • 文/蒼蘭香墨 我抬頭看了看天上的太陽。三九已至卵沉,卻和暖如春颠锉,著一層夾襖步出監(jiān)牢的瞬間,已是汗流浹背史汗。 一陣腳步聲響...
    開封第一講書人閱讀 32,846評(píng)論 1 269
  • 我被黑心中介騙來泰國打工琼掠, 沒想到剛下飛機(jī)就差點(diǎn)兒被人妖公主榨干…… 1. 我叫王不留,地道東北人停撞。 一個(gè)月前我還...
    沈念sama閱讀 47,819評(píng)論 2 370
  • 正文 我出身青樓瓷蛙,卻偏偏與公主長得像,于是被迫代替她去往敵國和親戈毒。 傳聞我的和親對(duì)象是個(gè)殘疾皇子艰猬,可洞房花燭夜當(dāng)晚...
    茶點(diǎn)故事閱讀 44,665評(píng)論 2 354

推薦閱讀更多精彩內(nèi)容