基於RocketMQ實現分散式事務

来源:https://www.cnblogs.com/ayic/p/18067431
-Advertisement-
Play Games

背景 在一個微服務架構的項目中,一個業務操作可能涉及到多個服務,這些服務往往是獨立部署,構成一個個獨立的系統。這種分散式的系統架構往往面臨著分散式事務的問題。為了保證系統數據的一致性,我們需要確保這些服務中的操作要麼全部成功,要麼全部失敗。通過使用RocketMQ實現分散式事務,我們可以協調這些服務 ...


image

背景

在一個微服務架構的項目中,一個業務操作可能涉及到多個服務,這些服務往往是獨立部署,構成一個個獨立的系統。這種分散式的系統架構往往面臨著分散式事務的問題。為了保證系統數據的一致性,我們需要確保這些服務中的操作要麼全部成功,要麼全部失敗。通過使用RocketMQ實現分散式事務,我們可以協調這些服務的操作,保證數據的一致性。

功能原理

RocketMQ的分散式事務消息功能,在普通消息基礎上,支持二階段的提交。將二階段提交和本地事務綁定,實現全局提交結果的一致性。

整個事務消息的詳細交互流程如下圖所示:

image

1、生產者將消息發送至RocketMQ服務端。

2、RocketMQ服務端將消息持久化成功之後,向生產者返回Ack確認消息已經發送成功,此時消息被標記為"暫不能投遞",這種狀態下的消息即為半事務消息。

3、生產者開始執行本地事務邏輯。

4、生產者根據本地事務執行結果向服務端提交二次確認結果(Commit或是Rollback),服務端收到確認結果後處理邏輯如下:

  • 二次確認結果為Commit:服務端將半事務消息標記為可投遞,並投遞給消費者。

  • 二次確認結果為Rollback:服務端將回滾事務,不會將半事務消息投遞給消費者。

5、在斷網或者是生產者應用重啟的特殊情況下,若服務端未收到生產者提交的二次確認結果,或服務端收到的二次確認結果為Unknown未知狀態,經過固定時間後,服務端將對消息生產者集群中任一生產者實例發起消息回查。

6、生產者收到消息回查後,需要檢查對應消息的本地事務執行的最終結果。

7、生產者根據檢查到的本地事務的最終狀態再次提交二次確認,服務端仍按照步驟4對半事務消息進行處理。

註意問題

消息類型
事務消息僅支持在MessageType為Transaction的主題使用,即事務消息只能發送至類型為事務消息的主題中。

消息消費
RocketMQ事務消息保證生產者本地事務和下游消息發送事務的一致性,但不保證消息消費結果和上游事務的一致性。因此需要下游業務自行保證消息正確處理,建議消費端做好消費重試。

中間狀態
RocketMQ事務消息一致性為最終一致性,即在消息提交到下游消費端處理完成之前,下游和上游事務之間的狀態會不一致。因此,事務消息僅適合能接受非同步執行的場景。

事務超時
RocketMQ事務消息的生命周期存在超時機制,即半事務消息被生產者發送服務端後,如果在指定時間內服務端無法確認提交或者回滾狀態,則消息預設會被回滾。

示例代碼

以下為RocketMQ 4.x版本事務消息示例代碼,

import org.apache.rocketmq.client.producer.LocalTransactionState;
import org.apache.rocketmq.client.producer.TransactionListener;
import org.apache.rocketmq.client.producer.TransactionMQProducer;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.common.message.MessageExt;

import java.util.concurrent.*;

public class RocketMqTransactionDemo {
	public static void main(String[] args) throws Exception {
		// 創建事務消息生產者
		TransactionMQProducer producer = new TransactionMQProducer("transaction_producer");
		producer.setNamesrvAddr("127.0.0.1:9876");

		// 設置事務監聽器
		TransactionListener transactionListener = new MyTransactionListener();
		producer.setTransactionListener(transactionListener);

		// 設置事務回查的線程池,可以不必設置,如果不設置也會預設生成一個
		ExecutorService executorService = new ThreadPoolExecutor(2, 5, 100, TimeUnit.SECONDS, new ArrayBlockingQueue <Runnable> (2000), new ThreadFactory() {
			@Override
			public Thread newThread(Runnable r) {
				Thread thread = new Thread(r);
				thread.setName("client-transaction-msg-check-thread");
				return thread;
			}
		});
		producer.setExecutorService(executorService);

		// 啟動生產者
		producer.start();

		// 發送事務消息
		Message message = new Message("transaction_topic", "test_tag", "test_key", "Hello RocketMQ".getBytes());
		producer.sendMessageInTransaction(message, null);

		// 關閉生產者
		producer.shutdown();
	}
}

