tcc分布式事务源码解析系列(四)之项目实战

通过之前的几篇文章我相信您已经搭建好了运行环境,本次的项目实战是依照happylifeplat-tcc-demo项目来演练,也是非常经典的分布式事务场景:支付成功,进行订单状态的更新,扣除用户账户,库存扣减这几个模块来进行tcc分布式事务。话不多说,让我们一起进入体验吧!

  • 首先我们找到 PaymentServiceImpl.makePayment 方法,这是tcc分布式事务的发起者

    //tcc分布式事务的发起者 @Tcc(confirmMethod = "confirmOrderStatus", cancelMethod = "cancelOrderStatus") public void makePayment(Order order) { order.setStatus(OrderStatusEnum.PAYING.getCode()); orderMapper.update(order); //扣除用户余额 远端的rpc调用 AccountDTO accountDTO = new AccountDTO(); accountDTO.setAmount(order.getTotalAmount()); accountDTO.setUserId(order.getUserId()); accountService.payment(accountDTO); //进入扣减库存操作 远端的rpc调用 InventoryDTO inventoryDTO = new InventoryDTO(); inventoryDTO.setCount(order.getCount()); inventoryDTO.setProductId(order.getProductId()); inventoryService.decrease(inventoryDTO); } //更新订单的confirm方法 public void confirmOrderStatus(Order order) { order.setStatus(OrderStatusEnum.PAY_SUCCESS.getCode()); orderMapper.update(order); LOGGER.info("=========进行订单confirm操作完成================");

    1} 2//更新订单的cancel方法 3public void cancelOrderStatus(Order order) { 4 order.setStatus(OrderStatusEnum.PAY_FAIL.getCode()); 5 orderMapper.update(order); 6 LOGGER.info("=========进行订单cancel操作完成================"); 7}
  • makePayment 方法是tcc分布式事务的发起者,它里面有更新订单状态(本地),进行库存扣减(rpc扣减),资金扣减(rpc调用)

  • confirmOrderStatus 为订单状态confirm方法,cancelOrderStatus 为订单状态cancel方法。

  • 我们重点关注 @Tcc(confirmMethod = "confirmOrderStatus", cancelMethod = "cancelOrderStatus") 这个注解,为什么加了这个注解就有神奇的作用。

@Tcc注解切面

  • 我们找到在happylifeplat-tcc-core包中找到 TccTransactionAspect,这里定义@Tcc的切点。

    @Aspect public abstract class TccTransactionAspect {

    1private TccTransactionInterceptor tccTransactionInterceptor; 2 3public void setTccTransactionInterceptor(TccTransactionInterceptor tccTransactionInterceptor) { 4 this.tccTransactionInterceptor = tccTransactionInterceptor; 5} 6 7@Pointcut("@annotation(com.happylifeplat.tcc.annotation.Tcc)") 8public void txTransactionInterceptor() { 9 10} 11 12@Around("txTransactionInterceptor()") 13public Object interceptCompensableMethod(ProceedingJoinPoint proceedingJoinPoint) throws Throwable { 14 return tccTransactionInterceptor.interceptor(proceedingJoinPoint); 15} 16 17public abstract int getOrder();

    }

  • 我们可以知道Spring实现类的方法凡是加了@Tcc注解的,在调用的时候,都会进行 tccTransactionInterceptor.interceptor 调用。

  • 其次我们发现该类是一个抽象类,肯定会有其他类继承它,你猜的没错,对应dubbo用户,他的继承类为:

    @Aspect @Component public class DubboTccTransactionAspect extends TccTransactionAspect implements Ordered {

    1@Autowired 2public DubboTccTransactionAspect(DubboTccTransactionInterceptor dubboTccTransactionInterceptor) { 3 super.setTccTransactionInterceptor(dubboTccTransactionInterceptor); 4} 5 6 7@Override 8public int getOrder() { 9 return Ordered.HIGHEST_PRECEDENCE; 10}

    }

  • springcloud的用户,它的继承类为:

    @Aspect @Component public class SpringCloudTxTransactionAspect extends TccTransactionAspect implements Ordered {

    1@Autowired 2public SpringCloudTxTransactionAspect(SpringCloudTxTransactionInterceptor springCloudTxTransactionInterceptor) { 3 this.setTccTransactionInterceptor(springCloudTxTransactionInterceptor); 4} 5 6 7public void init() { 8 9} 10 11@Override 12public int getOrder() { 13 return Ordered.HIGHEST_PRECEDENCE; 14}

    }

  • 我们注意到他们都实现了Spring的Ordered接口,并重写了 getOrder 方法,都返回了 Ordered.HIGHEST_PRECEDENCE 那么可以知道,他是优先级最高的切面。

