Fescar源碼閱讀-RPC和消息

TM、RM和TC之間如何通信壤巷。(源碼持續(xù)更新邑彪,本文僅供參考)



Fescar處理分布式事務(wù),本身也是分布式系統(tǒng)胧华。其中TM和RM可理解為客戶端Client寄症,TC事務(wù)協(xié)調(diào)器作為服務(wù)端Server。
Fescar中兩者之間的通信矩动,毫不意外是通過(guò)Netty實(shí)現(xiàn)的有巧。從以理解Fescar作為目標(biāo)這一角度來(lái)說(shuō),單純的網(wǎng)絡(luò)通信的模塊悲没,其實(shí)沒(méi)有必要去刨根究底篮迎,關(guān)注著一細(xì)節(jié)。
但是目前Fescar的實(shí)現(xiàn)中,請(qǐng)求甜橱、響應(yīng)逊笆、編碼解碼、消息處理等等聯(lián)系的比較緊密(耦合了岂傲。难裆。),如果想更好的跟蹤譬胎、理解Fescar是如何處理不同類(lèi)型的事務(wù)消息的,還是有必要把Fescar的rpc實(shí)現(xiàn)梳理一下命锄。


消息

requestMessage

Fescar為每一種交互都定義了一個(gè)消息類(lèi)(request和response為對(duì)應(yīng)一組)堰乔,為了容易理解,此圖只包含Request脐恩,镐侯,從上圖可以看出,
消息主要分三大類(lèi)

  • AbstractTransactionRequestToTC
    RM/TM發(fā)送到TC的消息
  • AbstractTransactionRequestToRM
    TC發(fā)送到TM的消息
  • AbstractIdentifyRequest
    TM/RM和TC連接時(shí)的注冊(cè)消息
    另外一個(gè)MergedWarpMessage是一個(gè)消息包裝類(lèi)驶冒,內(nèi)部包含多個(gè)消息苟翻,用于消息的異步批量發(fā)送。

RPC核心類(lèi)

核心類(lèi)圖如下:


TC骗污、TM崇猫、RM

主要實(shí)現(xiàn)類(lèi):

  • TC實(shí)現(xiàn) RpcServer extends AbstractRpcRemotingServer
  • TM實(shí)現(xiàn) TmRpcClient extends AbstractRpcRemotingClient
  • RM實(shí)現(xiàn) RMRpcClient extends AbstractRpcRemotingClient
    他們擁有共同父類(lèi)AbstractRpcRemoting
    AbstractRpcRemoting提供了基本的消息發(fā)送和消息處理實(shí)現(xiàn)需忿。
@Override
public void channelRead(final ChannelHandlerContext ctx, Object msg) throws Exception {
    if (msg instanceof RpcMessage) {
        final RpcMessage rpcMessage = (RpcMessage)msg;
        if (rpcMessage.isRequest()) {
            try {
                AbstractRpcRemoting.this.messageExecutor.execute(new Runnable() {
                    @Override
                    public void run() {
                        try {
                            dispatch(rpcMessage.getId(), ctx, rpcMessage.getBody());
                        } catch (Throwable th) {
                            LOGGER.error(FrameworkErrorCode.NetDispatch.errCode, th.getMessage(), th);
                        }
                    }
                });
            } catch (RejectedExecutionException e) {
                //...
            }
        } else {
            MessageFuture messageFuture = futures.remove(rpcMessage.getId());
            if (messageFuture != null) {
                messageFuture.setResultMessage(rpcMessage.getBody());
            } else {
                try {
                    AbstractRpcRemoting.this.messageExecutor.execute(new Runnable() {
                        @Override
                        public void run() {
                            try {
                                dispatch(rpcMessage.getId(), ctx, rpcMessage.getBody());
                            } catch (Throwable th) {
                                LOGGER.error(FrameworkErrorCode.NetDispatch.errCode, th.getMessage(), th);
                            }
                        }
                    });
                } catch (RejectedExecutionException e) {
                    //...
                }
            }
        }
    }
}

public abstract void dispatch(long msgId, ChannelHandlerContext ctx, Object msg);

