支付系统已从简单的支付通道对接演变为集交易、账务、风控、清算、对账于一体的关键业务中台。其核心挑战在于:如何保障交易最终一致性、账务准确性、高并发下的稳定性和实时风控的有效性。本文基于实际生产级实现,从交易模型、状态机、账务TCC、分布式事务、动态路由、熔断降级、幂等防重、对账引擎及可观测性等维度,系统阐述支付核心引擎的设计思路与代码落地细节。所有示例基于Spring Cloud Alibaba + Seata + RocketMQ,并已在大规模交易场景下验证。
我们将支付引擎划分为三层,确保职责清晰、变更隔离:
各层通过标准化API交互,避免业务逻辑对底层通道产生直接依赖。
支付订单是贯穿全流程的核心数据载体,定义如下:
@Data
@TableName("payment_order")
public class PaymentOrder {
private Long id;
private String orderNo; // 业务订单号
private String channelOrderNo; // 渠道流水号
private BigDecimal amount;
private Integer currency;
private String payerId;
private String payeeId;
private Integer status; // 0:待支付 1:支付中 2:成功 3:失败 4:关闭 5:待退款
private Integer bizType;
private LocalDateTime createTime;
private LocalDateTime updateTime;
private Integer version; // 乐观锁
}采用有限状态机(FSM)驱动订单生命周期,显式定义合法状态迁移路径:
状态变更必须满足原子性和幂等性。我们使用乐观锁version字段,每次更新时校验当前状态和版本号,若影响行数为0则说明存在并发冲突,需抛出异常并触发重试机制(如使用Spring Retry)。
public boolean updateStatus(Long id, Integer fromStatus, Integer toStatus, Integer version) {
return paymentOrderMapper.update(
new LambdaUpdateWrapper<PaymentOrder>()
.eq(PaymentOrder::getId, id)
.eq(PaymentOrder::getStatus, fromStatus)
.eq(PaymentOrder::getVersion, version)
.set(PaymentOrder::getStatus, toStatus)
.set(PaymentOrder::getVersion, version + 1)
) > 0;
}设计考量:为避免ABA问题,建议配合updateTime时间戳或全局递增序列,同时将状态机转换逻辑封装为独立服务,确保所有调用入口统一。
支付涉及付款方资产减少、收款方资产增加、平台手续费收入等科目变动。账务系统与交易主链路解耦,采用复式记账原则,保证借贷平衡。
为了防止账户资金超扣并避免长事务锁定数据库,我们引入TCC(Try-Confirm-Cancel)模式。以扣款为例,Try阶段冻结资金,Confirm阶段转为实际扣减,Cancel阶段释放冻结。
@TwoPhaseBusinessAction(name = "debitAccount", commitMethod = "confirm", rollbackMethod = "cancel")
public void tryDebit(AccountTransactionContext ctx) {
// 冻结资金:增加冻结字段,减少可用余额
accountMapper.freeze(ctx.getAccountId(), ctx.getAmount());
// 记录冻结流水(状态为TRY),用于对账和追溯
}
public void confirm(AccountTransactionContext ctx) {
// 将冻结转为实际扣减
accountMapper.confirmDebit(ctx.getAccountId(), ctx.getAmount());
// 更新流水状态为CONFIRM
}
public void cancel(AccountTransactionContext ctx) {
// 解冻资金
accountMapper.unfreeze(ctx.getAccountId(), ctx.getAmount());
// 流水状态CANCEL
}关键问题与对策:
transactionId防重,避免重复冻结。我们基于Seata TCC框架,并结合自定义@TccTransactional注解,通过AOP管理上下文传递和事务生命周期。
支付主链路无法将账务、积分、通知等下游操作置于同一本地事务,因此采用本地事务表 + 消息队列的最终一致性模式,避免分布式事务引入的性能开销。
OutboxMessage(待发送消息),同库事务保证原子性。主事务内插入订单和消息表:
@Transactional(rollbackFor = Exception.class)
public void processPayment(PaymentRequest request) {
PaymentOrder order = buildOrder(request);
paymentOrderMapper.insert(order);
OutboxMessage outbox = OutboxMessage.builder()
.aggregateType("PaymentOrder")
.aggregateId(order.getId())
.eventType("PAYMENT_SUCCESS")
.payload(JSON.toJSONString(order))
.status(0) // 待发送
.build();
outboxMessageMapper.insert(outbox);
// 渠道调用可能异步,若同步返回成功则直接更新订单状态并触发立即发送
}定时发送任务(需分布式锁防止重复执行):
@Scheduled(fixedDelay = 5000)
public void sendOutboxMessages() {
List<OutboxMessage> pending = outboxMessageMapper.selectList(
new LambdaQueryWrapper<OutboxMessage>().eq(OutboxMessage::getStatus, 0).last("limit 100")
);
for (OutboxMessage msg : pending) {
try {
rocketMQProducer.send(msg.getTopic(), msg.getPayload());
msg.setStatus(1); // 已发送
outboxMessageMapper.updateById(msg);
} catch (Exception e) {
msg.setRetryCount(msg.getRetryCount() + 1);
if (msg.getRetryCount() >= 5) {
msg.setStatus(2); // 死信
}
outboxMessageMapper.updateById(msg);
}
}
}消费端幂等性:下游消费者需依据messageId进行幂等处理,使用Redis记录已处理ID,防止重复消费。
不同支付渠道(微信、支付宝、银联等)的API差异大,稳定性不一。我们实现动态路由层,支持权重轮询、实时失败率、响应时间等多种策略,并可动态调整(通过配置中心推送)。
public interface ChannelRouter {
ChannelRoute selectRoute(PaymentRequest request);
}路由权重可基于历史成功率动态计算(如使用指数加权移动平均)。
每个渠道调用均使用Resilience4j的CircuitBreaker包装,并配置独立的熔断参数(失败率阈值、滑动窗口大小、重试超时等)。
@Bean
public CircuitBreaker channelBreaker() {
return CircuitBreaker.of("channelCB",
CircuitBreakerConfig.custom()
.failureRateThreshold(50)
.waitDurationInOpenState(Duration.ofSeconds(30))
.slidingWindowSize(100)
.build());
}
public ChannelResponse callChannel(ChannelRequest req) {
return circuitBreaker.executeSupplier(() -> {
// 实际HTTP调用,设置连接超时和读取超时
return httpClient.post(router.getEndpoint(), req);
});
}当熔断器开启时,自动降级至备用渠道(如从主通道切换至备通道)或返回明确错误码,同时触发告警。为避免瞬时流量导致频繁状态切换,需合理设置半开状态下的探测请求数。
实时风控需要多维特征:用户ID、交易金额、设备指纹、IP、小时级订单频次、黑名单状态等。
public class RiskContext {
private String userId;
private BigDecimal amount;
private String deviceId;
private String ip;
private Integer orderCountInHour;
private Boolean isBlacklist;
// 扩展字段
}我们采用轻量级表达式引擎(MVEL)而非重量级Drools,以降低运行时开销。规则从配置中心(Apollo/Nacos)动态加载,支持热更新。
规则定义示例(JSON格式,包含表达式、动作、优先级):
{
"condition": "amount > 50000 && orderCountInHour > 5",
"action": "REJECT",
"code": "RISK_001"
}执行引擎:
public class RiskRuleEngine {
private List<RiskRule> rules; // 定时刷新
public RiskResult evaluate(RiskContext ctx) {
for (RiskRule rule : rules) {
if (rule.matches(ctx)) {
return new RiskResult(rule.getAction(), rule.getCode());
}
}
return RiskResult.PASS;
}
}matches方法内使用MVEL.eval(condition, ctx),表达式编译结果可缓存以提升性能(实际测试单次评估 < 5ms)。风控日志异步写入Elasticsearch,用于离线模型训练和审计。
对账是资金安全的最终保障,我们将其分为四个阶段:
ReconciliationRecord(字段:渠道订单号、金额、状态、交易时间等)。由于单日文件可达数百万条,采用并行流 + ConcurrentHashMap加速查找:
public class ReconciliationEngine {
public List<DiffItem> compare(LocalRecords local, ChannelRecords channel) {
Map<String, LocalRecord> localMap = local.getRecords().stream()
.collect(Collectors.toConcurrentMap(LocalRecord::getChannelOrderNo, Function.identity()));
List<DiffItem> diffs = channel.getRecords().parallelStream().map(ch -> {
LocalRecord lr = localMap.get(ch.getChannelOrderNo());
if (lr == null) {
return DiffItem.missingLocal(ch);
}
if (!lr.getAmount().equals(ch.getAmount()) || !lr.getStatus().equals(ch.getStatus())) {
return DiffItem.amountMismatch(lr, ch);
}
return null;
}).filter(Objects::nonNull).collect(Collectors.toList());
// 处理本地有但渠道无的记录(长款)
// ...
return diffs;
}
}内存优化:对于超大文件,可采用分片加载或使用数据库临时表进行JOIN比对,避免OOM。
差异处理需生成唯一工单号,并记录处理状态(待处理、处理中、已调账、已驳回),确保可追溯。
客户端每次请求携带全局唯一requestId(UUID),服务端使用Redis的SETNX命令实现首次请求锁定,同时缓存响应结果。
public class IdempotentHandler {
@Autowired
private StringRedisTemplate redisTemplate;
public boolean tryAcquire(String requestId, String resultCacheKey) {
Boolean acquired = redisTemplate.opsForValue()
.setIfAbsent(requestId, "1", Duration.ofMinutes(5));
if (Boolean.TRUE.equals(acquired)) {
return true;
} else {
String cached = redisTemplate.opsForValue().get(resultCacheKey);
throw new DuplicateRequestException(cached);
}
}
public void cacheResult(String requestId, String result) {
redisTemplate.opsForValue().set("result:" + requestId, result, Duration.ofMinutes(5));
}
}在支付订单表上建立uk_request_id唯一索引,作为最终防线,确保即使Redis失效或并发穿透,也不会产生重复记录。
在消息消费端,通过messageId结合Redis分布式锁,保证同一消息仅被处理一次(处理成功后将ID存入已处理集合)。
集成Micrometer + Prometheus,自定义关键业务指标:
payment_request_total(按渠道、状态计数)payment_latency_seconds(分位数:50%、95%、99%)account_balance(实时余额,用于异常波动告警)reconciliation_diff_count(对账差异数)埋点示例:
@Autowired
private MeterRegistry meterRegistry;
public void recordPaymentResult(String channel, boolean success, long costMs) {
Counter.builder("payment.request.total")
.tag("channel", channel)
.tag("result", success ? "success" : "fail")
.register(meterRegistry)
.increment();
Timer.builder("payment.latency")
.publishPercentiles(0.5, 0.95, 0.99)
.register(meterRegistry)
.record(Duration.ofMillis(costMs));
}使用SkyWalking进行全链路追踪,每个请求注入traceId,贯穿网关、核心服务、渠道调用、消息队列,便于定位超时和异常节点。
配置Prometheus AlertManager规则,例如:
本文从支付核心引擎的实战角度,系统阐述了交易状态机、账务TCC柔性事务、基于本地消息表的最终一致性、动态路由与熔断降级、实时风控规则引擎、高效对账比对、全链路幂等防重以及可观测性设计。每一部分均以可落地的Java代码呈现,并涵盖常见异常场景的处理策略(超时、重试、死信、悬挂、空回滚等)。这套架构已在日均千万级交易量的生产环境中稳定运行,验证了其高可用性和数据一致性保障能力。
支付系统的本质是对一致性、可用性、安全性的持续追求,掌握上述核心模块的设计原理与实现细节,是构建可靠支付中台的基础。未来随着业务扩展,可在本架构上无缝接入更多新型支付方式,底层设计原则依然适用。
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。