基于rocketmq事务消息的分布式事务

先看下图


在这里插入图片描述

以上图例展示了mq事务消息解决分布式事务的producer环节,consumer正常消费即可。

show your code

根据以上流程我们可以用rocketmq很简单的实现如下代码。为了减少部分业务代码入侵做了一点点封装;
以下项目基于springboot2.1.3,此处引入jdbc,大家需要注意,可以和mybatis、mybatis-plus共用事务管理器(想了解jdbc与mybatis如何共用事务管理器,自行百度)。假如你是jpa 或者habinate 你就不能这样封装。

<dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-aop</artifactId>
        </dependency>
        <dependency>
            <groupId>org.springframework</groupId>
            <artifactId>spring-jdbc</artifactId>
        </dependency>
        <dependency>
            <groupId>org.apache.rocketmq</groupId>
            <artifactId>rocketmq-spring-boot-starter</artifactId>
            <version>2.1.1</version>
        </dependency>

需在你的业务库插入如下表记录事务日志。

CREATE TABLE `transaction_log` (
  `id` int(20) NOT NULL AUTO_INCREMENT,
  `trx_id` varchar(128) NOT NULL,
  PRIMARY KEY (`id`),
  UNIQUE KEY `idx_transaction_id` (`trx_id`) USING BTREE
) ENGINE=InnoDB AUTO_INCREMENT=1 DEFAULT CHARSET=utf8mb4;

由于我们使用rocketmq-spring-boot-starter这个自动配置,所以我们可以直接使用RocketMQTemplate来发送消息,非常方便
demo中配置

rocketmq.name-server=127.0.0.1:9876
rocketmq.producer.group=tx_group
rocketmq.producer.send-message-timeout=3000

package com.xxx.fw.rocketmq.trx.core;

import com.xxx.fw.rocketmq.trx.config.TrxContextHolder;
import javax.annotation.Resource;
import org.apache.rocketmq.spring.core.RocketMQLocalTransactionListener;
import org.apache.rocketmq.spring.core.RocketMQLocalTransactionState;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.messaging.Message;

/**
 * @Description 抽象类让业务继承此类,并实现doBusiness方法,实现不同的业务定制
 * @Author 姚仲杰 
 * * @Date 2020/12/21 20:51
 */
public abstract class AbstractTransactionListener implements RocketMQLocalTransactionListener {
    private final static Logger LOGGER= LoggerFactory.getLogger(AbstractTransactionListener.class);
    @Resource
    TransactionLogService transactionLogService;

    public abstract void doBusiness(Object o);

    /**这个方法中执行本地事务
     * @param message 已发送到mq的事务消息
     * @param o 要保存到库的对象
     * @return
     */
    @Override
    public RocketMQLocalTransactionState executeLocalTransaction(Message message, Object o) {
        RocketMQLocalTransactionState state;
        TrxContextHolder.setTrxId(message.getHeaders().getId().toString());
        try{
            LOGGER.info("执行业务逻辑,trx_id:[{}]",message.getHeaders().getId());
            doBusiness(o);
            state=RocketMQLocalTransactionState.COMMIT;
            LOGGER.info("执行业务逻辑---[COMMIT]");
        }catch(Exception ex){
            LOGGER.info("消息事务回滚[ROLLBACK] {}",ex);
            state=RocketMQLocalTransactionState.ROLLBACK;
        }
        return state;
    }

    /**此方法提供统一回查
     * @param message 回查的消息数据
     * @return
     */
    @Override
    public RocketMQLocalTransactionState checkLocalTransaction(Message message) {
        RocketMQLocalTransactionState state=RocketMQLocalTransactionState.UNKNOWN;
            if (transactionLogService.query(message.getHeaders().getId().toString())>0){
                LOGGER.info("commit msg trx_id:{}",message.getHeaders().getId());
                state=RocketMQLocalTransactionState.COMMIT;
            }else{
                LOGGER.info("LocalTransaction review [UNKNOWN] will try again");
                state=RocketMQLocalTransactionState.UNKNOWN;
            }
        return state;
    }
}

自定义注解

package com.xxx.fw.rocketmq.trx.config;

import java.lang.annotation.ElementType;
import java.lang.annotation.Inherited;
import java.lang.annotation.Retention;
import java.lang.annotation.RetentionPolicy;
import java.lang.annotation.Target;

/**
 * @Description TODO
 * @Author 姚仲杰
 * @Date 2020/12/28 11:28
 */
@Target(ElementType.METHOD)
@Retention(RetentionPolicy.RUNTIME)
@Inherited
public @interface MqTrx {
}

切面处理

package com.xxx.fw.rocketmq.trx.aspect;

import com.xxx.fw.rocketmq.trx.config.MqTrx;
import com.xxx.fw.rocketmq.trx.config.TrxContextHolder;
import com.xxx.fw.rocketmq.trx.core.TransactionLogService;
import javax.annotation.Resource;
import org.aspectj.lang.JoinPoint;
import org.aspectj.lang.annotation.After;
import org.aspectj.lang.annotation.Aspect;
import org.aspectj.lang.annotation.Pointcut;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.core.annotation.Order;
import org.springframework.stereotype.Component;
import org.springframework.transaction.annotation.EnableTransactionManagement;

