大家好,我是此林。

支付系统-支付微服务落地设计1:支付渠道开发

支付系统-支付微服务落地设计2:扫码支付流程

前面我们讲述了支付系统的支付渠道设计,包括支付渠道的概要、redis 缓存、工厂+策略模式解耦等、扫码支付流程的内容。

今天我们来继续讲述支付宝支付和微信支付具体实现、支付交易查询、退款具体流程。

1. 支付宝扫码支付实现

1.1. AlipayConfig 支付配置

/**
 * 支付宝支付的配置
 */
public class AlipayConfig {

    /**
     * 将支付渠道配置转化为支付宝的配置
     *
     * @param enterpriseId 商户ID
     * @return 支付宝的配置
     */
    public static Config getConfig(Long enterpriseId) {
        // 查询配置
        PayChannelService payChannelService = SpringUtil.getBean(PayChannelService.class);
        PayChannelEntity payChannel = payChannelService.findByEnterpriseId(enterpriseId, TradingConstant.TRADING_CHANNEL_ALI_PAY);

        if (ObjectUtil.isEmpty(payChannel)) {
            throw new WYException(TradingEnum.CONFIG_EMPTY);
        }

        Config config = new Config();
        config.protocol = "https";
        config.gatewayHost = payChannel.getDomain();
        config.signType = "RSA2";
        config.appId = payChannel.getAppId();
        //配置应用私钥
        config.merchantPrivateKey = payChannel.getMerchantPrivateKey();
        //配置支付宝公钥
        config.alipayPublicKey = payChannel.getPublicKey();
        //可设置异步通知接收服务地址(可选)
        config.notifyUrl = StrUtil.replace(payChannel.getNotifyUrl(), "{enterpriseId}", Convert.toStr(enterpriseId));
        //设置AES密钥,调用AES加解密相关接口时需要(可选)
        config.encryptKey = payChannel.getEncryptKey();
        return config;
    }

}

这里还是一样,通过工厂模式拿到数据库中对应的支付渠道(appId、支付公钥、支付私钥等)信息,然后封装为 config 返回。

注意这里的签名方式是 RSA2,我们在支付平台回调的时候会验签,保证请求是支付平台来的。因为我们的回调地址会暴露在网关中,公网可以访问,所以要做验签机制。

AES 加密解密是可选项,主要是为了防止明文传输造成支付信息泄露,保障支付安全。

1.2. 具体实现

/**
 * 支付宝的扫描支付的具体实现
 */
@Slf4j
@Component("aliNativePayHandler")
@PayChannel(type = PayChannelEnum.ALI_PAY)
public class AliNativePayHandler implements NativePayHandler {

    @Override
    public void createDownLineTrading(TradingEntity tradingEntity) throws SLException {
        //查询配置
        Config config = AlipayConfig.getConfig(tradingEntity.getEnterpriseId());
        //Factory使用配置
        Factory.setOptions(config);
        AlipayTradePrecreateResponse response;
        try {
            //调用支付宝API面对面支付
            response = Factory
                    .Payment
                    .FaceToFace()
                    .preCreate(tradingEntity.getMemo(), //订单描述
                            Convert.toStr(tradingEntity.getTradingOrderNo()), //业务订单号
                            Convert.toStr(tradingEntity.getTradingAmount())); //金额
        } catch (Exception e) {
            log.error("支付宝统一下单创建失败:tradingEntity = {}", tradingEntity, e);
            throw new WYException(TradingEnum.NATIVE_PAY_FAIL, e);
        }

        //受理结果【只表示请求是否成功,而不是支付是否成功】
        boolean isSuccess = ResponseChecker.success(response);
        //6.1、受理成功:修改交易单
        if (isSuccess) {
            String subCode = response.getSubCode();
            String subMsg = response.getQrCode();
            tradingEntity.setPlaceOrderCode(subCode); //返回的编码
            tradingEntity.setPlaceOrderMsg(subMsg); //二维码需要展现的信息
            tradingEntity.setPlaceOrderJson(JSONUtil.toJsonStr(response));
            tradingEntity.setTradingState(TradingStateEnum.FKZ);
            return;
        }
        throw new WYException(JSONUtil.toJsonStr(response), TradingEnum.NATIVE_PAY_FAIL.getCode(), TradingEnum.NATIVE_PAY_FAIL.getStatus());
    }

}

这里我们用了支付宝的 Easy-SDK,只需通过 Factory 的 FaceToFace() 即可完成支付。