TccCoordinatorMethodAspect 切面

1@Aspect 2@Component 3public class TccCoordinatorMethodAspect implements Ordered { 4 5 6 private final TccCoordinatorMethodInterceptor tccCoordinatorMethodInterceptor; 7 8 @Autowired 9 public TccCoordinatorMethodAspect(TccCoordinatorMethodInterceptor tccCoordinatorMethodInterceptor) { 10 this.tccCoordinatorMethodInterceptor = tccCoordinatorMethodInterceptor; 11 } 12 13 14 @Pointcut("@annotation(com.happylifeplat.tcc.annotation.Tcc)") 15 public void coordinatorMethod() { 16 17 } 18 19 @Around("coordinatorMethod()") 20 public Object interceptCompensableMethod(ProceedingJoinPoint proceedingJoinPoint) throws Throwable { 21 return tccCoordinatorMethodInterceptor.interceptor(proceedingJoinPoint); 22 } 23 24 25 @Override 26 public int getOrder() { 27 return Ordered.HIGHEST_PRECEDENCE + 1; 28 } 29} 30
  • 该切面是第二优先级的切面,意思就是当我们对有@Tcc注解的方法发起调用的时候,首先会进入 TccTransactionAspect 切面,然后再进入 TccCoordinatorMethodAspect 切面。 该切面主要是用来获取@Tcc注解上的元数据,比如confrim方法名称等等。

  • 到这里如果大家基本能明白的话,差不多就了解了整个框架原理,现在让我们来跟踪调用吧。

