噓!非同步事件這樣用真的好麽?

来源:https://www.cnblogs.com/coderjava/archive/2020/06/29/13210405.html
-Advertisement-
Play Games

為了方便大家理解我把之前方案的圖片複製過來了,如下: 上圖的方案存在一個問題,就是我們今天文章要聊的內容。 這個問題就是當 MQ Consumer 收到消息後,就直接發佈 Event 了,如果是同步的,沒有問題。如果某個 EventListener 中處理失敗了,那麼這條消息將不會 ACK。 如果是 ...


為了方便大家理解我把之前方案的圖片複製過來了,如下:

img

上圖的方案存在一個問題,就是我們今天文章要聊的內容。

這個問題就是當 MQ Consumer 收到消息後,就直接發佈 Event 了,如果是同步的,沒有問題。如果某個 EventListener 中處理失敗了,那麼這條消息將不會 ACK。

如果是非同步發佈 Event 的場景,發佈完消息馬上就 ACK 了。就算某個 EventListener 中處理失敗了,MQ 也感知不到,不會進行消息的重新投遞,這就是存在的問題。

img

解決方案

方案一

既然消息已經 ACK 了,那就不利用 MQ 的重試功能了,使用方自己重試是不是也可以呢?

可肯定是可以的,內部處理是否成功肯定是可以知道的,如果處理失敗了可以預設重試,或者有一定策略的重試。實在不行還可以落庫,保存記錄。

這樣的問題在於太煩了呀,每個使用的地方都要去做這件事情,而且對於未來接手你代碼的程式小哥哥來說,這很有可能讓小哥哥頭髮慢慢脫落啊。。。。

脫落不要緊,關鍵他還不知道要做這個處理,說不定哪天就背鍋了,慘兮兮。。。。

方案二

要保證消息和業務處理的一致性,就不能立馬進行 ACK 操作。而是要等業務處理完成後再決定是否要 ACK。

如果有處理失敗的就不應該 ACK,這樣就能復用 MQ 的重試機制了。

分析下來,這就是一個典型的非同步轉同步的場景。像 Dubbo 中也有這個場景,所以我們可以借鑒 Dubbo 中的實現思路。

創建一個 DefaultFuture 用於同步等待獲取任務執行結果。然後在 MQ 消費的地方使用 DefaultFuture。

 @Service
 @RocketMQMessageListener(topic = "${rocketmq.topic.data_change}", consumerGroup = "${rocketmq.group.data_change_consumer}")
 public class DataChangeConsume implements RocketMQListener<DataChangeRequest> {
     @Autowired
     private ApplicationContext applicationContext;
     @Autowired
     private CustomApplicationContextAware customApplicationContextAware;
     @Override
     public void onMessage(DataChangeRequest dataChangeRequest) {
         log.info("received message {} , Thread {}", dataChangeRequest, Thread.currentThread().getName());
         DataChangeEvent event = new DataChangeEvent(this);
         event.setChangeType(dataChangeRequest.getChangeType());
         event.setTable(dataChangeRequest.getTable());
         event.setMessageId(dataChangeRequest.getMessageId());
         DefaultFuture defaultFuture = DefaultFuture.newFuture(dataChangeRequest, customApplicationContextAware.getTaskCount(), 6000 * 10);
         applicationContext.publishEvent(event);
         Boolean result = defaultFuture.get();
         log.info("MessageId {} 處理結果 {}", dataChangeRequest.getMessageId(), result);
         if (!result) {
             throw new RuntimeException("處理失敗,不進行消息ACK,等待下次重試");
         }
     }
 }

 

newFuture() 會傳入事件參數,超時時間,任務數量幾個參數。任務數量是用於判斷所有 EventListener 是否全部執行完成。

defaultFuture.get(); 這不就會阻塞,等待所有任務執行完成才會返回結果,如果所有業務都處理成功了,那麼會返回 true,流程結束,消息自動 ACK。

如果返回了 false 證明有處理失敗的或者超時的,就不需要 ACK 了,拋出異常等待重試。

 public Boolean get() {
     if (isDone()) {
         return true;
     }
     long start = System.currentTimeMillis();
     lock.lock();
     try {
         while (!isDone()) {
             done.await(timeout, TimeUnit.MILLISECONDS);
             // 有失敗的任務反饋
             if (!isSuccessDone()) {
                 return false;
             }
             // 全部執行成功
             if (isDone()) {
                 return true;
             }
             // 超時
             if (System.currentTimeMillis() - start > timeout) {
                 return false;
             }
         }
     } catch (InterruptedException e) {
         throw new RuntimeException(e);
     } finally {
         lock.unlock();
     }
     return true;
 }

 