Factory 里先设置刚刚获取到的 AlipayConfig,然后调用 pay() ,传入参数订单描述、交易号、金额,即可完成统一下单。

调用成功后,更新交易流水的状态为付款中,以及二维码相关链接等信息。

那其实微信的实现也差不多,官方也提供了相应的 SDK,我们只需要根据官方文档配置调用即可,这里就不多赘述了。

2. 支付查询

2.1. 异步回调

在支付平台创建交易单后,如果用户支付成功,我们怎么知道支付成功了呢?一般的做法有两种,分别是【异步通知】和【主动查询】。

异步通知是指一笔订单支付完成后,支付平台会将该笔订单的变更信息,沿着商家调用支付请求时所传入的异步通知地址 notify_url,通过 POST 请求的形式将支付结果作为参数通知到商家系统。

关于异步通知,支付宝和微信都提供了异步通知功能,具体参考官方文档:

支付宝异步回调

微信异步回调

无论哪种异步回调,支付平台都需要我们提供给它公网可以访问的接口,这个接口我们通常会暴露在支付网关里。

如果我们做练手项目,没有公网IP,那么就要使用内网穿透技术,把我们的服务映射到公网。像 ngrokcpolar 等都是比较常见的内网穿透工具。

回调 NotifyController 设计:

/**
 * 支付结果的通知
 */
@RestController
@Api(tags = "支付通知")
@RequestMapping("notify")
public class NotifyController {

    @Resource
    private NotifyService notifyService;

    /**
     * 微信支付成功回调(成功后无需响应内容)
     *
     * @param httpEntity   微信请求信息
     * @param enterpriseId 商户id
     * @return 正常响应200,否则响应500
     */
    @PostMapping("wx/{enterpriseId}")
    public ResponseEntity<Object> wxPayNotify(HttpEntity<String> httpEntity, @PathVariable("enterpriseId") Long enterpriseId) {
        try {
            //获取请求头
            HttpHeaders headers = httpEntity.getHeaders();

            //构建微信请求数据对象
            NotificationRequest request = new NotificationRequest.Builder()
                    .withSerialNumber(headers.getFirst("Wechatpay-Serial")) //证书序列号(微信平台)
                    .withNonce(headers.getFirst("Wechatpay-Nonce"))  //随机串
                    .withTimestamp(headers.getFirst("Wechatpay-Timestamp")) //时间戳
                    .withSignature(headers.getFirst("Wechatpay-Signature")) //签名字符串
                    .withBody(httpEntity.getBody())
                    .build();

            //微信通知的业务处理
            this.notifyService.wxPayNotify(request, enterpriseId);

        } catch (WYException e) {
            Map<String, Object> result = MapUtil.<String, Object>builder()
                    .put("code", "FAIL")
                    .put("message", e.getMsg())
                    .build();
            //响应500
            return ResponseEntity.status(HttpStatus.INTERNAL_SERVER_ERROR).body(result);
        }
        return ResponseEntity.ok(null);
    }

    /**
     * 支付宝支付成功回调(成功后需要响应success)
     *
     * @param enterpriseId 商户id
     * @return 正常响应200,否则响应500
     */
    @PostMapping("alipay/{enterpriseId}")
    public ResponseEntity<String> aliPayNotify(HttpServletRequest request,
                                               @PathVariable("enterpriseId") Long enterpriseId) {
        try {
            //支付宝通知的业务处理
            this.notifyService.aliPayNotify(request, enterpriseId);
        } catch (WYException e) {
            //响应500
            return ResponseEntity.status(HttpStatus.INTERNAL_SERVER_ERROR).build();
        }
        return ResponseEntity.ok("success");
    }
}

NotifyServiceImpl 设计:

/**
 * 支付成功的通知处理
 */
@Slf4j
@Service
public class NotifyServiceImpl implements NotifyService {

    @Resource
    private TradingService tradingService;
    @Resource
    private RedissonClient redissonClient;
    @Resource
    private MQFeign mqFeign;