/**
 * @Description TODO
 * @Author 姚仲杰
 * @Date 2020/12/28 11:29
 */
 //此处需介绍下注解顺序,当前注解需包在事务注解中执行,注解先进后出所以让事务注解order为0 本注解顺序为1 则先进事务注解在进本注解
@EnableTransactionManagement(order = 0)
@Order(1)
@Component
@Aspect
public class TrxAspect {
    public static  final Logger LOGGER= LoggerFactory.getLogger(TrxAspect.class);
    @Resource
    TransactionLogService transactionLogService;

    @Pointcut("@annotation(com.xxx.fw.rocketmq.trx.config.MqTrx)")
    public void pointcut(){};

    @After("pointcut() && @annotation(mqTrx)")
    public void after(JoinPoint joinPoint, MqTrx mqTrx){
        try {
            String id = TrxContextHolder.getTrxId();
            transactionLogService.insert(id);
            LOGGER.info("事务消息日志入库成功,trx_id:[{}]", id);
        }finally {
            TrxContextHolder.remove();
        }
    }

}

mq事务上下文管理

package com.xxx.fw.rocketmq.trx.config;

import java.util.HashMap;
import java.util.Map;

/**
 * @Description TODO
 * @Author 姚仲杰
 * @Date 2020/12/28 11:42
 */
public class TrxContext {

    private ThreadLocal<Map<String,String>> threadLocal=new ThreadLocal<Map<String,String>>(){
        @Override
        protected Map<String, String> initialValue() {
            return new HashMap<String, String>();
        }
    };

    public String put(String key, String value) {
        return threadLocal.get().put(key, value);
    }

    public String get(String key) {
        return threadLocal.get().get(key);
    }

    public String remove(String key) {
        return threadLocal.get().remove(key);
    }

    public Map<String, String> entries() {
        return threadLocal.get();
    }
}

package com.xxx.fw.rocketmq.trx.config;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

/**
 * @Description 此类用于传递事务id到切面中,即发送完成后执行本地事务前,将trx_id放入本地线程,然后通过切面把事务id写入到事务日志表
 * @Author 姚仲杰
 * @Date 2020/12/28 11:46
 */
public class TrxContextHolder {
    private static final Logger LOGGER = LoggerFactory.getLogger(TrxContextHolder.class);

    public static final TrxContext TRX_CONTEXT_HOLDER=new TrxContext();

    public static final String TRX_ID="TRX_ID";

    public static String getTrxId(){
        return TRX_CONTEXT_HOLDER.get(TRX_ID);
    }

    public static void setTrxId(String trxId){
        if (LOGGER.isDebugEnabled()) {
            LOGGER.debug("set trx_id:[{}]", trxId);
        }
        TRX_CONTEXT_HOLDER.put(TRX_ID, trxId);

    }

    public static String remove() {
        String trxId = TRX_CONTEXT_HOLDER.remove(TRX_ID);
        if (LOGGER.isDebugEnabled()) {
            LOGGER.debug("remove trx_id:[{}] ", trxId);
        }
        return trxId;
    }

}

package com.xxx.fw.rocketmq.trx.core;

import javax.annotation.Resource;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Component;
import org.springframework.util.StringUtils;

/**
 * @Description 这里使用jdbcTemplate内置日志入库和查询操作让业务开发无需关注这些代码只需要添加一个注解即可
 * @Author 姚仲杰
 * @Date 2020/12/28 19:32
 */
@Component
public class TransactionLogService {
    private final static String INSERT_TRX_LOG="insert into transaction_log (trx_id)value('%s')";
    private final static String CHECK_TRX_LOG="select count(*) from transaction_log where trx_id='%s'";

    @Resource
    JdbcTemplate jdbcTemplate;

    public int insert(String trxId){
        if (StringUtils.isEmpty(trxId)){
            throw new IllegalArgumentException("trxId must not be null");
        }
        int update = jdbcTemplate.update(String.format(INSERT_TRX_LOG, trxId));
        return update;
    }

    public int query(String trxId){
        if (StringUtils.isEmpty(trxId)){
            throw new IllegalArgumentException("trxId must not be null");
        }
        int count = jdbcTemplate
            .queryForObject(String.format(CHECK_TRX_LOG, trxId), int.class);

        return count;
    }

}

提供封装的发送服务以及工具类

package com.xxx.fw.rocketmq.trx.core;

import com.xxx.fw.rocketmq.trx.util.TagsUtil;
import javax.annotation.Resource;
import org.apache.rocketmq.client.producer.TransactionSendResult;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Component;

/**
 * @Description 为了拼接tag实现不同业务
 * @Author 姚仲杰
 * @Date 2020/12/28 20:21
 */

@Component
public class TrxService {
    @Resource
    RocketMQTemplate rocketMQTemplate;