可以看到RpcMessage就是Fescar系統(tǒng)間交互的信息載體诅炉,所有的消息都是在此進(jìn)行處理(RpcServerOveride了此方法,對(duì)一些特殊消息進(jìn)行了提前處理屋厘,但最終還是調(diào)用此處的方法)涕烧。
所有的消息都是異步消費(fèi),交由dispatch方法處理汗洒。

TC處理消息

RpcServer的dispatch實(shí)現(xiàn)议纯,所有消息都委托給ServerMessageListener處理。

 // RpcServer
@Override
public void dispatch(long msgId, ChannelHandlerContext ctx, Object msg) {
    if (msg instanceof RegisterRMRequest) {
        serverMessageListener.onRegRmMessage(msgId, ctx, (RegisterRMRequest)msg, this,
            checkAuthHandler);
    } else {
        if (ChannelManager.isRegistered(ctx.channel())) {
            serverMessageListener.onTrxMessage(msgId, ctx, msg, this);
        } else {
                closeChannelHandlerContext(ctx);
        }
    }
}
  • 注冊(cè)類(lèi)消息
    對(duì)注冊(cè)TM和注冊(cè)RM的消息的處理主要是將channel和請(qǐng)求綁定(IP溢谤,PORT瞻凤,ResourceId等等),綁定關(guān)系使用ChannelManagerRpcContext管理世杀。最終ChannelManager將TM和RM的Channel緩存在兩個(gè)很復(fù)雜的map中鲫构。
    具體綁定過(guò)程比較復(fù)雜(有點(diǎn)亂。玫坛。)先忽略
// resourceId -> applicationId -> ip -> port -> RpcContext
private static final ConcurrentMap<String, ConcurrentMap<String, ConcurrentMap<String, ConcurrentMap<Integer, RpcContext>>>>
        RM_CHANNELS = new ConcurrentHashMap<>();
//ip+appname,port
private static final ConcurrentMap<String, ConcurrentMap<Integer, RpcContext>> TM_CHANNELS = new ConcurrentHashMap<>();

  • 事務(wù)消息
@Override
public void onTrxMessage(long msgId, ChannelHandlerContext ctx, Object message, ServerMessageSender sender) {
    RpcContext rpcContext = ChannelManager.getContextFromIdentified(ctx.channel());
    if (message instanceof MergedWarpMessage) {
        AbstractResultMessage[] results = new AbstractResultMessage[((MergedWarpMessage)message).msgs.size()];
        for (int i = 0; i < results.length; i++) {
            final AbstractMessage subMessage = ((MergedWarpMessage)message).msgs.get(i);
            results[i] = transactionMessageHandler.onRequest(subMessage, rpcContext);
        }
        MergeResultMessage resultMessage = new MergeResultMessage();
        resultMessage.setMsgs(results);
        sender.sendResponse(msgId, ctx.channel(), resultMessage);
    } else if (message instanceof AbstractResultMessage) {
        transactionMessageHandler.onResponse((AbstractResultMessage)message, rpcContext);
    }
}

層層委托后结笨,所有消息最終委托給DefaultCoordinatorDefaultCore處理,具體邏輯不在此文分析。

TM處理消息

準(zhǔn)確的說(shuō)TM主要是發(fā)送消息到TC(register炕吸、commit和rollback等)伐憾,對(duì)于入站消息的處理沒(méi)有特殊實(shí)現(xiàn)。

RM處理消息

RM除了需要發(fā)送消息到TC外赫模,需要處理TC的branchCommit和branchRollback消息树肃,參見(jiàn)RmMessageListener

    @Override
    public void onMessage(long msgId, String serverAddress, Object msg, ClientMessageSender sender) {
        if (msg instanceof BranchCommitRequest) {
            handleBranchCommit(msgId, serverAddress, (BranchCommitRequest)msg, sender);
        } else if (msg instanceof BranchRollbackRequest) {
            handleBranchRollback(msgId, serverAddress, (BranchRollbackRequest)msg, sender);
        }
    }

合并消息