    @Override
    public void wxPayNotify(NotificationRequest request, Long enterpriseId) throws SLException {
        // 查询配置
        WechatPayHttpClient client = WechatPayHttpClient.get(enterpriseId);

        JSONObject jsonData;

        //验证签名,确保请求来自微信
        try {
            //确保在管理器中存在自动更新的商户证书
            client.createHttpClient();

            CertificatesManager certificatesManager = CertificatesManager.getInstance();
            Verifier verifier = certificatesManager.getVerifier(client.getMchId());

            //验签和解析请求数据
            NotificationHandler notificationHandler = new NotificationHandler(verifier, client.getApiV3Key().getBytes(StandardCharsets.UTF_8));
            Notification notification = notificationHandler.parse(request);

            if (!StrUtil.equals("TRANSACTION.SUCCESS", notification.getEventType())) {
                //非成功请求直接返回,理论上都是成功的请求
                return;
            }

            //获取解密后的数据
            jsonData = JSONUtil.parseObj(notification.getDecryptData());
        } catch (Exception e) {
            throw new WYException("验签失败");
        }

        if (!StrUtil.equals(jsonData.getStr("trade_state"), TradingConstant.WECHAT_TRADE_SUCCESS)) {
            return;
        }

        //交易单号
        Long tradingOrderNo = jsonData.getLong("out_trade_no");
        log.info("微信支付通知:tradingOrderNo = {}, data = {}", tradingOrderNo, jsonData);

        //更新交易单
        this.updateTrading(tradingOrderNo, jsonData.getStr("trade_state_desc"), jsonData.toString());
    }

    private void updateTrading(Long tradingOrderNo, String resultMsg, String resultJson) {
        String key = TradingCacheConstant.CREATE_PAY + tradingOrderNo;
        RLock lock = redissonClient.getFairLock(key);
        try {
            //获取锁
            if (lock.tryLock(TradingCacheConstant.REDIS_WAIT_TIME, TimeUnit.SECONDS)) {
                TradingEntity trading = this.tradingService.findTradByTradingOrderNo(tradingOrderNo);
                if (trading.getTradingState() == TradingStateEnum.YJS) {
                    // 已付款
                    return;
                }

                //设置成付款成功
                trading.setTradingState(TradingStateEnum.YJS);
                //清空二维码数据
                trading.setQrCode("");
                trading.setResultMsg(resultMsg);
                trading.setResultJson(resultJson);
                this.tradingService.saveOrUpdate(trading);

                // 发消息通知其他系统支付成功
                TradeStatusMsg tradeStatusMsg = TradeStatusMsg.builder()
                        .tradingOrderNo(trading.getTradingOrderNo())
                        .productOrderNo(trading.getProductOrderNo())
                        .statusCode(TradingStateEnum.YJS.getCode())
                        .statusName(TradingStateEnum.YJS.name())
                        .build();

                String msg = JSONUtil.toJsonStr(Collections.singletonList(tradeStatusMsg));
                this.mqFeign.sendMsg(Constants.MQ.Exchanges.TRADE, Constants.MQ.RoutingKeys.TRADE_UPDATE_STATUS, msg);
                return;
            }
        } catch (Exception e) {
            throw new WYException("处理业务失败");
        } finally {
            lock.unlock();
        }
        throw new WYException("处理业务失败");
    }

    @Override
    public void aliPayNotify(HttpServletRequest request, Long enterpriseId) throws SLException {
        //获取参数
        Map<String, String[]> parameterMap = request.getParameterMap();
        Map<String, String> param = new HashMap<>();
        for (Map.Entry<String, String[]> entry : parameterMap.entrySet()) {
            param.put(entry.getKey(), StrUtil.join(",", entry.getValue()));
        }

        String tradeStatus = param.get("trade_status");
        if (!StrUtil.equals(tradeStatus, TradingConstant.ALI_TRADE_SUCCESS)) {
            return;
        }

        //查询配置
        Config config = AlipayConfig.getConfig(enterpriseId);
        Factory.setOptions(config);
        try {
            Boolean result = Factory
                    .Payment
                    .Common().verifyNotify(param);
            if (!result) {
                throw new WYException("验签失败");
            }
        } catch (Exception e) {
            throw new WYException("验签失败");
        }

        //获取交易单号
        Long tradingOrderNo = Convert.toLong(param.get("out_trade_no"));
        //更新交易单
        this.updateTrading(tradingOrderNo, "支付成功", JSONUtil.toJsonStr(param));
    }
}

代码看着有点多,其实流程不难。

第一步:验签。使用 RSA 的验签方法,通过签名字符串、签名参数(经过 base64 解码)及支付平台公钥验证签名。

为什么要验签?我们要确保这个请求是支付平台发过来的,因为我们这个回调接口是暴露在公网的,如果不验签被恶意利用,那势必造成支付安全问题。