现在我们回过头,我们来调用 PaymentServiceImpl.makePayment 方法

  • dubbo用户,我们会进入如下拦截器:

    @Component public class DubboTccTransactionInterceptor implements TccTransactionInterceptor {

    1private final TccTransactionAspectService tccTransactionAspectService; 2 3@Autowired 4public DubboTccTransactionInterceptor(TccTransactionAspectService tccTransactionAspectService) { 5 this.tccTransactionAspectService = tccTransactionAspectService; 6} 7 8 9@Override 10public Object interceptor(ProceedingJoinPoint pjp) throws Throwable { 11 final String context = RpcContext.getContext().getAttachment(Constant.TCC_TRANSACTION_CONTEXT); 12 TccTransactionContext tccTransactionContext = null; 13 if (StringUtils.isNoneBlank(context)) { 14 tccTransactionContext = 15 GsonUtils.getInstance().fromJson(context, TccTransactionContext.class); 16 } 17 return tccTransactionAspectService.invoke(tccTransactionContext, pjp); 18}

    }

  • 我们继续跟踪 tccTransactionAspectService.invoke(tccTransactionContext, pjp) ,发现通过判断过后,我们会进入 StartTccTransactionHandler.handler 方法,我们已经进入了分布式事务处理的入口了!

    /** * 分布式事务处理接口 * * @param point point 切点 * @param context 信息 * @return Object * @throws Throwable 异常 */ @Override public Object handler(ProceedingJoinPoint point, TccTransactionContext context) throws Throwable { Object returnValue; try { //开启分布式事务 tccTransactionManager.begin(); try { //发起调用 执行try方法 returnValue = point.proceed();

    1 } catch (Throwable throwable) { 2 //异常执行cancel 3 4 tccTransactionManager.cancel(); 5 6 throw throwable; 7 } 8 //try成功执行confirm confirm 失败的话,那就只能走本地补偿 9 tccTransactionManager.confirm(); 10 } finally { 11 tccTransactionManager.remove(); 12 } 13 return returnValue;

    }

  • 我们进行跟进 tccTransactionManager.begin() 方法:

    /** * 该方法为发起方第一次调用 * 也是tcc事务的入口 */ void begin() { LogUtil.debug(LOGGER, () -> "开始执行tcc事务!start"); TccTransaction tccTransaction = CURRENT.get(); if (Objects.isNull(tccTransaction)) { tccTransaction = new TccTransaction(); tccTransaction.setStatus(TccActionEnum.TRYING.getCode()); tccTransaction.setRole(TccRoleEnum.START.getCode()); } //保存当前事务信息 coordinatorCommand.execute(new CoordinatorAction(CoordinatorActionEnum.SAVE, tccTransaction));

    1 CURRENT.set(tccTransaction); 2 3 //设置tcc事务上下文,这个类会传递给远端 4 TccTransactionContext context = new TccTransactionContext(); 5 context.setAction(TccActionEnum.TRYING.getCode());//设置执行动作为try 6 context.setTransId(tccTransaction.getTransId());//设置事务id 7 TransactionContextLocal.getInstance().set(context);

    }

  • 这里我们保存了事务信息,并且开启了事务上下文,并把它保存在了ThreadLoacl里面,大家想想这里为什么一定要保存在ThreadLocal里面。

  • begin方法执行完后,我们回到切面,现在我们来执行 point.proceed(),当执行这一句代码的时候,会进入第二个切面,即进入了 TccCoordinatorMethodInterceptor

    public Object interceptor(ProceedingJoinPoint pjp) throws Throwable {

    1 final TccTransaction currentTransaction = tccTransactionManager.getCurrentTransaction(); 2 3 if (Objects.nonNull(currentTransaction)) { 4 final TccActionEnum action = TccActionEnum.acquireByCode(currentTransaction.getStatus()); 5 switch (action) { 6 case TRYING: 7 registerParticipant(pjp, currentTransaction.getTransId()); 8 break; 9 case CONFIRMING: 10 break; 11 case CANCELING: 12 break; 13 } 14 } 15 return pjp.proceed(pjp.getArgs());

    }

  • 这里由于是在try阶段,直接就进入了 registerParticipant 方法

    private void registerParticipant(ProceedingJoinPoint point, String transId) throws NoSuchMethodException {

    1 MethodSignature signature = (MethodSignature) point.getSignature(); 2 Method method = signature.getMethod(); 3 4 Class<?> clazz = point.getTarget().getClass(); 5 6 Object[] args = point.getArgs(); 7 8 final Tcc tcc = method.getAnnotation(Tcc.class); 9 10 //获取协调方法 11 String confirmMethodName = tcc.confirmMethod(); 12 13 /* if (StringUtils.isBlank(confirmMethodName)) { 14 confirmMethodName = method.getName(); 15 }*/ 16 17 String cancelMethodName = tcc.cancelMethod(); 18 19 /* if (StringUtils.isBlank(cancelMethodName)) { 20 cancelMethodName = method.getName(); 21 }

    */ //设置模式 final TccPatternEnum pattern = tcc.pattern();

    1 tccTransactionManager.getCurrentTransaction().setPattern(pattern.getCode()); 2 3 4 TccInvocation confirmInvocation = null; 5 if (StringUtils.isNoneBlank(confirmMethodName)) { 6 confirmInvocation = new TccInvocation(clazz, 7 confirmMethodName, method.getParameterTypes(), args); 8 } 9 10 TccInvocation cancelInvocation = null; 11 if (StringUtils.isNoneBlank(cancelMethodName)) { 12 cancelInvocation = new TccInvocation(clazz, 13 cancelMethodName, 14 method.getParameterTypes(), args); 15 } 16 17 18 //封装调用点 19 final Participant participant = new Participant( 20 transId, 21 confirmInvocation, 22 cancelInvocation); 23 24 tccTransactionManager.enlistParticipant(participant); 25 26}
  • 这里就获取了 @Tcc(confirmMethod = "confirmOrderStatus", cancelMethod = "cancelOrderStatus") 的信息,并把他封装成了 Participant,存起来,然后真正的调用 PaymentServiceImpl.makePayment 业务方法,在业务方法里面,我们首先执行的是更新订单状态(相信现在你还记得0.0):

    1 order.setStatus(OrderStatusEnum.PAYING.getCode()); 2 orderMapper.update(order);
  • 接下来,我们进行扣除用户余额,注意扣除余额这里是一个rpc方法:

    1 AccountDTO accountDTO = new AccountDTO(); 2 accountDTO.setAmount(order.getTotalAmount()); 3 accountDTO.setUserId(order.getUserId()); 4 accountService.payment(accountDTO)
  • 现在我们来关注 accountService.payment(accountDTO) ,这个接口的定义。

dubbo接口:

1public interface AccountService { 2 3 4 /** 5 * 扣款支付 6 * 7 * @param accountDTO 参数dto 8 * @return true 9 */ 10 @Tcc 11 boolean payment(AccountDTO accountDTO); 12} 13

springcloud接口

1@FeignClient(value = "account-service", configuration = MyConfiguration.class) 2public interface AccountClient { 3 4 @PostMapping("/account-service/account/payment") 5 @Tcc 6 Boolean payment(@RequestBody AccountDTO accountDO); 7 8}

很明显这里我们都在接口上加了@Tcc注解,我们知道springAop的特性,在接口上加注解,是无法进入切面的,所以我们在这里,要采用rpc框架的某些特性来帮助我们获取到 @Tcc注解信息。 这一步很重要。当我们发起

accountService.payment(accountDTO) 调用的时候:

  • dubbo用户,会走dubbo的filter接口,TccTransactionFilter:

    private void registerParticipant(Class clazz, String methodName, Object[] arguments, Class... args) throws TccRuntimeException { try { Method method = clazz.getDeclaredMethod(methodName, args); Tcc tcc = method.getAnnotation(Tcc.class); if (Objects.nonNull(tcc)) {

    1 //获取事务的上下文 2 final TccTransactionContext tccTransactionContext = 3 TransactionContextLocal.getInstance().get(); 4 if (Objects.nonNull(tccTransactionContext)) { 5 //dubbo rpc传参数 6 RpcContext.getContext() 7 .setAttachment(Constant.TCC_TRANSACTION_CONTEXT, 8 GsonUtils.getInstance().toJson(tccTransactionContext)); 9 } 10 if (Objects.nonNull(tccTransactionContext)) { 11 if(TccActionEnum.TRYING.getCode()==tccTransactionContext.getAction()){ 12 //获取协调方法 13 String confirmMethodName = tcc.confirmMethod(); 14 15 if (StringUtils.isBlank(confirmMethodName)) { 16 confirmMethodName = method.getName(); 17 } 18 19 String cancelMethodName = tcc.cancelMethod(); 20 21 if (StringUtils.isBlank(cancelMethodName)) { 22 cancelMethodName = method.getName(); 23 } 24 //设置模式 25 final TccPatternEnum pattern = tcc.pattern(); 26 27 tccTransactionManager.getCurrentTransaction().setPattern(pattern.getCode()); 28 TccInvocation confirmInvocation = new TccInvocation(clazz, 29 confirmMethodName, 30 args, arguments); 31 32 TccInvocation cancelInvocation = new TccInvocation(clazz, 33 cancelMethodName, 34 args, arguments); 35 36 //封装调用点 37 final Participant participant = new Participant( 38 tccTransactionContext.getTransId(), 39 confirmInvocation, 40 cancelInvocation); 41 42 tccTransactionManager.enlistParticipant(participant); 43 } 44 45 } 46 47 48 49 } 50 } catch (NoSuchMethodException e) { 51 throw new TccRuntimeException("not fount method " + e.getMessage()); 52 }

    }

  • springcloud 用户,则会进入 TccFeignHandler

    public class TccFeignHandler implements InvocationHandler { /** * logger */ private static final Logger LOGGER = LoggerFactory.getLogger(TccFeignHandler.class);

    1private Target<?> target; 2private Map<Method, MethodHandler> handlers; 3 4public Object invoke(Object proxy, Method method, Object[] args) throws Throwable { 5 if (Object.class.equals(method.getDeclaringClass())) { 6 return method.invoke(this, args); 7 } else { 8 9 final Tcc tcc = method.getAnnotation(Tcc.class); 10 if (Objects.isNull(tcc)) { 11 return this.handlers.get(method).invoke(args); 12 } 13 14 final TccTransactionContext tccTransactionContext = 15 TransactionContextLocal.getInstance().get(); 16 if (Objects.nonNull(tccTransactionContext)) { 17 final TccTransactionManager tccTransactionManager = 18 SpringBeanUtils.getInstance().getBean(TccTransactionManager.class); 19 if (TccActionEnum.TRYING.getCode() == tccTransactionContext.getAction()) { 20 //获取协调方法 21 String confirmMethodName = tcc.confirmMethod(); 22 23 if (StringUtils.isBlank(confirmMethodName)) { 24 confirmMethodName = method.getName(); 25 } 26 27 String cancelMethodName = tcc.cancelMethod(); 28 29 if (StringUtils.isBlank(cancelMethodName)) { 30 cancelMethodName = method.getName(); 31 } 32 33 //设置模式 34 final TccPatternEnum pattern = tcc.pattern(); 35 36 tccTransactionManager.getCurrentTransaction().setPattern(pattern.getCode()); 37 38 final Class<?> declaringClass = method.getDeclaringClass(); 39 40 TccInvocation confirmInvocation = new TccInvocation(declaringClass, 41 confirmMethodName, 42 method.getParameterTypes(), args); 43 44 TccInvocation cancelInvocation = new TccInvocation(declaringClass, 45 cancelMethodName, 46 method.getParameterTypes(), args); 47 48 //封装调用点 49 final Participant participant = new Participant( 50 tccTransactionContext.getTransId(), 51 confirmInvocation, 52 cancelInvocation); 53 54 tccTransactionManager.enlistParticipant(participant); 55 } 56 57 } 58 59 60 return this.handlers.get(method).invoke(args); 61 } 62} 63 64 65public void setTarget(Target<?> target) { 66 this.target = target; 67} 68 69 70public void setHandlers(Map<Method, MethodHandler> handlers) { 71 this.handlers = handlers; 72}

    }

重要提示 在这里我们获取了在第一步设置的事务上下文,然后进行dubbo的rpc传参数,细心的你可能会发现,这里获取远端confirm,cancel方法其实就是自己本身啊?,那发起者还怎么来调用远端的confrim,cancel方法呢?这里的调用,我们通过事务上下文的状态来控制,比如在confrim阶段,我们就调用confrim方法等。。

到这里我们就完成了整个消费端的调用,下一篇分析提供者的调用流程!

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

2020年前端实用代码段,为你的工作保驾护航

有空的时候,自己总结了几个代码段,在开发中也经常使用,谢谢。1、使用解构获取json数据let jsonData  id: 1,status: "OK",data: 'a', 'b';let  id, status, data: number   jsonData;console.log(id, status, number )