發(fā)送消息時(shí),為了提高吞吐量瀑罗,F(xiàn)escar通過(guò)MergedSendRunnable來(lái)合并消息胸嘴,批量異步發(fā)送;合并后的Request消息為MergedWarpMessage


啟動(dòng)Fescar

最后提一下Fescar的啟動(dòng)方式斩祭。

public interface RemotingService {
    void start();
    void shutdown();
}

RemotingService定義了TM劣像、RM、TC的啟動(dòng)和關(guān)閉摧玫;不過(guò)實(shí)際的實(shí)現(xiàn)中耳奕,F(xiàn)escar是通過(guò)AbstractRpcRemoting的init()方法去完成服務(wù)的初始化,各種任務(wù)線程池的啟動(dòng),以及服務(wù)長(zhǎng)連接的建立诬像。
RmCpcCLient為例:

//RmRpcClient
@Override
public void init() {
    if (initialized.compareAndSet(false, true)) {
        super.init(); 
        timerExecutor.scheduleAtFixedRate(new Runnable() {
            @Override
            public void run() {
                reconnect(); //讀取配置屋群,獲取TC地址,連接TC
            }
        }, SCHEDULE_INTERVAL_MILLS, SCHEDULE_INTERVAL_MILLS, TimeUnit.SECONDS);
        ExecutorService mergeSendExecutorService = new ThreadPoolExecutor(MAX_MERGE_SEND_THREAD,
            MAX_MERGE_SEND_THREAD,
            KEEP_ALIVE_TIME, TimeUnit.MILLISECONDS,
            new LinkedBlockingQueue<Runnable>(),
            new NamedThreadFactory(getThreadPrefix(MERGE_THREAD_PREFIX), MAX_MERGE_SEND_THREAD));
        mergeSendExecutorService.submit(new MergedSendRunnable()); //啟動(dòng)消息合并發(fā)送的任務(wù)
    }
}

OK坏挠,對(duì)Fescar的RPC和消息模塊大致有了總體認(rèn)識(shí)了芍躏,后面再去看看Fescar是具體是怎么流程化去處理各種事務(wù)消息的。