第二步:幂等性校验。

1. 如果发现支付平台告知支付状态不成功,直接 return 不修改支付状态。

2. 如果交易状态为已结算,那么无需再次修改订单状态,直接 return。

3. 更进一步校验:

  1. 商家需要验证该通知数据中的 out_trade_no 是否为商家系统中创建的订单号。
  2. 判断 total_amount 是否确实为该订单的实际金额(即商家订单创建时的金额)。
  3. 校验通知中的 seller_id(或者 seller_email ) 是否为 out_trade_no 这笔单据的对应的操作方(有的时候,一个商家可能有多个seller_id/seller_email)。
  4. 验证 app_id 是否为该商家本身。
  5. 上述 1、2、3、4 有任何一个验证不通过,则表明本次通知是异常通知,务必忽略。

问:为什么支付平台告知支付状态不成功,直接 return 不修改支付状态?

答:根据支付宝官方文档,默认情况下只有在支付成功的条件下才会触发异步回调,所以对于告知我们支付不成功的,视为异常回调,不做处理。

第三步:更新交易单状态。

再贴一下更新状态的代码。

    private void updateTrading(Long tradingOrderNo, String resultMsg, String resultJson) {
        String key = TradingCacheConstant.CREATE_PAY + tradingOrderNo;
        RLock lock = redissonClient.getFairLock(key);
        try {
            //获取锁
            if (lock.tryLock(TradingCacheConstant.REDIS_WAIT_TIME, TimeUnit.SECONDS)) {
                TradingEntity trading = this.tradingService.findTradByTradingOrderNo(tradingOrderNo);
                if (trading.getTradingState() == TradingStateEnum.YJS) {
                    // 已付款
                    return;
                }

                //设置成付款成功
                trading.setTradingState(TradingStateEnum.YJS);
                //清空二维码数据
                trading.setQrCode("");
                trading.setResultMsg(resultMsg);
                trading.setResultJson(resultJson);
                this.tradingService.saveOrUpdate(trading);

                // 发消息通知其他系统支付成功
                TradeStatusMsg tradeStatusMsg = TradeStatusMsg.builder()
                        .tradingOrderNo(trading.getTradingOrderNo())
                        .productOrderNo(trading.getProductOrderNo())
                        .statusCode(TradingStateEnum.YJS.getCode())
                        .statusName(TradingStateEnum.YJS.name())
                        .build();

                String msg = JSONUtil.toJsonStr(Collections.singletonList(tradeStatusMsg));
                this.mqFeign.sendMsg(Constants.MQ.Exchanges.TRADE, Constants.MQ.RoutingKeys.TRADE_UPDATE_STATUS, msg);
                return;
            }
        } catch (Exception e) {
            throw new WYException("处理业务失败");
        } finally {
            lock.unlock();
        }
        throw new WYException("处理业务失败");
    }

我们这里同样地对交易单号加了分布式锁,为什么要加分布式锁,只要涉及到交易单的状态更改我们基本上都加了分布式锁,其实还是幂等性。

因为我们的幂等性逻辑是如果支付平台返回 “交易成功” 且 数据库查询出支付状态不是 “已结算”,那么就要更新支付状态,然后清除二维码数据等操作(我们代码里读操作和写操作是分开的,不具有原子性)。

实际生产中,由于网络或者一些不可知的原因,我们不能保证回调请求只会发一次,或者没有异常情况,所以我们系统这边要尽可能做好容错,避免出现并发写问题。

更进一步,我们前端每10s轮询后端更新支付状态,我们后端也有XXL-JOB定时任务,这些都会造成并发修改支付状态。

何况我们后续还要发送 MQ 通知订单系统去更新订单状态,这里如果不加分布式锁,多次回调请求就会导致发送多次 MQ,造成重复生产和消费问题。

不过个人觉得这里改为数据库级别的乐观锁或悲观锁也可以。

2.2. XXL-JOB 定时任务

如果因为某些原因导致支付 异步回调通知没有接收到,我们会怎么处理?

首先无论有没有异步通知,前端都会每隔 10s 查询后端接口。

    @PostMapping("query/{tradingOrderNo}")
    @ApiOperation(value = "查询统一收单线下交易", notes = "查询统一收单线下交易")
    @ApiImplicitParam(name = "tradingOrderNo", value = "交易单", required = true)
    public TradingDTO queryTrading(@PathVariable("tradingOrderNo") Long tradingOrderNo) {
        return this.basicPayService.queryTrading(tradingOrderNo);
    }

