分布式锁与分布式事务
一、分布式锁
简介
锁的粒度。在开发时尽量缩小粒度,参考concarentHashMap的分段锁提高并发量。
异常处理。业务出异常后能释放锁吗,服务器异常宕机能释放锁吗
释放自己。释放锁时只能释放自己的锁
锁的续命。注意锁超时续命问题
X.1 Redis实现

整合redission框架
redisson开源项目地址。
redission的tryLock()底层怎么实现可重入的?

1、获取可重入/可重试锁,默认开启看门狗机制
- waitTime!=-1才开启重试机制,默认情况锁释放时间leaseTime=-1才开启看门狗机制;
- 获取锁成功后,看门狗默认是30秒后过期,递归地每到10秒续时一次;
- 没设置锁的最大超时时间可能会发生死锁;
- redission通过订阅信号量/等待/通知的方式等上一个线程释放锁后再尝试获取锁,获取失败重试,并且设置了锁的重试等待时间。




案例-释放可重入锁




源码-获取锁








源码-释放锁




获取可重入/不可重试锁,默认开启看门狗机制
- waitTime!=-1才开启重试机制,默认情况锁释放时间leaseTime=-1才开启看门狗机制;
- 获取锁成功后,看门狗默认是30秒后过期,递归地每到10秒续时一次;
- 没设置锁的最大超时时间可能会发生死锁。












获取可重入/可重试/超时释放锁,没得看门狗机制超时就关闭线程










4、主从一致性 - multiLock联锁原理
传统的主从架构redis的master宕机后通过哨兵选择主节点,会导致重复获取锁问题;此时redis架构应该采用多节点方式,通过多节点获取锁成功才算真正获取锁成功方式解决了该问题,也解决了redis的ha问题。实在不够再多节点模式加个主从模式。但这种模式需要至少部署3台主redis。











源码解读




异步秒杀案例
使用lua脚本在redis层面处理一部分逻辑后把结果放到消息队列,最后异步写入数据库。
在redis 的官方文档中有描述lua脚本在执行的时候不会执行其他脚本或Redis命令所以lua脚本具有排它性、原子性。但是存在的另一个问题是,它在执行的过程中如果一个命令报错不会回滚已执行的命令,所以要保证lua脚本的正确性。
而且他们有两个运行的函数call()和pcall(),两者区别在于,call在执行命令的时候如果报错,会抛出reisd错误终端执行,而pcall不会中断会记录下来错误信息。



1 下面缓存了一个商品编号为9的秒杀商品信息在redis



2 编写lua脚本操作redis


执行lua脚本,set保存抢到优惠券用户的信息







3 stream实现消息队列(也可用rocketmq),将消息放在队列中,基本介绍详见博文




4 开启线程去消息队列中拉取消息来处理





对于抛出异常导致未确认的消息会放到pendingList中,因此下面来再次处理未确认的消息,注意事务问题,下面下单时才对mysql进行库存操作。


5 单个测试




6 高并发测试




JMT压测
准备1000个用户

用户的token

秒杀库存为200

订单表清空

jmt配置






分段锁提高并发量
在开发时尽量缩小粒度,参考concarentHashMap的分段锁提高并发量。如下图将一个库存为200的商品拆分成多个key,每个key平分200的库存,此时并发时只需在每个key处10加锁即可,这样性能提高n倍。需要注意每段key的库存不够减时此时可以去操作其它key进行扣除。

