摘要:本文結合 RocketMQ 原始碼,分享并發編程三大神器的相關知識點,
本文分享自華為云社區《讀 RocketMQ 原始碼,學習并發編程三大神器》,作者:勇哥java實戰分享,
這篇文章,筆者結合 RocketMQ 原始碼,分享并發編程三大神器的相關知識點,
1 CountDownLatch 實作網路同步請求
CountDownLatch 是一個同步工具類,用來協調多個執行緒之間的同步,它能夠使一個執行緒在等待另外一些執行緒完成各自作業之后,再繼續執行,
下圖是 CountDownLatch 的核心方法:
我們可以認為它內置一個計數器,建構式初始化計數值,每當執行緒執行 countDown 方法,計數器的值就會減一,當計數器的值為 0 時,表示所有的任務都執行完成,然后在 CountDownLatch 上等待的執行緒就可以恢復執行接下來的任務,
舉例,資料庫有100萬條資料需要處理,單執行緒執行比較慢,我們可以將任務分為5個批次,執行緒池按照每個批次執行,當5個批次整體執行完成后,列印出任務執行的時間 ,
long start = System.currentTimeMillis(); ExecutorService executorService = Executors.newFixedThreadPool(10); int batchSize = 5; CountDownLatch countDownLatch = new CountDownLatch(batchSize); for (int i = 0; i < batchSize; i++) { final int batchNumber = i; executorService.execute(new Runnable() { @Override public void run() { try { doSomething(batchNumber); } catch (Exception e) { e.printStackTrace(); } finally { countDownLatch.countDown(); } } }); } countDownLatch.await(); System.out.println("任務執行耗時:" + (System.currentTimeMillis() - start) + "毫秒");
溫習完 CountDownLatch 的知識點,回到 RocketMQ 原始碼,
筆者在沒有接觸網路編程之前,一直很疑惑,網路同步請求是如何實作的?
同步請求指:客戶端執行緒發起呼叫后,需要在指定的超時時間內,等到回應結果,才能完成本次呼叫,如果超時時間內沒有得到結果,那么會拋出超時例外,
RocketMQ 的同步發送訊息介面見下圖:
追蹤原始碼,真正發送請求的方法是通訊模塊的同步請求方法 invokeSyncImpl ,
整體流程:
- 發送訊息執行緒 Netty channel 物件呼叫 writeAndFlush 方法后 ,它的本質是通過 Netty 的讀寫執行緒將資料包發送到內核 , 這個程序本身就是異步的;
- ResponseFuture 類中內置一個 CountDownLatch 物件 ,responseFuture 物件呼叫 waitRepsone 方法,發送訊息執行緒會阻塞 ;
3.客戶端收到回應命令后, 執行 processResponseCommand 方法,核心邏輯是執行 ResponseFuture 的 putResponse 方法,
該方法的本質就是填充回應物件,并呼叫 countDownLatch 的 countDown 方法 , 這樣發送訊息執行緒就不再阻塞,
CountDownLatch 實作網路同步請求是非常實用的技巧,在很多開源中間件里,比如 Metaq ,Xmemcached 都有類似的實作,
2 ReadWriteLock 名字服務路由管理
讀寫鎖是一把鎖分為兩部分:讀鎖和寫鎖,其中讀鎖允許多個執行緒同時獲得,而寫鎖則是互斥鎖,
它的規則是:讀讀不互斥,讀寫互斥,寫寫互斥,適用于讀多寫少的業務場景,
我們一般都使用 ReentrantReadWriteLock ,該類實作了 ReadWriteLock ,ReadWriteLock 介面也很簡單,其內部主要提供了兩個方法,分別回傳讀鎖和寫鎖 ,
public interface ReadWriteLock { //獲取讀鎖 Lock readLock(); //獲取寫鎖 Lock writeLock(); }
讀寫鎖的使用方式如下所示:
1.創建 ReentrantReadWriteLock 物件 , 當使用 ReadWriteLock 的時候,并不是直接使用,而是獲得其內部的讀鎖和寫鎖,然后分別呼叫 lock / unlock 方法 ;
private ReadWriteLock readWriteLock = new ReentrantReadWriteLock();
2.讀取共享資料 ;
Lock readLock = readWriteLock.readLock(); readLock.lock(); try { // TODO 查詢共享資料 } finally { readLock.unlock(); }
3.寫入共享資料;
Lock writeLock = readWriteLock.writeLock(); writeLock.lock(); try { // TODO 修改共享資料 } finally { writeLock.unlock(); }
RocketMQ架構上主要分為四部分,如下圖所示 :

- Producer :訊息發布的角色,Producer 通過 MQ 的負載均衡模塊選擇相應的 Broker 集群佇列進行訊息投遞,投遞的程序支持快速失敗并且低延遲,
- Consumer :訊息消費的角色,支持以 push 推,pull 拉兩種模式對訊息進行消費,
- BrokerServer :Broker主要負責訊息的存盤、投遞和查詢以及服務高可用保證,
- NameServer :名字服務是一個非常簡單的 Topic 路由注冊中心,其角色類似 Dubbo 中的zookeeper,支持Broker的動態注冊與發現,
NameServer 是一個幾乎無狀態節點,可集群部署,節點之間無任何資訊同步,Broker 啟動之后會向所有 NameServer 定期(每 30s)發送心跳包(路由資訊),NameServer 會定期掃描 Broker 存活串列,如果超過 120s 沒有心跳則移除此 Broker 相關資訊,代表下線,
那么 NameServer 如何保存路由資訊呢?
路由資訊通過幾個 HashMap 來保存,當 Broker 向 Nameserver 發送心跳包(路由資訊),Nameserver 需要對 HashMap 進行資料更新,但我們都知道 HashMap 并不是執行緒安全的,高并發場景下,容易出現 CPU 100% 問題,所以更新 HashMap 時需要加鎖,RocketMQ 使用了 JDK 的讀寫鎖 ReentrantReadWriteLock ,
1.更新路由資訊,操作寫鎖
2.查詢主題資訊,操作讀鎖
讀寫鎖適用于讀多寫少的場景,比如名字服務,配置服務等,
3 CompletableFuture 異步訊息處理
RocketMQ 主從架構中,主節點與從節點之間資料同步/復制的方式有同步雙寫和異步復制兩種模式,
異步復制是指訊息在主節點落盤成功后就告訴客戶端訊息發送成功,無需等待訊息從主節點復制到從節點,訊息的復制由其他執行緒完成,
同步雙寫是指主節點將訊息成功落盤后,需要等待從節點復制成功,再告訴客戶端訊息發送成功,
同步雙寫模式是阻塞的,筆者按照 RocketMQ 4.6.1 原始碼,整理出主節點處理一個發送訊息的請求的時序圖,
整體流程:
- 生產者將訊息發送到 Broker , Broker 接收到訊息后,發送訊息處理器 SendMessageProcessor 的執行執行緒池 SendMessageExecutor 執行緒池來處理發送訊息命令;
- 執行 ComitLog 的 putMessage 方法;
- ComitLog 內部先執行 appendMessage 方法;
- 然后提交一個 GroupCommitRequest 到同步復制服務 HAService ,等待 HAService 通知 GroupCommitRequest 完成;
- 回傳寫入結果并回應客戶端 ,
我們可以看到:發送訊息的執行執行緒需要等待訊息復制從節點 , 并將訊息回傳給生產者才能開始處理下一個訊息,
RocketMQ 4.6.1 原始碼中,執行執行緒池的執行緒數量是 1 ,假如執行緒處理主從同步速度慢了,系統在這一瞬間無法處理新的發送訊息請求,造成 CPU 資源無法被充分利用 , 同時系統的吞吐量也會降低,
那么優化同步雙寫呢 ?
從 RocketMQ 4.7 開始,RocketMQ 引入了 CompletableFuture 實作了異步訊息處理 ,
- 發送訊息的執行執行緒不再等待訊息復制到從節點后再處理新的請求,而是提前生成 CompletableFuture 并回傳 ;
- HAService 中的執行緒在復制成功后,呼叫 CompletableFuture 的 complete 方法,通知 remoting 模塊回應客戶端(執行緒池:PutMessageExecutor ) ,
我們分析下 RocketMQ 4.9.4 核心代碼:
1.Broker 接收到訊息后,發送訊息處理器 SendMessageProcessor 的執行執行緒池 SendMessageExecutor 執行緒池來處理發送訊息命令;
2.呼叫 SendMessageProcessor 的 asyncProcessRequest 方法;
3.呼叫 Commitlog 的 aysncPutMessage 方法寫入訊息 ;
這段代碼中,當 commitLog 執行完 appendMessage 后, 需要執行刷盤任務和同步復制兩個任務,
但這兩個任務并不是同步執行,而是異步的方式,
4.復制執行緒復制訊息后,喚醒 future ;
5.組裝回應命令 ,并將回應命令回傳給客戶端,
為了便于理解這一段訊息發送處理程序的執行緒模型,筆者在 RocketMQ 原始碼中做了幾處埋點,修改 Logback 的日志配置,發送一條普通的訊息,觀察服務端日志,
從日志中,我們可以觀察到:
- 發送訊息的執行執行緒(圖中紅色)在執行完創建刷盤 Future 和同步復制 future 之后,并沒有等待這兩個任務執行完成,而是在結束 asyncProcessRequest 方法后就可以處理發送訊息請求了 ;
- 刷盤執行緒和復制執行緒執行完各自的任務后,喚醒 future,然后通過刷盤執行緒組裝存盤結果,最后通過 PutMessageExecutor 執行緒池(圖中黃色)將回應命令回傳給客戶端,
筆者一直認為:異步是更細粒度的使用系統資源的一種方式,在異步訊息處理的程序中,通過 CompletableFuture 這個神器,各個執行緒各司其職,優雅且高效的提升了 RocketMQ 的性能,
點擊關注,第一時間了解華為云新鮮技術~
轉載請註明出處,本文鏈接:https://www.uj5u.com/qita/539032.html
標籤:其他
上一篇:利用Seagate service獲得system shell
下一篇:編譯器優化丨Cache優化