这个 queryTrading() 里,会调用支付平台 API 查询并更新支付状态。

但是即使是前端有轮询机制,我们后端依然要做好兜底曾略,

所以我们这边还用 XXL-JOB 分片广播定时轮询支付平台、以更新支付状态来做了兜底。

这个定时任务的时间间隔就要长一点,比如每10分钟执行一次。

问:为什么要使用xxl-job呢?

一般在项目中实现定时任务主要是两种技术方案,一种是Spring Task,另一种是xxl-job,其中Spring Task是适合单体项目中使用,而xxl-job是分布式任务调度框架,更适合在分布式项目中使用,所以在支付微服务中我们将采用xxl-job来实现。

问:那什么是 XXL-JOB 分片广播?

分片 (Sharding):当某个任务很大时,可以将任务拆分成多个小片段(shard),由调度中心分配给集群中的执行器节点去并行处理。

 • 广播 (Broadcast):调度中心会将同一个任务下发到所有在线的执行器节点上,让每个节点都执行一次。

上图为 xxl-job 架构。

关于xxl-job 的具体使用,这里不多介绍,它的底层原理在下面这篇文章已经详细叙述。

https://blog.csdn.net/2401_82540083/article/details/145392738?spm=1001.2014.3001.5502

我们重点来看它的调度器的路由策略,即调度器每次把任务委派给各个执行器的策略:

● FIRST(第一个):固定选择第一个机器;

● LAST(最后一个):固定选择最后一个机器;

● ROUND(轮询):在线的机器按照顺序一次执行一个

● RANDOM(随机):随机选择在线的机器;

● CONSISTENT_HASH(一致性HASH):每个任务按照Hash算法固定选择某一台机器,且所有任务均匀散列在不同机器上。

● LEAST_FREQUENTLY_USED(最不经常使用):使用频率最低的机器优先被选举;

● LEAST_RECENTLY_USED(最近最久未使用):最久未使用的机器优先被选举;

● FAILOVER(故障转移):按照顺序依次进行心跳检测,第一个心跳检测成功的机器选定为目标执行器并发起调度;

● BUSYOVER(忙碌转移):按照顺序依次进行空闲检测,第一个空闲检测成功的机器选定为目标执行器并发起调度;

● SHARDING_BROADCAST(分片广播):广播触发对应集群中所有机器执行一次任务,同时系统自动传递分片参数;可根据分片参数开发分片任务;

我们这里使用分片广播模式。

/**
 * 交易任务,主要是查询订单的支付状态 和 退款的成功状态
 */
@Slf4j
@Component
public class TradeJob {

    @Value("${wy.job.trading.count:100}")
    private Integer tradingCount;
    @Value("${wy.job.refund.count:100}")
    private Integer refundCount;
    @Resource
    private TradingService tradingService;
    @Resource
    private RefundRecordService refundRecordService;
    @Resource
    private BasicPayService basicPayService;
    @Resource
    private MQFeign mqFeign;

    /**
     * 分片广播方式查询支付状态
     * 逻辑:每次最多查询{tradingCount}个未完成的交易单,交易单id与shardTotal取模,值等于shardIndex进行处理
     */
    @XxlJob("tradingJob")
    public void tradingJob() {
        // 分片参数
        int shardIndex = NumberUtil.max(XxlJobHelper.getShardIndex(), 0);
        int shardTotal = NumberUtil.max(XxlJobHelper.getShardTotal(), 1);

        List<TradingEntity> list = this.tradingService.findListByTradingState(TradingStateEnum.FKZ, tradingCount);
        if (CollUtil.isEmpty(list)) {
            XxlJobHelper.log("查询到交易单列表为空!shardIndex = {}, shardTotal = {}", shardIndex, shardTotal);
            return;
        }

        //定义消息通知列表,只要是状态不为【付款中】就需要通知其他系统
        List<TradeStatusMsg> tradeMsgList = new ArrayList<>();
        for (TradingEntity trading : list) {
            if (trading.getTradingOrderNo() % shardTotal != shardIndex) {
                continue;
            }
            try {
                //查询交易单
                TradingDTO tradingDTO = this.basicPayService.queryTrading(trading.getTradingOrderNo());
                if (TradingStateEnum.FKZ != tradingDTO.getTradingState()) {
                    TradeStatusMsg tradeStatusMsg = TradeStatusMsg.builder()
                            .tradingOrderNo(trading.getTradingOrderNo())
                            .productOrderNo(trading.getProductOrderNo())
                            .statusCode(tradingDTO.getTradingState().getCode())
                            .statusName(tradingDTO.getTradingState().name())
                            .build();
                    tradeMsgList.add(tradeStatusMsg);
                }
            } catch (Exception e) {
                XxlJobHelper.log("查询交易单出错!shardIndex = {}, shardTotal = {}, trading = {}", shardIndex, shardTotal, trading, e);
            }
        }

        if (CollUtil.isEmpty(tradeMsgList)) {
            return;
        }

        //发送消息通知其他系统
        String msg = JSONUtil.toJsonStr(tradeMsgList);
        this.mqFeign.sendMsg(Constants.MQ.Exchanges.TRADE, Constants.MQ.RoutingKeys.TRADE_UPDATE_STATUS, msg);
    }