    /**
     * @param msg 发送给消息队列得消息
     * @param tag 消息队列得tag,组装后形成trx_topic:tag
     * @param o 要入库得数据
     * @return
     */
    public TransactionSendResult send(String msg,String tag,Object o){
        Message message = MessageBuilder.withPayload(msg).build();
        TransactionSendResult result = rocketMQTemplate
            .sendMessageInTransaction(TagsUtil.bindTag(tag), message, o);
        return result;
    }
}

工具类

package com.xxx.fw.rocketmq.trx.util;

import org.springframework.util.StringUtils;

/**
 * @Description 拼接tag
 * @Author 姚仲杰
 * @Date 2020/12/28 19:59
 */
public class TagsUtil {
    public static final String TRX_TOPIC="trx_topic:";

    public static String bindTag(String tag){
        if (StringUtils.isEmpty(tag)){
            throw new IllegalArgumentException("trx_tag must not be empty");
        }
        return TRX_TOPIC+tag;

    }
}

©著作权归作者所有,转载或内容合作请联系作者
  • 序言:七十年代末,一起剥皮案震惊了整个滨河市,随后出现的几起案子,更是在滨河造成了极大的恐慌,老刑警刘岩,带你破解...
    沈念sama阅读 213,992评论 6 493
  • 序言:滨河连续发生了三起死亡事件,死亡现场离奇诡异,居然都是意外死亡,警方通过查阅死者的电脑和手机,发现死者居然都...
    沈念sama阅读 91,212评论 3 388
  • 文/潘晓璐 我一进店门,熙熙楼的掌柜王于贵愁眉苦脸地迎上来,“玉大人,你说我怎么就摊上这事。” “怎么了?”我有些...
    开封第一讲书人阅读 159,535评论 0 349
  • 文/不坏的土叔 我叫张陵,是天一观的道长。 经常有香客问我,道长,这世上最难降的妖魔是什么? 我笑而不...
    开封第一讲书人阅读 57,197评论 1 287
  • 正文 为了忘掉前任,我火速办了婚礼,结果婚礼上,老公的妹妹穿的比我还像新娘。我一直安慰自己,他们只是感情好,可当我...
    茶点故事阅读 66,310评论 6 386
  • 文/花漫 我一把揭开白布。 她就那样静静地躺着,像睡着了一般。 火红的嫁衣衬着肌肤如雪。 梳的纹丝不乱的头发上,一...
    开封第一讲书人阅读 50,383评论 1 292
  • 那天,我揣着相机与录音,去河边找鬼。 笑死,一个胖子当着我的面吹牛,可吹牛的内容都是我干的。 我是一名探鬼主播,决...
    沈念sama阅读 39,409评论 3 412
  • 文/苍兰香墨 我猛地睁开眼,长吁一口气:“原来是场噩梦啊……” “哼!你这毒妇竟也来了?” 一声冷哼从身侧响起,我...
    开封第一讲书人阅读 38,191评论 0 269
  • 序言:老挝万荣一对情侣失踪,失踪者是张志新(化名)和其女友刘颖,没想到半个月后,有当地人在树林里发现了一具尸体,经...
    沈念sama阅读 44,621评论 1 306
  • 正文 独居荒郊野岭守林人离奇死亡,尸身上长有42处带血的脓包…… 初始之章·张勋 以下内容为张勋视角 年9月15日...
    茶点故事阅读 36,910评论 2 328
  • 正文 我和宋清朗相恋三年,在试婚纱的时候发现自己被绿了。 大学时的朋友给我发了我未婚夫和他白月光在一起吃饭的照片。...
    茶点故事阅读 39,084评论 1 342
  • 序言:一个原本活蹦乱跳的男人离奇死亡,死状恐怖,灵堂内的尸体忽然破棺而出,到底是诈尸还是另有隐情,我是刑警宁泽,带...
    沈念sama阅读 34,763评论 4 337
  • 正文 年R本政府宣布,位于F岛的核电站,受9级特大地震影响,放射性物质发生泄漏。R本人自食恶果不足惜,却给世界环境...
    茶点故事阅读 40,403评论 3 322
  • 文/蒙蒙 一、第九天 我趴在偏房一处隐蔽的房顶上张望。 院中可真热闹,春花似锦、人声如沸。这庄子的主人今日做“春日...
    开封第一讲书人阅读 31,083评论 0 21
  • 文/苍兰香墨 我抬头看了看天上的太阳。三九已至,却和暖如春,着一层夹袄步出监牢的瞬间,已是汗流浃背。 一阵脚步声响...
    开封第一讲书人阅读 32,318评论 1 267
  • 我被黑心中介骗来泰国打工, 没想到刚下飞机就差点儿被人妖公主榨干…… 1. 我叫王不留,地道东北人。 一个月前我还...
    沈念sama阅读 46,946评论 2 365
  • 正文 我出身青楼,却偏偏与公主长得像,于是被迫代替她去往敌国和亲。 传闻我的和亲对象是个残疾皇子,可洞房花烛夜当晚...
    茶点故事阅读 43,967评论 2 351

推荐阅读更多精彩内容