isDone() 會判斷反饋結果了的任務數量跟總數量是否一致,如果一直就說明全部執行完成了。

 public boolean isDone() {
     return feedbackResultCount.get() == taskCount;
 }

 

那麼任務執行完了怎麼反饋呢? 不可能讓每個使用的方法去關心,所以我們定義了一個切麵來做這件事情。

 @Aspect
 @Component
 public class EventListenerAspect {
     @Around(value = "@annotation(eventListener)")
     public Object aroundAdvice(ProceedingJoinPoint joinpoint, EventListener eventListener) throws Throwable {
       DataChangeEvent event = null;
       boolean executeResult = true;
        try {
          event = (DataChangeEvent)joinpoint.getArgs()[0];
          Object result = joinpoint.proceed();
          return result;
       } catch (Exception e) {
          executeResult = false;
           throw e;
       } finally {
          DefaultFuture.received(event.getMessageId(), executeResult);
       }
     }
 }

 

通過 DefaultFuture.received() 反饋執行結果。

 
public static void received(String id, boolean result) {
     DefaultFuture future = FUTURES.get(id);
     if (future != null) {
         // 累加失敗任務數量
         if (!result) {
             future.feedbackFailResultCount.incrementAndGet();
         }
         // 累加執行完成任務數量
         future.feedbackResultCount.incrementAndGet();
         if (future.isDone()) {
             FUTURES.remove(id);
             future.doReceived();
         }
     }
 }
 private void doReceived() {
     lock.lock();
     try {
         if (done != null) {
             // 喚醒阻塞的線程
             done.signal();
         }
     } finally {
         lock.unlock();
     }
 }

 

下麵我們來總結整個流程:

  • 收到 MQ 消息,組裝成 DefaultFuture,通過 get 方法獲取執行結果,未執行完的時候此方法阻塞。

  • 通過切麵切入加了 EventListener 的方法,判斷是否有異常來判斷任務的執行結果。

  • 通過 DefaultFuture.received() 反饋結果。

  • 反饋時計算是否全部完成,全部完成則喚醒阻塞的線程。DefaultFuture.get()就能獲取到結果。

  • 是否要進行 ACK 操作。

需要註意的是每個 EventListener 內部消費的邏輯都要做冪等控制。

 


您的分享是我們最大的動力!

-Advertisement-
Play Games
更多相關文章
  • 13 約定 A common problem that arises when wrapping C libraries is that of maintaining reliability and checking for errors. The fact of the matter is tha ...
  • 構建生命周期 Maven的生命周期(lifecycle)可以理解為由Maven的各種plugin按照一定的順序執行來完成java項目清理、編譯、打包、測試、佈署等整個項目的流程的一個過程。 Maven內置了各種插件,如果再pom中沒有顯示配置就會調用預設的內置插件,如果pom中配置了就會調用配置的插 ...
  • 介面 恰當的原則是優先使用類而不是介面。從類開始,如果使用介面的必要性變得很明確,那麼就重構。介面是一個偉大的工具,但它們容易被濫用。 介面中可添加靜態方法與預設方法 一個類實現一個介面的同時必須實現該介面的所有方法(可以不用實現預設方法即關鍵詞為為 default的方法) extends 只能用於 ...
  • 小白是一名.net程式員,之前小白介紹了過了自己的博客系統http://www.ttblog.site/,用.net寫厭了,所以想學下java嘗嘗鮮,於是小白準備用spring boot來實現一個博客內容管理系統。 因為管理系統要有自己的數據源,但是又要從博客系統獲取博客內容,所以第一反應是要弄一個 ...
  • 初次看到原文我是有一些震撼的,原來作為開發人員,閑暇時間還算可以做這麼多有趣程式的開發。閱讀時暫且拋棄你所使用的語言的限制,你是否也能夠在“無聊”之時找到一個開發者的樂趣。 閱讀以下內容時重點關註項目的創意性,並結合自己的獨特經歷進行拓展,你一定也能夠找到編程的樂趣所在。很多項目都可以通過不同的技術 ...
  • 前言 本文的文字及圖片來源於網路,僅供學習、交流使用,不具有任何商業用途,版權歸原作者所有,如有問題請及時聯繫我們以作處理。 知識點 • 企業資產介紹 • 財務分析方法 • 企業資產數據爬取 • 企業資產數據展示 企業資產介紹 企業的資產包括流動資產、固定資產、無形資產、股東權益等等,本次給大家介紹 ...
  • 主要梳理一下SpringBoot2.x的依賴關係和依賴的版本管理,依賴版本管理是開發和管理一個SpringBoot項目的前提。 SpringBoot其實是通過starter的形式,對spring-framework進行裝箱,消除了(但是相容和保留)原來的XML配置,目的是更加便捷地集成其他框架,打造 ...
  • 1.迭代器 迭代器為我們提供了統一遍歷容器(List/Map/Set)的方式 1.遍歷List或Set 2.遍歷Map 2.Collections工具類 類java.util.Collections提供對Set、List、Map進行排序、填充、查找元素的輔助方法 1.void sort(List): ...