    /**
     * 分片广播方式查询退款状态
     */
    @XxlJob("refundJob")
    public void refundJob() {
        // 分片参数
        int shardIndex = NumberUtil.max(XxlJobHelper.getShardIndex(), 0);
        int shardTotal = NumberUtil.max(XxlJobHelper.getShardTotal(), 1);

        List<RefundRecordEntity> list = this.refundRecordService.findListByRefundStatus(RefundStatusEnum.SENDING, refundCount);
        if (CollUtil.isEmpty(list)) {
            XxlJobHelper.log("查询到退款单列表为空!shardIndex = {}, shardTotal = {}", shardIndex, shardTotal);
            return;
        }

        //定义消息通知列表,只要是状态不为【退款中】就需要通知其他系统
        List<TradeStatusMsg> tradeMsgList = new ArrayList<>();

        for (RefundRecordEntity refundRecord : list) {
            if (refundRecord.getRefundNo() % shardTotal != shardIndex) {
                continue;
            }
            try {
                //查询退款单
                RefundRecordDTO refundRecordDTO = this.basicPayService.queryRefundTrading(refundRecord.getRefundNo());
                if (RefundStatusEnum.SENDING != refundRecordDTO.getRefundStatus()) {
                    TradeStatusMsg tradeStatusMsg = TradeStatusMsg.builder()
                            .tradingOrderNo(refundRecord.getTradingOrderNo())
                            .productOrderNo(refundRecord.getProductOrderNo())
                            .refundNo(refundRecord.getRefundNo())
                            .statusCode(refundRecord.getRefundStatus().getCode())
                            .statusName(refundRecord.getRefundStatus().name())
                            .build();
                    tradeMsgList.add(tradeStatusMsg);
                }
            } catch (Exception e) {
                XxlJobHelper.log("查询退款单出错!shardIndex = {}, shardTotal = {}, refundRecord = {}", shardIndex, shardTotal, refundRecord, e);
            }
        }

        if (CollUtil.isEmpty(tradeMsgList)) {
            return;
        }

        //发送消息通知其他系统
        String msg = JSONUtil.toJsonStr(tradeMsgList);
        this.mqFeign.sendMsg(Constants.MQ.Exchanges.TRADE, Constants.MQ.RoutingKeys.REFUND_UPDATE_STATUS, msg);
    }
}

在分布式定时任务里,如果同一个 Job 要在多个执行器节点上跑,就涉及 任务分片

  • shardIndex:当前执行器节点的分片编号(0、1、2…)

  • shardTotal:总的分片数(即有多少个执行器节点一起跑这个任务)

这样,每个节点就可以只处理一部分数据,避免重复计算。

我们代码里:

  • shardIndex 意思是当前机器负责第几片

  • shardTotal 意思是一共有多少片

接下来查询所有 付款中(FKZ) 的交易单:

List<TradingEntity> list = this.tradingService.findListByTradingState(TradingStateEnum.FKZ, tradingCount);

然后通过 取模分片 来决定某个订单由哪个节点处理:

if (trading.getTradingOrderNo() % shardTotal != shardIndex) {
    continue;
}

比如:

  • 有 3 台机器(shardTotal = 3),分别是 0、1、2

  • 一个交易单号是 1005

  • 计算 1005 % 3 = 2,所以应该由 分片 2 的节点 来处理

这样每个订单只会被其中一个节点处理到,保证不会重复。

今天的介绍就到这里了。

我是此林,关注我吧!

带你看不一样的世界!

Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