支付系统-支付微服务落地设计3:支付宝微信支付具体实现、交易查询 (异步回调+定时轮询)
大家好,我是此林。
前面我们讲述了支付系统的支付渠道设计,包括支付渠道的概要、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,那么就要使用内网穿透技术,把我们的服务映射到公网。像 ngrok、cpolar 等都是比较常见的内网穿透工具。
回调 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. 更进一步校验:
- 商家需要验证该通知数据中的 out_trade_no 是否为商家系统中创建的订单号。
- 判断 total_amount 是否确实为该订单的实际金额(即商家订单创建时的金额)。
- 校验通知中的 seller_id(或者 seller_email ) 是否为 out_trade_no 这笔单据的对应的操作方(有的时候,一个商家可能有多个seller_id/seller_email)。
- 验证 app_id 是否为该商家本身。
-
上述 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 的节点 来处理
这样每个订单只会被其中一个节点处理到,保证不会重复。
今天的介绍就到这里了。
我是此林,关注我吧!
带你看不一样的世界!
更多推荐




所有评论(0)