/**
 * 事務監聽器
 */
class MyTransactionListener implements TransactionListener {
	@Override
	public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
		// 執行本地事務操作
		System.out.println("執行本地事務操作,消息內容:" + new String(msg.getBody()));
		return LocalTransactionState.COMMIT_MESSAGE; // 提交事務,允許消費者消費該消息
		// return LocalTransactionState.ROLLBACK_MESSAGE;// 回滾事務,消息將被丟棄不允許消費。
		// return LocalTransactionState.UNKNOW;// 暫時無法判斷狀態,等待固定時間以後Broker端根據回查規則向生產者進行消息回查。
	}

	@Override
	public LocalTransactionState checkLocalTransaction(MessageExt msg) {
		// 檢查本地事務狀態
		System.out.println("檢查本地事務狀態,消息內容:" + new String(msg.getBody()));
		return LocalTransactionState.COMMIT_MESSAGE;
	}
}

代碼解釋:
1、事務消息的生產者使用TransactionMQProducer創建。
2、MyTransactionListener作為事務監聽器,實現了介面TransactionListener,該介面有兩個方法,分別是:

  • executeLocalTransaction
    半事務消息發送成功後,執行本地事務的方法,具體執行完本地事務後,可以在該方法中返回以下三種狀態:
    LocalTransactionState.COMMIT_MESSAGE: 提交事務,允許消費者消費該消息。
    LocalTransactionState.ROLLBACK_MESSAGE: 回滾事務,消息將被丟棄不允許消費。
    LocalTransactionState.UNKNOW: 暫時無法判斷狀態,等待固定時間以後RocketMQ服務端根據回查規則向生產者進行消息回查。

  • checkLocalTransaction
    二次確認消息沒有收到,RocketMQ服務端回查生產者端事務結果的方法。回查規則:本地事務執行完成後,若RocketMQ服務端收到的本地事務返回狀態為LocalTransactionState.UNKNOW,或生產者應用退出導致本地事務未提交任何狀態。則RocketMQ服務端會向消息生產者發起事務回查,第一次回查後仍未獲取到事務狀態,則之後每隔一段時間會再次回查。

本文來自博客園,作者:Y00,轉載請註明原文鏈接:https://www.cnblogs.com/ayic/p/18067431

聊聊技術,聊聊人生。歡迎關註我的公眾號!^_^


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