一周排行
    -Advertisement-
    Play Games
  • 概述:在C#中,++i和i++都是自增運算符,其中++i先增加值再返回,而i++先返回值再增加。應用場景根據需求選擇,首碼適合先增後用,尾碼適合先用後增。詳細示例提供清晰的代碼演示這兩者的操作時機和實際應用。 在C#中,++i 和 i++ 都是自增運算符,但它們在操作上有細微的差異,主要體現在操作的 ...
  • 上次發佈了:Taurus.MVC 性能壓力測試(ap 壓測 和 linux 下wrk 壓測):.NET Core 版本,今天計劃準備壓測一下 .NET 版本,來測試並記錄一下 Taurus.MVC 框架在 .NET 版本的性能,以便後續持續優化改進。 為了方便對比,本文章的電腦環境和測試思路,儘量和... ...
  • .NET WebAPI作為一種構建RESTful服務的強大工具,為開發者提供了便捷的方式來定義、處理HTTP請求並返迴響應。在設計API介面時,正確地接收和解析客戶端發送的數據至關重要。.NET WebAPI提供了一系列特性,如[FromRoute]、[FromQuery]和[FromBody],用 ...
  • 原因:我之所以想做這個項目,是因為在之前查找關於C#/WPF相關資料時,我發現講解圖像濾鏡的資源非常稀缺。此外,我註意到許多現有的開源庫主要基於CPU進行圖像渲染。這種方式在處理大量圖像時,會導致CPU的渲染負擔過重。因此,我將在下文中介紹如何通過GPU渲染來有效實現圖像的各種濾鏡效果。 生成的效果 ...
  • 引言 上一章我們介紹了在xUnit單元測試中用xUnit.DependencyInject來使用依賴註入,上一章我們的Sample.Repository倉儲層有一個批量註入的介面沒有做單元測試,今天用這個示例來演示一下如何用Bogus創建模擬數據 ,和 EFCore 的種子數據生成 Bogus 的優 ...
  • 一、前言 在自己的項目中,涉及到實時心率曲線的繪製,項目上的曲線繪製,一般很難找到能直接用的第三方庫,而且有些還是定製化的功能,所以還是自己繪製比較方便。很多人一聽到自己畫就害怕,感覺很難,今天就分享一個完整的實時心率數據繪製心率曲線圖的例子;之前的博客也分享給DrawingVisual繪製曲線的方 ...
  • 如果你在自定義的 Main 方法中直接使用 App 類並啟動應用程式,但發現 App.xaml 中定義的資源沒有被正確載入,那麼問題可能在於如何正確配置 App.xaml 與你的 App 類的交互。 確保 App.xaml 文件中的 x:Class 屬性正確指向你的 App 類。這樣,當你創建 Ap ...
  • 一:背景 1. 講故事 上個月有個朋友在微信上找到我,說他們的軟體在客戶那邊隔幾天就要崩潰一次,一直都沒有找到原因,讓我幫忙看下怎麼回事,確實工控類的軟體環境複雜難搞,朋友手上有一個崩潰的dump,剛好丟給我來分析一下。 二:WinDbg分析 1. 程式為什麼會崩潰 windbg 有一個厲害之處在於 ...
  • 前言 .NET生態中有許多依賴註入容器。在大多數情況下,微軟提供的內置容器在易用性和性能方面都非常優秀。外加ASP.NET Core預設使用內置容器,使用很方便。 但是筆者在使用中一直有一個頭疼的問題:服務工廠無法提供請求的服務類型相關的信息。這在一般情況下並沒有影響,但是內置容器支持註冊開放泛型服 ...
  • 一、前言 在項目開發過程中,DataGrid是經常使用到的一個數據展示控制項,而通常表格的最後一列是作為操作列存在,比如會有編輯、刪除等功能按鈕。但WPF的原始DataGrid中,預設只支持固定左側列,這跟大家習慣性操作列放最後不符,今天就來介紹一種簡單的方式實現固定右側列。(這裡的實現方式參考的大佬 ...