最后編輯于
?著作權(quán)歸作者所有,轉(zhuǎn)載或內(nèi)容合作請(qǐng)聯(lián)系作者
  • 序言:七十年代末降狠,一起剝皮案震驚了整個(gè)濱河市纸肉,隨后出現(xiàn)的幾起案子,更是在濱河造成了極大的恐慌喊熟,老刑警劉巖柏肪,帶你破解...
    沈念sama閱讀 212,542評(píng)論 6 493
  • 序言:濱河連續(xù)發(fā)生了三起死亡事件,死亡現(xiàn)場(chǎng)離奇詭異芥牌,居然都是意外死亡烦味,警方通過(guò)查閱死者的電腦和手機(jī),發(fā)現(xiàn)死者居然都...
    沈念sama閱讀 90,596評(píng)論 3 385
  • 文/潘曉璐 我一進(jìn)店門(mén)壁拉,熙熙樓的掌柜王于貴愁眉苦臉地迎上來(lái)谬俄,“玉大人,你說(shuō)我怎么就攤上這事弃理±B郏” “怎么了?”我有些...
    開(kāi)封第一講書(shū)人閱讀 158,021評(píng)論 0 348
  • 文/不壞的土叔 我叫張陵痘昌,是天一觀的道長(zhǎng)钥勋。 經(jīng)常有香客問(wèn)我炬转,道長(zhǎng),這世上最難降的妖魔是什么算灸? 我笑而不...
    開(kāi)封第一講書(shū)人閱讀 56,682評(píng)論 1 284
  • 正文 為了忘掉前任扼劈,我火速辦了婚禮,結(jié)果婚禮上菲驴,老公的妹妹穿的比我還像新娘荐吵。我一直安慰自己,他們只是感情好赊瞬,可當(dāng)我...
    茶點(diǎn)故事閱讀 65,792評(píng)論 6 386
  • 文/花漫 我一把揭開(kāi)白布先煎。 她就那樣靜靜地躺著,像睡著了一般巧涧。 火紅的嫁衣襯著肌膚如雪薯蝎。 梳的紋絲不亂的頭發(fā)上,一...
    開(kāi)封第一講書(shū)人閱讀 49,985評(píng)論 1 291
  • 那天褒侧,我揣著相機(jī)與錄音良风,去河邊找鬼谊迄。 笑死闷供,一個(gè)胖子當(dāng)著我的面吹牛,可吹牛的內(nèi)容都是我干的统诺。 我是一名探鬼主播歪脏,決...
    沈念sama閱讀 39,107評(píng)論 3 410
  • 文/蒼蘭香墨 我猛地睜開(kāi)眼,長(zhǎng)吁一口氣:“原來(lái)是場(chǎng)噩夢(mèng)啊……” “哼粮呢!你這毒婦竟也來(lái)了婿失?” 一聲冷哼從身側(cè)響起,我...
    開(kāi)封第一講書(shū)人閱讀 37,845評(píng)論 0 268
  • 序言:老撾萬(wàn)榮一對(duì)情侶失蹤啄寡,失蹤者是張志新(化名)和其女友劉穎豪硅,沒(méi)想到半個(gè)月后,有當(dāng)?shù)厝嗽跇?shù)林里發(fā)現(xiàn)了一具尸體挺物,經(jīng)...
    沈念sama閱讀 44,299評(píng)論 1 303
  • 正文 獨(dú)居荒郊野嶺守林人離奇死亡懒浮,尸身上長(zhǎng)有42處帶血的膿包…… 初始之章·張勛 以下內(nèi)容為張勛視角 年9月15日...
    茶點(diǎn)故事閱讀 36,612評(píng)論 2 327
  • 正文 我和宋清朗相戀三年,在試婚紗的時(shí)候發(fā)現(xiàn)自己被綠了识藤。 大學(xué)時(shí)的朋友給我發(fā)了我未婚夫和他白月光在一起吃飯的照片砚著。...
    茶點(diǎn)故事閱讀 38,747評(píng)論 1 341
  • 序言:一個(gè)原本活蹦亂跳的男人離奇死亡,死狀恐怖痴昧,靈堂內(nèi)的尸體忽然破棺而出稽穆,到底是詐尸還是另有隱情,我是刑警寧澤赶撰,帶...
    沈念sama閱讀 34,441評(píng)論 4 333
  • 正文 年R本政府宣布舌镶,位于F島的核電站柱彻,受9級(jí)特大地震影響,放射性物質(zhì)發(fā)生泄漏乎折。R本人自食惡果不足惜绒疗,卻給世界環(huán)境...
    茶點(diǎn)故事閱讀 40,072評(píng)論 3 317
  • 文/蒙蒙 一、第九天 我趴在偏房一處隱蔽的房頂上張望骂澄。 院中可真熱鬧吓蘑,春花似錦、人聲如沸坟冲。這莊子的主人今日做“春日...
    開(kāi)封第一講書(shū)人閱讀 30,828評(píng)論 0 21
  • 文/蒼蘭香墨 我抬頭看了看天上的太陽(yáng)健提。三九已至琳猫,卻和暖如春,著一層夾襖步出監(jiān)牢的瞬間私痹,已是汗流浹背脐嫂。 一陣腳步聲響...
    開(kāi)封第一講書(shū)人閱讀 32,069評(píng)論 1 267
  • 我被黑心中介騙來(lái)泰國(guó)打工, 沒(méi)想到剛下飛機(jī)就差點(diǎn)兒被人妖公主榨干…… 1. 我叫王不留紊遵,地道東北人账千。 一個(gè)月前我還...
    沈念sama閱讀 46,545評(píng)論 2 362
  • 正文 我出身青樓,卻偏偏與公主長(zhǎng)得像暗膜,于是被迫代替她去往敵國(guó)和親匀奏。 傳聞我的和親對(duì)象是個(gè)殘疾皇子,可洞房花燭夜當(dāng)晚...
    茶點(diǎn)故事閱讀 43,658評(píng)論 2 350

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