-Advertisement-
Play Games
更多相關文章
  • 本文詳細介紹了Libcomm通信庫及其原理,讓我們更好的理解GaussDB(DWS)集群通信中的具體邏輯,對於GaussDB(DWS)通信運維也具備一定的參考意義。 ...
  • 本文主要講解了 Compose 中狀態的概念。最後做個小結, - Compose UI 依賴狀態變化,觸發重組,驅動界面更新。 - 使用 remember 和 rememberSaveable 進行狀態持久化。remember 保證在 recompose 過程中狀態穩定,rememberSaveab... ...
  • 這裡給大家分享我在網上總結出來的一些知識,希望對大家有所幫助 一、介紹 模塊,(Module),是能夠單獨命名並獨立地完成一定功能的程式語句的集合(即程式代碼和數據結構的集合體)。 兩個基本的特征:外部特征和內部特征 外部特征是指模塊跟外部環境聯繫的介面(即其他模塊或程式調用該模塊的方式,包括有輸入 ...
  • 隨著 Vue 3 的發佈,組件通信成為了前端開發中一個值得關註的話題。本文介紹了 Vue 3 中幾種常見的組件通信方式,包括 Props 和 Events、事件匯流排、Provide 和 Inject,以及 Vuex 狀態管理。每種方式都有其適用場景和優缺點,開發者需要根據具體情況選擇最合適的方式。 ... ...
  • "node-sass": "^4.12.0", "sass-loader": "^8.0.2", 本地和local環境正常,pre和生產環境編譯報錯 local、pre、生產的編譯環境一樣,node版本都是14.16.1。拷貝本地node_modules文件夾到pre同樣報錯。 應該是node-sa ...
  • v-model 是 Vue.js 框架中用於實現雙向數據綁定的指令。它充分體現了 MVVM(Model-View-ViewModel)模式中的雙向數據綁定特性。下麵我們將詳細解釋 v-model 如何體現 MVVM 和雙向綁定: 1.MVVM 模式 MVVM 模式是一種軟體架構設計模式,它將應用程式 ...
  • 在你的 TypeScript 代碼中,當調用 nextPage_TopSelling() 或 prevPage_TopSelling() 方法時,雖然你更新了 currentPage_TopSelling 的值並調用了 reloadTopSelling() 方法,但是 Angular 並不會自動檢測 ...
  • 零售商家為什麼要建設線上商城 傳統的實體門店服務範圍有限,只能吸引周邊500米內的消費者。因此,如何拓展服務範圍,吸引更多消費者到店,成為了店家迫切需要解決的問題。 缺乏忠實顧客,客戶基礎不穩,往往是一次性購物,門店無法形成有效的顧客迴流。在當前的市場環境下,構建並維護粉絲群體,成為了商家的核心競爭 ...
一周排行
    -Advertisement-
    Play Games
  • 最近做項目過程中,使用到了海康相機,官方只提供了C/C++的SDK,沒有搜尋到一個合適的封裝了的C#庫,故自己動手,簡單的封裝了一下,方便大家也方便自己使用和二次開發 ...
  • 前言 MediatR 是 .NET 下的一個實現消息傳遞的庫,輕量級、簡潔高效,用於實現進程內的消息傳遞機制。它基於中介者設計模式,支持請求/響應、命令、查詢、通知和事件等多種消息傳遞模式。通過泛型支持,MediatR 可以智能地調度不同類型的消息,非常適合用於領域事件處理。 在本文中,將通過一個簡 ...
  • 前言 今天給大家推薦一個超實用的開源項目《.NET 7 + Vue 許可權管理系統 小白快速上手》,DncZeus的願景就是做一個.NET 領域小白也能上手的簡易、通用的後臺許可權管理模板系統基礎框架。 不管你是技術小白還是技術大佬或者是不懂前端Vue 的新手,這個項目可以快速上手讓我們從0到1,搭建自 ...
  • 第1章:WPF概述 本章目標 瞭解Windows圖形演化 瞭解WPF高級API 瞭解解析度無關性概念 瞭解WPF體繫結構 瞭解WPF 4.5 WPF概述 ​ 歡迎使用 Windows Presentation Foundation (WPF) 桌面指南,這是一個與解析度無關的 UI 框架,使用基於矢 ...
  • 在日常開發中,並不是所有的功能都是用戶可見的,還在一些背後默默支持的程式,這些程式通常以服務的形式出現,統稱為輔助角色服務。今天以一個簡單的小例子,簡述基於.NET開發輔助角色服務的相關內容,僅供學習分享使用,如有不足之處,還請指正。 ...
  • 第3章:佈局 本章目標 理解佈局的原則 理解佈局的過程 理解佈局的容器 掌握各類佈局容器的運用 理解 WPF 中的佈局 WPF 佈局原則 ​ WPF 視窗只能包含單個元素。為在WPF 視窗中放置多個元素並創建更貼近實用的用戶男面,需要在視窗上放置一個容器,然後在這個容器中添加其他元素。造成這一限制的 ...
  • 前言 在平時項目開發中,定時任務調度是一項重要的功能,廣泛應用於後臺作業、計劃任務和自動化腳本等模塊。 FreeScheduler 是一款輕量級且功能強大的定時任務調度庫,它支持臨時的延時任務和重覆迴圈任務(可持久化),能夠按秒、每天/每周/每月固定時間或自定義間隔執行(CRON 表達式)。 此外 ...
  • 目錄Blazor 組件基礎路由導航參數組件參數路由參數生命周期事件狀態更改組件事件 Blazor 組件 基礎 新建一個項目命名為 MyComponents ,項目模板的交互類型選 Auto ,其它保持預設選項: 客戶端組件 (Auto/WebAssembly): 最終解決方案裡面會有兩個項目:伺服器 ...
  • 先看一下效果吧: isChecked = false 的時候的效果 isChecked = true 的時候的效果 然後我們來實現一下這個效果吧 第一步:創建一個空的wpf項目; 第二步:在項目裡面添加一個checkbox <Grid> <CheckBox HorizontalAlignment=" ...
  • 在編寫上位機軟體時,需要經常處理命令拼接與其他設備進行通信,通常對不同的命令封裝成不同的方法,擴展稍許麻煩。 本次擬以特性方式實現,以兼顧維護性與擴展性。 思想: 一種命令對應一個類,其類中的各個屬性對應各個命令段,通過特性的方式,實現其在這包數據命令中的位置、大端或小端及其轉換為對應的目標類型; ...