X.2 Mysql实现
定段锁
一般用于定时任务到高可用而产生的锁,特点是记录锁名和锁的拥有者。锁名明确了要加/解的锁;锁的拥有者记录了拥有者信息,确保释放的锁是这任务加的锁。
源码
CREATE TABLE `zl_task_lock` (
`id` int(11) NOT NULL,
`lock_name` varchar(255) NOT NULL COMMENT '任务锁名称',
`lock_status` char(10) NOT NULL COMMENT '锁状态',
`lock_time` datetime DEFAULT NULL COMMENT '锁开始时间',
`lock_time_out` datetime DEFAULT NULL COMMENT '锁超时时间',
`lock_own` varchar(64) DEFAULT NULL COMMENT '锁拥有者',
PRIMARY KEY (`id`),
UNIQUE KEY `idx_lock_name` (`lock_name`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8;
@Mapper
public interface TaskLockDao extends BaseMapper<TaskLockEntity> {
/**
* 尝试获取锁
* @param lockName
* @param own
* @return
*/
boolean tryLock(@Param(value = "lockName") String lockName, @Param(value = "own") String own);
/**
* 获取超时时间锁
* @param lockName
* @param own
* @return
*/
boolean tryTimeOutLock(@Param(value = "lockName") String lockName, @Param(value = "own") String own);
/**
* 尝试释放锁
* @param lockName
* @param own
* @return
*/
boolean tryUnlock(@Param(value = "lockName") String lockName, @Param(value = "own") String own);
}
<mapper namespace="com.zlinepay.qr.cloud.services.withdrawal.dao.TaskLockDao">
<update id="tryLock" >
update zl_task_lock set lock_status = 'LOCK' ,lock_Time = NOW(), lock_time_out = date_add(now(), interval 30 MINUTE),lock_own = #{own} where
lock_name = #{lockName} and lock_status = 'UNLOCK'
</update>
<update id="tryTimeOutLock" >
update zl_task_lock set lock_status = 'LOCK' ,lock_Time = NOW(), lock_time_out = date_add(now(), interval 30 MINUTE),lock_own = #{own} where
lock_name = #{lockName} and now() <![CDATA[>=]]> lock_time_out
</update>
<update id="tryUnlock" >
update zl_task_lock set lock_status = 'UNLOCK' ,lock_Time = null , lock_time_out = null ,lock_own = null where
lock_name = #{lockName} and lock_status = 'LOCK' and lock_own = #{own}
</update>
</mapper>
public interface TaskLockService extends IService<TaskLockEntity> {
/**
* 获取锁
* @param response
* @param lockName
* @param action
* @param owner
* @return
* @throws IOException
*/
public boolean tryLock(HttpServletResponse response, String lockName, String action, String owner) throws IOException;
/**
* 释放锁
* @param owner
* @param taskLockStatus
* @param lockName
*/
public void tryUnlock(String owner, boolean taskLockStatus, String lockName);
}
@Service
@Slf4j
public class TaskLockServiceImpl extends ServiceImpl<TaskLockDao,TaskLockEntity> implements TaskLockService {
@Autowired
private TaskLockDao taskLockDao;
/**
* 获取锁
* @param response
* @param lockName
* @param action
* @param owner
* @return
* @throws IOException
*/
@Override
public boolean tryLock(HttpServletResponse response, String lockName, String action, String owner) throws IOException {
boolean taskLockStatus = false;
Rsp myRsp = new Rsp();
log.info("{} - {}",lockName,action);
taskLockStatus = tryLock(lockName,owner);
log.info("开始获取{}任务锁",action);
if(taskLockStatus){
log.info("获取{}锁成功",action);
myRsp.setData(action + "获取锁成功");
response.getWriter().write(myRsp.toString());
}else{
log.error("获取{}锁失败,任务正在被执行",action);
myRsp.setData(action + "获取锁失败,等待下一次获取");
response.getWriter().write(myRsp.toString());
}
return taskLockStatus;
}
/**
* 释放锁
* @param owner
* @param taskLockStatus
*/
@Override
public void tryUnlock(String owner, boolean taskLockStatus,String lockName) {
if(taskLockStatus){
boolean unlockResult = tryUnlock(lockName,owner);
if(!unlockResult){
log.error("释放锁失败,请检查失败原因");
}
}
}
public boolean tryLock(String lockName, String own) {
boolean tryLockResult = taskLockDao.tryLock(lockName, own);
if(!tryLockResult){
tryLockResult = taskLockDao.tryTimeOutLock(lockName,own);
}
return tryLockResult;
}
public boolean tryUnlock(String lockName, String own) {
TaskLockEntity query = new TaskLockEntity();
query.setLockOwn(own);
query.setLockName(lockName);
TaskLockEntity taskLockEntity = taskLockDao.selectOne(query);
if(taskLockEntity == null){
return false;
}else{
return taskLockDao.tryUnlock(lockName,own);
}
}
}
用法

扩展
public String getOwner(){
String owner = CodeUtils.generateCode();
if(StringUtils.isBlank(owner)){
log.error("获取序号失败");
throw new RuntimeException("获取锁序号失败");
}
return owner;
}
/**
* 生成流水工具类
*/
public class CodeUtils {
/**
* 生成流水号
* @param codePrefix 前缀
* @return
*/
public static String generateCode(String codePrefix) {
StringBuilder sb = new StringBuilder();
sb.append(codePrefix);
final IdGenerator idg = IdGenerator.INSTANCE;
sb.append( idg.nextId());
return sb.toString();
}
public static String generateCode() {
final IdGenerator idg = IdGenerator.INSTANCE;
return idg.nextId();
}
public static Long generateId(){
final IdGenerator idg = IdGenerator.INSTANCE;
String id = idg.nextId();
return Long.valueOf(id);
}
}
import org.apache.commons.lang.time.DateFormatUtils;
import java.net.InetAddress;
import java.net.NetworkInterface;
import java.net.SocketException;
import java.net.UnknownHostException;
import java.util.Enumeration;
import static java.net.InetAddress.getLocalHost;
/**
* 与snowflake算法区别,返回字符串id
*
* @Project concurrency
*/
public enum IdGenerator {
INSTANCE;
public static final String ip;
private long sequence = 0L;
private long sequenceBits = 12L; //序列号12位
private long sequenceMask = -1L ^ (-1L << sequenceBits); //4095
private long lastTimestamp = -1L;
IdGenerator(){
}
static {
ip = getLocalIP();
}
public synchronized String nextId() {
long timestamp = timeGen(); //获取当前毫秒数
//如果服务器时间有问题(时钟后退) 报错。
if (timestamp < lastTimestamp) {
throw new RuntimeException(String.format(
"Clock moved backwards. Refusing to generate id for %d milliseconds", lastTimestamp - timestamp));
}
//如果上次生成时间和当前时间相同,在同一毫秒内
if (lastTimestamp == timestamp) {
//sequence自增,因为sequence只有12bit,所以和sequenceMask相与一下,去掉高位
sequence = (sequence + 1) & sequenceMask;
//判断是否溢出,也就是每毫秒内超过4095,当为4096时,与sequenceMask相与,sequence就等于0
if (sequence == 0) {
timestamp = tilNextMillis(lastTimestamp); //自旋等待到下一毫秒
}
} else {
sequence = 0L; //如果和上次生成时间不同,重置sequence,就是下一毫秒开始,sequence计数重新从0开始累加
}
lastTimestamp = timestamp;
long suffix = sequence;
String suffixStr =String.valueOf(suffix);
// while(suffixStr.length() < 4){
// suffixStr = "0" + suffixStr ;
// }
//yyyy-MM-dd\'T\'HH:mm:ssZZ ,linux必须是标准的符号
String datePrefix = DateFormatUtils.format(timestamp, "yyyyMMddHHmmssSSS");
return datePrefix + ip + suffixStr;
}
protected long tilNextMillis(long lastTimestamp) {
long timestamp = timeGen();
while (timestamp <= lastTimestamp) {
timestamp = timeGen();
}
return timestamp;
}
protected long timeGen() {
return System.currentTimeMillis();
}
private String getLastIP(){
String lastip = null;
try{
lastip = getLocalHost().getHostAddress();
String[] split = lastip.split("\\.");
lastip = split[split.length-1];
} catch (UnknownHostException e) {
e.printStackTrace();
}
return lastip;
}
public static String getLocalIP() {
String lastip = null;
try{
if (isWindowsOS()) {
lastip = InetAddress.getLocalHost().getHostAddress();
} else {
lastip = getLinuxLocalIp();
}
String[] split = lastip.split("\\.");
lastip = split[split.length-1];
} catch (Exception e) {
e.printStackTrace();
System.out.println("获取Ip异常");
}
while (lastip.length()<3){
lastip = "0"+lastip;
}
return lastip;
}
private static String getLinuxLocalIp() throws SocketException {
String ip = "";
try {
for (Enumeration<NetworkInterface> en = NetworkInterface.getNetworkInterfaces(); en.hasMoreElements();) {
NetworkInterface intf = en.nextElement();
String name = intf.getName();
if (!name.contains("docker") && !name.contains("lo")) {
for (Enumeration<InetAddress> enumIpAddr = intf.getInetAddresses(); enumIpAddr.hasMoreElements();) {
InetAddress inetAddress = enumIpAddr.nextElement();
if (!inetAddress.isLoopbackAddress()) {
String ipaddress = inetAddress.getHostAddress().toString();
if (!ipaddress.contains("::") && !ipaddress.contains("0:0:") && !ipaddress.contains("fe80")) {
ip = ipaddress;
System.out.println(ipaddress);
}
}
}
}
}
} catch (SocketException ex) {
System.out.println("获取ip地址异常");
ip = "127.0.0.1";
ex.printStackTrace();
}
System.out.println("IP:"+ip);
return ip;
}
public static boolean isWindowsOS() {
boolean isWindowsOS = false;
String osName = System.getProperty("os.name");
if (osName.toLowerCase().indexOf("windows") > -1) {
isWindowsOS = true;
}
return isWindowsOS;
}
}
分段锁
核心是用商户号来做锁的标识,如果这个商户号没退过款(没生成过锁)就创建锁,如果退过则尝试获取锁。
CREATE TABLE `zl_merchant_available_refund_amount` (
`id` varchar(32) NOT NULL,
`merchant_no` varchar(64) DEFAULT NULL,
`trade_status` varchar(255) DEFAULT NULL COMMENT '商户交易状态 NORMAL 正常,LOCK 锁定',
`last_lock_time` datetime DEFAULT NULL ON UPDATE CURRENT_TIMESTAMP COMMENT '最后锁定时间',
PRIMARY KEY (`id`),
UNIQUE KEY `merchant_index` (`merchant_no`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8;
<?xml version="1.0" encoding="UTF-8"?>
<!DOCTYPE mapper PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" "http://mybatis.org/dtd/mybatis-3-mapper.dtd">
<mapper namespace="com.qr.order.dao.AvailableRefundAmountDao">
<!-- (上锁)-->
<update id="lock">
UPDATE
zl_merchant_available_refund_amount
SET
trade_status = 'LOCK',
last_lock_time = now()
WHERE merchant_no = #{merchantNo}
AND trade_status = 'NORMAL'
</update>
<update id="lockNowTime">
UPDATE
zl_merchant_available_refund_amount
SET
trade_status = 'LOCK',
last_lock_time = now()
WHERE merchant_no = #{merchantNo}
AND DATE_SUB(now(), INTERVAL 2 minute) > last_lock_time
</update>
<!-- (放锁)-->
<update id="unlock">
UPDATE
zl_merchant_available_refund_amount
SET
trade_status = 'NORMAL'
WHERE merchant_no = #{merchantNo}
AND trade_status = 'LOCK'
</update>
</mapper>
@Mapper
public interface AvailableRefundAmountDao extends BaseMapper<MerchantAvailableRefundAmount> {
/**
* 加锁
* @param merchantNo 商户号
* @return
*/
int lock( @Param(value = "merchantNo") String merchantNo);
int lockNowTime( @Param(value = "merchantNo") String merchantNo,@Param(value = "nowTime") Date nowTime);
/**
* 解锁
* @param merchantNo
*/
int unlock( String merchantNo);
}
@TableName("zl_merchant_available_refund_amount")
public class MerchantAvailableRefundAmount implements Serializable {
/*** */
@TableId
private String id;
/*** 商户编号*/
private String merchantNo;
/*** 商户交易状态 NORMAL 正常,LOCK 锁定*/
private String tradeStatus;
/*** 最后锁定时间*/
private Date lastLockTime;
public String getId() {
return id;
}
public void setId(String id) {
this.id = id;
}
public String getMerchantNo() {
return merchantNo;
}
public void setMerchantNo(String merchantNo) {
this.merchantNo = merchantNo;
}
public String getTradeStatus() {
return tradeStatus;
}
public void setTradeStatus(String tradeStatus) {
this.tradeStatus = tradeStatus;
}
public Date getLastLockTime() {
return lastLockTime;
}
public void setLastLockTime(Date lastLockTime) {
this.lastLockTime = lastLockTime;
}
}
import com.baomidou.mybatisplus.service.impl.ServiceImpl;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import java.math.BigDecimal;
import java.util.ArrayList;
import java.util.Date;
import java.util.List;
/**
* @version: v1.0
* @jdk version used: JDK1.8
*/
@Service
public class AvailableRefundAmountServiceImpl extends ServiceImpl<AvailableRefundAmountDao,MerchantAvailableRefundAmount> implements AvailableRefundAmountService {
private Logger logger = LoggerFactory.getLogger(getClass());
@Autowired
private AvailableRefundAmountDao availableRefundAmountDao;
/**
* @author wwei
* @description 单线程判断是否可退款
* @dateTime 2021/12/3/003 21:39
* @params
* @return true:允许退款,false:不允许退款
*/
private boolean getLockToRefund(TradePaymentRecordEntity tradePaymentRecordEntity) {
try {
// 1.如果这个商户没退过款(没生成过锁)就创建一个锁
MerchantAvailableRefundAmount queryAvailableRefundAmount = new MerchantAvailableRefundAmount();
queryAvailableRefundAmount.setMerchantNo(tradePaymentRecordEntity.getMerchantNo());
MerchantAvailableRefundAmount availableRefundAmount = availableRefundAmountDao.selectOne(queryAvailableRefundAmount);
if (availableRefundAmount == null) {
availableRefundAmount = new MerchantAvailableRefundAmount();
availableRefundAmount.setId(CodeUtils.generateCode());
availableRefundAmount.setMerchantNo(tradePaymentRecordEntity.getMerchantNo());//唯一索引
availableRefundAmount.setTradeStatus("NORMAL");
availableRefundAmount.setLastLockTime(new Date());
int result = availableRefundAmountDao.insert(availableRefundAmount);
if(result != 1){
availableRefundAmount = availableRefundAmountDao.selectOne(queryAvailableRefundAmount);//为获取超时锁做铺垫
if(availableRefundAmount == null){
return false;
}
}
}
int lockResult = availableRefundAmountDao.lock(tradePaymentRecordEntity.getMerchantNo());
if (lockResult != 1) {
lockResult = availableRefundAmountDao.lockNowTime(tradePaymentRecordEntity.getMerchantNo(), availableRefundAmount.getLastLockTime());
}
if (lockResult == 1){
// 2.获取锁后此处编写业务代码,返回true表示处理成功
return true;
}
}catch (Exception e){
logger.info("获取锁失败");
e.printStackTrace();
}finally {
availableRefundAmountDao.unlock(tradePaymentRecordEntity.getMerchantNo());
}
return false;
}
}
二、分布式事务
X.1 分布式事务方案
(1)Atomikos 是一个为Java平台提供增值服务的并且开源类事务管理器。表现形式为JTA两阶段提交,和seata的AT模式类似。
(2)Seata 是一款Java写的开源的分布式事务解决方案,致力于在微服务架构下提供高性能和简单易用的分布式事务服务。表现形式为AT和TCC模式。
X.1 选择AT还是TCC
个人建议使用TCC因为支持复杂sql且效率高,而且关系型数据库和非关系型数据库都很好的实现分布式事务的控制。


根据上述AT模式与TCC模式的对比,我们可以了解到,如果业务功能对性能要求很高,并且属于公司的核心业务,那么建议采用TCC模式;如果该功能对性能要求一般,建议采用AT模式;有一点需要提醒的是,AT模式和TCC模式是可以在一个项目中共存的,所以小伙伴们完全可以根据具体的业务需求选择任意的分布式事务解决方案;感兴趣的小伙伴还可以上github下载我准备好的相关案例:awesome-seata。
X.1 简介

源码
Seata生命周期


Seata阶段提交

一阶段过程-记录前后快照并提交


blob另存为拷贝下来看看

将文件中的内容拷贝到json在线解析
改之前

改之后

第二阶段过程-提交或回滚


可能存在的问题
- 如果服务宕机了回滚还没删完数据就gg了。
- 回滚时有个数据库超时,那么就会默认重试30次。
- 如何避免aba问题,他考虑的还是严谨,滴滴都在用,一般不会产生这一问题。
- 在下单时A线程还没提交,B线程可以修改库存吗?答案是不能,因为有个本地锁和和最终提交时向TC获取的全局事务锁在本地锁没释放其它线程是不能操作。
X.2 Seata TC Server单机环境部署
Seata TC server使用java写的程序。







这里可以配置存储方式


X.3 Seata TC Server集群环境部署





复制两个seata服务

拉取seata目录中的建表语句执行

配置数据库方式持久化事务信息




配置注册中心,服务注册可以不用密码

调整jvm

将上面配置好的seata服务再复制一份重命名(个人偷懒建议,如果不行则解压后分别配置即可)然后启动两台seata服务,查看已经注册到nacos
X.4 springboot-Seata基于AT事务模式-单个服务对应多数据源
案例:有三张表,分别在订单库的订单表、产品库的产品表、账户余额库的余额表。相当于是一个微服务连接了多个数据源。


表结构
有三张表,每个库都有undo_log表,是官网提供的,用来做回滚的。


依赖
数据源
配置分布式事务相关

分析

代码

扣库存-RM
减余额-RM
说明

测试
1、子数据库表操作的impl加了spring事务注解,如果在加了全局事务注解的方法中抛出异常则全局事务触发回滚有效。

2、子数据库表操作的impl去掉spring事务注解,如果在加了全局事务注解的方法中抛出异常则全局事务触发回滚有效。

3、在去掉spring事务注解的子数据库impl中抛出异常,全局事务回滚有效。
X.5 springcloud-Seat基于AT事务模式-微服务对应单个数据源
下面是基于spring alibaba微服务实现的一个分布式事务。就一个注解@GlobalTansction即可。
架构


依赖

下面是springcloud的seta依赖
配置
代码
发起方

被调用方
测试
1、fegin调用超时分布式事务不会回滚。
2、在调用方中触发异常,分布式事务触发回滚有效。
3、在被调用方内部抛异常,分布式事务回滚有效。
4、如果在try里面捕获了异常,那么就不会触发分布式事务的回滚这个特性和spring事务回滚机制一样。 也可以知道异常触发回滚。
X.6 seate的读写隔离
写隔离
防止并发,不会发生脏写。原理:
在二阶段提交或回滚全局事务后才会释放全局锁,tx1要回滚全局事务那么就必须先获取行数据的本地锁,由于本地锁被tx2占用而tx2要获取到全局事务锁才会释放本地锁,因此就会产生互相持有对方需要的资源,最终tx1会一直重试获取本地锁直到tx2获取全局锁超时从而释放本地锁使得tx1回滚成功。行数据的本地锁和全局锁互斥。

![]()

在二阶段提交或回滚全局事务后才会释放全局锁,tx1要回滚全局事务那么就必须先获取行数据的本地锁,由于本地锁被tx2占用而tx2要获取到全局事务锁才会释放本地锁,因此就会产生互相持有对方需要的资源,最终tx1会一直重试获取本地锁直到tx2获取全局锁超时从而释放本地锁使得tx1回滚成功。行数据的本地锁和全局锁互斥。


案例:
代码中有安照扣库存(微服务1)、减余额(微服务2)、下单(微服务3)的顺序操作都放在全局事务里,且减余额的微服务代码报错,现在来10个并发。执行结果如下:
此时第一个线程获取全局锁,执行到第二个微服务(减余额)就报异常,此时第一个线程就进入第二阶段进行回滚操作该操作要获取本地锁,但是由于是并发,扣库存的本地锁一直被线程二持有,而线程二一直获取全局锁,此时它们互相持有对方想要的锁,那么第一个线程会一直把线程获取全局锁的操作消耗超时。同理线程三/四/五都超时后线程二/三/四/五就释放了本地锁,从而线程一就得到了本地锁,最终线程一回滚成功。
读隔离
读已经提交的数据,如果读的数据在被其它事务修改中则一直等待该事务释放后才能读到数据,前提是select语句后面要加for update。


如果不加后面的for update就会直接读

X.7 springboot-Seata基于AT事务模式-单个服务对应多数据源
主要用于单应用连接多个数据源数据库操作的一个回滚



测试:重启服务测试成功

X.8 springcloud-Seata基于AT事务模式-微服务对应单个数据源
用于分布式微服务的回滚。
给每个服务加上配置




测试
- fegin超时,回滚成功
测试如果发现获取全局锁超时可以清掉数据库

X.9 分布式事务Seata-TCC原理
TCC工作原理和AT类似


![]()
X.10 springboot-Seat基于TCC事务模式-单个服务对应多数据源
TC服务搭建详见“AT事务模式”部分。
①配置
依赖和配置都参照“AT事务模式”Seata基于AT事务模式-单个服务对应多数据源”集群部分,变动的详见下面
service



impl




把要操作的数据库的service和impl如上写法即可。
注意:cancelTcc里面要控制幂等性,防止加多了。
②测试
所有操作弄完了,最后来个异常,这时触发自定义的回滚方法

X.11 springcloud-Seat基于TCC事务模式-微服务对应单数据源
TC服务搭建详见“AT事务模式”部分。
配置
依赖和配置都参照“AT事务模式”的“Seata基于AT事务模式-微服务对应单个数据源”集群部分,变动的详见下面
将上面的tcc单应用写法搬到每个微服务的service和impl即可。
测试
测试结果同上面的"单应用"测试结果
幂等性
在自定义的回滚中检验幂等性,防止回滚时加多了库存,解决方案是用分布式锁。


X.11 分布式事务Atomikos-JTA模式
下面是单应用操作多个数据库的案例。
CREATE TABLE `t_student` (
`n_id` int(11) NOT NULL AUTO_INCREMENT,
`c_name` varchar(255) DEFAULT NULL,
`c_age` int(12) DEFAULT NULL,
PRIMARY KEY (`n_id`) USING BTREE
) ENGINE=InnoDB AUTO_INCREMENT=2 DEFAULT CHARSET=utf8;
CREATE TABLE `t_teacher` (
`n_id` int(11) NOT NULL AUTO_INCREMENT,
`c_name` varchar(255) DEFAULT NULL,
PRIMARY KEY (`n_id`)
) ENGINE=InnoDB AUTO_INCREMENT=2 DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_0900_ai_ci;

<!--spring boot 版本依赖-->
<parent>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-parent</artifactId>
<version>2.1.4.RELEASE</version>
<relativePath/> <!-- lookup parent from repository -->
</parent>
<dependencies>
<!-- spring boot 对 web的依赖 可用 spring-boot:run启动 -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
</dependency>
<!-- spring boot 测试 -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>
<!-- 测试包 -->
<dependency>
<groupId>junit</groupId>
<artifactId>junit</artifactId>
<version>4.12</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.apache.commons</groupId>
<artifactId>commons-lang3</artifactId>
</dependency>
<!-- mybatis -->
<dependency>
<groupId>org.mybatis.spring.boot</groupId>
<artifactId>mybatis-spring-boot-starter</artifactId>
<version>1.3.1</version>
</dependency>
<!--mapper -->
<!-- https://mvnrepository.com/artifact/tk.mybatis/mapper-spring-boot-starter-->
<dependency>
<groupId>tk.mybatis</groupId>
<artifactId>mapper-spring-boot-starter</artifactId>
<version>2.1.5</version>
</dependency>
<!--pagehelper -->
<!-- https://mvnrepository.com/artifact/com.github.pagehelper/pagehelper-spring-boot-starter-->
<dependency>
<groupId>com.github.pagehelper</groupId>
<artifactId>pagehelper-spring-boot-starter</artifactId>
<version>1.2.10</version>
</dependency>
<!-- MySQL 连接驱动依赖 -->
<!-- https://mvnrepository.com/artifact/mysql/mysql-connector-java -->
<dependency>
<groupId>mysql</groupId>
<artifactId>mysql-connector-java</artifactId>
<version>5.1.47</version>
</dependency>
<dependency>
<groupId>com.alibaba</groupId>
<artifactId>druid</artifactId>
<version>1.1.10</version>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-configuration-processor</artifactId>
<optional>true</optional>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-jta-atomikos</artifactId>
</dependency>
</dependencies>
#配置端口
server:
port: 8080
#配置数据源
spring:
#student的配置
student:
uniqueResourceName: studentDatasource
jdbcUrl: jdbc:mysql://localhost:3306/db_student?characterEncoding=utf-8&useSSL=false
username: root
password: 123456
poolSize: 10
dataSourceClassName: com.mysql.jdbc.jdbc2.optional.MysqlXADataSource
#teacher配置
teacher:
uniqueResourceName: teacherDatasource
jdbcUrl: jdbc:mysql://localhost:3306/db_teacher?characterEncoding=utf-8&useSSL=false
username: root
password: 123456
poolSize: 10
dataSourceClassName: com.mysql.jdbc.jdbc2.optional.MysqlXADataSource
student数据源配置
@ConfigurationProperties(prefix = "spring.student")
@Data
public class StudentDataSourceProperties {
/**
* 数据源唯一资源名
*/
private String uniqueResourceName;
/**
* jdbc链接URL
*/
private String jdbcUrl;
/**
* 用户名
*/
private String username;
/**
* 密码
*/
private String password;
/**
* 连接池
*/
private Integer poolSize;
/**
* 数据源类名
*/
private String dataSourceClassName;
}
@Configuration
@Import(StudentDataSourceProperties.class)
@MapperScan(basePackages = "com.spring.student.mapper", sqlSessionFactoryRef = "studentSqlSessionFactory")
public class StudentDatasourceConfig {
@Autowired
private StudentDataSourceProperties studentProperties;
/**
* 配置AtomikosDataSourceBean数据源
*
*/
@Bean(name = "studentDatasource")
@Primary
public DataSource studentDataSource() {
//属性
Properties properties = new Properties();
properties.put("URL",studentProperties.getJdbcUrl());
properties.put("user", studentProperties.getUsername());
properties.put("password", studentProperties.getPassword());
//数据源
AtomikosDataSourceBean atomikosDataSourceBean = new AtomikosDataSourceBean();
atomikosDataSourceBean.setUniqueResourceName(studentProperties.getUniqueResourceName());
atomikosDataSourceBean.setXaDataSourceClassName(studentProperties.getDataSourceClassName());
atomikosDataSourceBean.setPoolSize(studentProperties.getPoolSize());
atomikosDataSourceBean.setXaProperties(properties);
return atomikosDataSourceBean;
}
/**
* 获取sqlSessionFactory
*/
@Bean
public SqlSessionFactory studentSqlSessionFactory() throws Exception {
SqlSessionFactoryBean factoryBean = new SqlSessionFactoryBean();
factoryBean.setDataSource(studentDataSource());
return factoryBean.getObject();
}
/**
* 获取会话模板
*/
@Bean
public SqlSessionTemplate studentSqlSessionTemplate() throws Exception {
return new SqlSessionTemplate(studentSqlSessionFactory());
}
/*
* 使用这个来做总事务 后面的数据源就不用设置事务了
* */
@Bean(name = "transactionManager")
@Primary
public JtaTransactionManager regTransactionManager () {
UserTransactionManager userTransactionManager = new UserTransactionManager();
UserTransaction userTransaction = new UserTransactionImp();
return new JtaTransactionManager(userTransaction, userTransactionManager);
}
}
Teacher数据源配置
@ConfigurationProperties(prefix = "spring.teacher")
@Data
public class TeacherDataSourceProperties {
/**
* 数据源唯一资源名
*/
private String uniqueResourceName;
/**
* jdbc链接URL
*/
private String jdbcUrl;
/**
* 用户名
*/
private String username;
/**
* 密码
*/
private String password;
/**
* 连接池
*/
private Integer poolSize;
/**
* 数据源类名
*/
private String dataSourceClassName;
}
@Configuration
@Import(TeacherDataSourceProperties.class)
@MapperScan(basePackages = "com.spring.teacher.mapper", sqlSessionFactoryRef = "teacherSqlSessionFactory")
public class TeacherDatasourceConfig {
@Autowired
private TeacherDataSourceProperties teacherProperties;
/**
* 获取数据源
* ConfigurationProperties:读取spring.datasource01的数据源
*/
@Bean(name = "teacherDatasource")
public DataSource teacherDataSource() {
//属性
Properties properties = new Properties();
properties.put("URL",teacherProperties.getJdbcUrl());
properties.put("user", teacherProperties.getUsername());
properties.put("password", teacherProperties.getPassword());
//数据源
AtomikosDataSourceBean atomikosDataSourceBean = new AtomikosDataSourceBean();
atomikosDataSourceBean.setUniqueResourceName(teacherProperties.getUniqueResourceName());
atomikosDataSourceBean.setXaDataSourceClassName(teacherProperties.getDataSourceClassName());
atomikosDataSourceBean.setPoolSize(teacherProperties.getPoolSize());
atomikosDataSourceBean.setXaProperties(properties);
return atomikosDataSourceBean;
}
/**
* 获取sqlSessionFactory
*/
@Bean
public SqlSessionFactory teacherSqlSessionFactory() throws Exception {
SqlSessionFactoryBean factoryBean = new SqlSessionFactoryBean();
factoryBean.setDataSource(teacherDataSource());
return factoryBean.getObject();
}
/**
* 获取会话模板
*/
@Bean
public SqlSessionTemplate teacherSqlSessionTemplate() throws Exception {
return new SqlSessionTemplate(teacherSqlSessionFactory());
}
}
entity
@Data
@Table(name = "t_student")
@AllArgsConstructor
@NoArgsConstructor
public class Student {
@Id
@Column(name = "n_id")
private Integer id;
@Column(name = "c_name")
private String name;
@Column(name = "c_age")
private Integer age;
}
@Data
@AllArgsConstructor
@NoArgsConstructor
@Table(name = "t_teacher")
public class Teacher {
@Id
@Column(name = "n_id")
private Integer id;
@Column(name = "c_name")
private String name;
}
mapper
@Mapper
public interface StudentMapper extends BaseMapper<Student> {
}
@Mapper
public interface TeacherMapper extends BaseMapper<Teacher> {
}
service
@Service
public class CommonService {
@Autowired
private StudentMapper studentMapper;
@Autowired
private TeacherMapper teacherMapper;
@Transactional
public void addSuccessTest(){
studentMapper.insertSelective(new Student(null, "张三",15));
teacherMapper.insertSelective(new Teacher(null,"张三"));
}
@Transactional
public void addRollbackTest(){
studentMapper.insertSelective(new Student(null, "张三",15));
teacherMapper.insertSelective(new Teacher(null,"张三"));
int a=1/0;
}
}
启动类
@SpringBootApplication
public class MybatisApplicationContext {
public static void main(String[] args) {
SpringApplication.run(MybatisApplicationContext.class);
}
}
测试 addSuccessTest
@RunWith(SpringRunner.class)
//主application方法
@SpringBootTest(classes=MybatisApplicationContext.class)
public class MybatisTest {
@Autowired
private CommonService commonService;
@Test
public void addSuccessTest(){
commonService.addSuccessTest();
}
@Test
public void addRollbackTest(){
commonService.addRollbackTest();
}
}



测试 addRollbackTest,回滚成功



更多推荐







所有评论(0)