一、 那个让全行停机三天的“完美”重构,到底错在哪?
我们先不聊理论,先讲一个发生在几年前、业内心照不宣的真实故事。
那是一家头部股份制银行,决心重构其核心账务系统。旧系统是上世纪90年代的单体架构,耦合严重,每次发版都要全行停机维护至少4小时。新方案看起来无懈可击:基于微服务架构,数据库采用最新的分布式数据库,承诺“双写零停机”,并且引入了强大的分布式事务框架(类似Seata的AT模式)。
上线那天,技术总监信心满满。然而,凌晨2点,系统开始报错。
问题出在一个看似普通的“转账”业务上。用户A向用户B转账100元。
- 服务A(账户服务)扣减A的余额100元,状态标记为“处理中”。
- 服务B(账户服务)增加B的余额100元,状态标记为“处理中”。
- 服务C(风控服务)异步检查这笔交易,发现A的账户存在可疑大额交易历史,触发风控阻断,回滚了服务A的操作。
但是,服务B的操作已经提交成功了。
更糟糕的是,由于分布式事务框架的网络超时重试用例设计缺陷,当服务A回滚时,框架试图通知服务B也回滚,但此时服务B的数据库连接池已满,响应超时。框架判定“部分成功”,将事务标记为“未知”而非“回滚”。
第二天早上9点,晨会还没开始,行长接到了监管机构的电话:某用户账户显示余额正常,但资金去向不明,经查实是被重复扣款且另一方未入账。
那家银行最终不得不回滚整个系统,损失惨重。
这个故事告诉我们一个血淋淋的真理:在支付领域,架构的“优雅”必须让位于资金的“安全”。任何事务一致性问题的代价,都是真金白银,甚至是金融牌照。
二、 为什么传统数据库事务在分布式环境下会“失效”?
要理解支付系统的复杂性,我们必须先回顾一下为什么单机数据库的事务如此可靠,而分布式环境下却如此脆弱。
2.1 ACID的破灭
在单体数据库中,事务的ACID特性(原子性Atomicity、一致性Consistency、隔离性Isolation、持久性Durability)由数据库引擎保证。当你执行BEGIN TRANSACTION; UPDATE account SET balance = balance - 100 WHERE id = 'A'; UPDATE account SET balance = balance + 100 WHERE id = 'B'; COMMIT;时,数据库内部通过undo log和redo log确保了这一过程的绝对安全。
然而,在微服务架构下,服务A和服务B可能部署在不同的服务器上,甚至连接不同的数据库实例。当你尝试跨服务调用时,你面对的是网络不可靠性和分布式事务难题。
2.2 分布式事务的三大困境
- 网络延迟与超时:服务B可能因为网络抖动暂时不可达,服务A已经扣款,但服务B的入账请求还在路上。如果服务A直接提交事务,资金就“消失”了;如果等待,系统性能会大幅下降。
- 幂等性缺失:在分布式环境中,消息可能会重复投递。如果支付网关发送了一次“扣款”指令,但由于网络问题没有收到ACK,网关可能会重试。如果代码没有做幂等处理,用户可能会被扣款两次。
- 最终一致性的挑战:既然强一致性(如2PC两阶段提交)会导致性能瓶颈和单点故障,我们通常选择“最终一致性”。但这要求开发者必须手动处理“如何保证最终一致”的复杂逻辑,稍有不慎就会出现数据不一致。
三、 第三方支付高并发架构实战:我们是如何构建“资金安全网”的
基于上述教训,我们在设计第三方支付库时,没有选择任何黑盒的分布式事务框架,而是构建了一套“以消息队列为核心、本地消息表为兜底、幂等性为基石”的可靠支付架构。
以下是核心模块的详细设计与代码实现。
3.1 核心设计原则:本地消息表(Local Message Table)
这是解决分布式事务最经典、最可靠的模式之一。其核心思想是:将分布式事务转化为本地事务+异步消息。
流程图解:
用户发起支付
|
v
[1. 业务数据库] 写入订单状态 "PENDING",同时写入一条"发送消息"的记录到[本地消息表]
|
v
[2. 本地事务] 提交数据库事务(保证订单和消息要么同时成功,要么同时失败)
|
v
[3. 定时任务/消息监听] 扫描本地消息表中"未发送"的记录
|
v
[4. 发送MQ消息] 将支付指令发送到MQ(如RocketMQ/Kafka)
|
v
[5. MQ消费] 下游服务(如账户服务)消费消息,执行扣款/入账
|
v
[6. 确认] 消费成功后,更新本地消息表状态为"已发送"
3.2 代码实现:订单创建与本地消息表
假设我们使用Spring Boot + MyBatis + RocketMQ。
Step 1: 数据库设计
我们需要两个表:order表和local_message表。
-- 订单表
CREATE TABLE `order` (
`id` bigint(20) NOT NULL AUTO_INCREMENT,
`order_no` varchar(64) NOT NULL COMMENT '订单号',
`user_id` bigint(20) NOT NULL,
`amount` decimal(10,2) NOT NULL COMMENT '金额',
`status` tinyint(1) NOT NULL DEFAULT 0 COMMENT '0:待支付, 1:支付中, 2:支付成功, 3:支付失败',
PRIMARY KEY (`id`),
UNIQUE KEY `uk_order_no` (`order_no`)
);
-- 本地消息表
CREATE TABLE `local_message` (
`id` bigint(20) NOT NULL AUTO_INCREMENT,
`msg_id` varchar(64) NOT NULL COMMENT '消息唯一ID',
`topic` varchar(64) NOT NULL COMMENT 'MQ Topic',
`content` text NOT NULL COMMENT '消息内容(JSON)',
`status` tinyint(1) NOT NULL DEFAULT 0 COMMENT '0:待发送, 1:已发送, 2:发送失败',
`create_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP,
`update_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
PRIMARY KEY (`id`),
UNIQUE KEY `uk_msg_id` (`msg_id`)
);
Step 2: 创建订单并插入本地消息(关键!)
这里的关键是使用@Transactional注解,确保订单创建和本地消息插入在同一个数据库事务中。
@Service
public class OrderService {
@Autowired
private OrderMapper orderMapper;
@Autowired
private LocalMessageMapper localMessageMapper;
@Autowired
private RocketMQTemplate rocketMQTemplate;
/**
* 创建订单,并确保本地消息表写入成功
*/
@Transactional(rollbackFor = Exception.class)
public OrderDTO createOrder(CreateOrderRequest request) {
// 1. 创建订单
Order order = new Order();
order.setOrderNo(UUID.randomUUID().toString().replace("-", ""));
order.setUserId(request.getUserId());
order.setAmount(request.getAmount());
order.setStatus(OrderStatus.PENDING);
orderMapper.insert(order);
// 2. 构造本地消息
// 注意:此时订单状态仍是PENDING,我们发送的消息是“处理支付”
String messageContent = JSON.toJSONString(new PayMessage(order.getOrderNo(), order.getAmount()));
LocalMessage localMessage = new LocalMessage();
localMessage.setMsgId(UUID.randomUUID().toString().replace("-", ""));
localMessage.setTopic("PAYMENT_TOPIC");
localMessage.setContent(messageContent);
localMessage.setStatus(MessageStatus.PENDING);
localMessageMapper.insert(localMessage);
// 3. 更新订单状态为“支付中”
orderMapper.updateStatus(order.getId(), OrderStatus.PAYING);
// 返回订单信息
OrderDTO dto = new OrderDTO();
dto.setOrderNo(order.getOrderNo());
dto.setStatus(order.getStatus());
return dto;
}
}
这里有一个非常重要的细节: 我们并没有在创建订单后立即发送MQ消息,而是将消息写入本地表。这是为了防止服务宕机后消息丢失。
3.3 异步发送MQ:定时任务兜底
由于消息写入了本地表,我们需要一个机制将这些消息发送到MQ。通常有两种方式:
- 事务消息(RocketMQ):更优雅,但配置复杂。
- 定时任务扫描:更可控,适合复杂业务。
我们选择定时任务扫描,因为我们可以精确控制重试逻辑。
@Component
public class LocalMessageTask {
@Autowired
private LocalMessageMapper localMessageMapper;
@Autowired
private RocketMQTemplate rocketMQTemplate;
/**
* 每5秒扫描一次,发送状态为“待发送”的消息
*/
@Scheduled(fixedDelay = 5000)
public void sendPendingMessages() {
List<LocalMessage> pendingMessages = localMessageMapper.selectByStatus(MessageStatus.PENDING);
for (LocalMessage msg : pendingMessages) {
try {
// 发送MQ消息
rocketMQTemplate.syncSend("PAYMENT_TOPIC", MessageBuilder.withPayload(msg.getContent()).build());
// 发送成功后,更新本地消息状态为“已发送”
localMessageMapper.updateStatus(msg.getId(), MessageStatus.SENT);
log.info("消息发送成功, msgId: {}", msg.getMsgId());
} catch (Exception e) {
log.error("消息发送失败, msgId: {}", msg.getMsgId(), e);
// 更新状态为“发送失败”,等待下次重试
localMessageMapper.updateStatus(msg.getId(), MessageStatus.FAILED);
}
}
}
}
关键点: 如果MQ发送失败,我们只更新本地消息状态,而不会删除消息。下次定时任务会再次尝试发送。这保证了消息不丢失。
3.4 幂等性:防止重复扣款
这是支付系统最核心的安全屏障。在分布式系统中,网络重试、消息重复投递是常态。我们必须确保同一笔订单,无论收到多少次支付请求,只处理一次。
数据库层面的幂等性
我们可以在local_message表或支付流水表中添加唯一索引。
-- 支付流水表,确保同一订单号只有一条成功记录
CREATE TABLE `payment_flow` (
`id` bigint(20) NOT NULL AUTO_INCREMENT,
`order_no` varchar(64) NOT NULL,
`amount` decimal(10,2) NOT NULL,
`status` tinyint(1) NOT NULL DEFAULT 0,
`create_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP,
PRIMARY KEY (`id`),
UNIQUE KEY `uk_order_no` (`order_no`) -- 关键!唯一索引保证幂等
);
代码层面的幂等性检查
在消费MQ消息时,我们必须先检查是否已经处理过该订单。
@Component
public class PaymentConsumer {
@Autowired
private OrderService orderService;
@Autowired
private PaymentFlowMapper paymentFlowMapper;
@RocketMQMessageListener(topic = "PAYMENT_TOPIC", consumerGroup = "payment_consumer_group")
public void onMessage(String messageContent) {
PayMessage payMessage = JSON.parseObject(messageContent, PayMessage.class);
String orderNo = payMessage.getOrderNo();
try {
// 1. 幂等性检查:查询是否已有成功处理的支付流水
PaymentFlow existingFlow = paymentFlowMapper.selectByOrderNo(orderNo);
if (existingFlow != null && existingFlow.getStatus() == PaymentStatus.SUCCESS) {
log.warn("订单已处理,跳过重复消息, orderNo: {}", orderNo);
return;
}
// 2. 开启数据库事务,执行扣款和入账
// 注意:这里使用分布式锁或乐观锁防止并发重复处理
executePayment(orderNo, payMessage.getAmount());
// 3. 记录支付流水
PaymentFlow flow = new PaymentFlow();
flow.setOrderNo(orderNo);
flow.setAmount(payMessage.getAmount());
flow.setStatus(PaymentStatus.SUCCESS);
paymentFlowMapper.insert(flow);
} catch (DuplicateKeyException e) {
// 捕获唯一索引冲突,说明并发情况下已被其他线程处理
log.warn("并发重复处理, orderNo: {}", orderNo);
} catch (Exception e) {
log.error("支付处理失败, orderNo: {}", orderNo, e);
// 可以根据业务需要进行重试或人工介入
}
}
@Transactional(rollbackFor = Exception.class)
private void executePayment(String orderNo, BigDecimal amount) {
// 1. 检查订单状态,防止状态不一致
Order order = orderMapper.selectByOrderNo(orderNo);
if (order == null || order.getStatus() != OrderStatus.PAYING) {
throw new BusinessException("订单状态异常");
}
// 2. 扣减用户余额(假设有一个account服务)
// 这里可以是直接调RPC,也可以是利用分布式事务
// 为了简化,假设我们直接操作账户表(同一数据库)
int affectedRows = accountMapper.deductBalance(order.getUserId(), amount);
if (affectedRows == 0) {
throw new BusinessException("余额不足");
}
// 3. 更新订单状态为支付成功
orderMapper.updateStatus(orderNo, OrderStatus.SUCCESS);
}
}
关键点:
- 幂等性检查在前:在真正执行扣款前,先查数据库,如果已处理则直接返回。
- 事务包裹核心逻辑:扣款和订单状态更新必须在同一个事务中。
- 捕获并发异常:即使有幂等检查,高并发下仍可能出现两个请求同时通过检查。此时,数据库唯一索引的
DuplicateKeyException是最好的兜底。
3.5 对账系统:最后一道防线
无论代码写得多么完美,都可能存在意外。因此,对账系统是支付安全的终极保障。
对账系统的核心逻辑是:每日凌晨,拉取第三方支付渠道(如支付宝、微信、银行)的昨日交易流水,与本地数据库的订单进行比对,找出差异。
@Service
public class ReconciliationService {
@Autowired
private PaymentChannelService channelService; // 第三方支付渠道接口
@Autowired
private OrderMapper orderMapper;
@Autowired
private DiscrepancyRecordMapper discrepancyMapper;
/**
* 执行每日对账
*/
public void dailyReconciliation(LocalDate date) {
// 1. 从第三方支付渠道拉取昨日所有成功交易
List<ChannelTransaction> channelTransactions = channelService.queryTransactions(date);
// 2. 从本地数据库拉取昨日所有成功订单
List<Order> localOrders = orderMapper.selectByDateAndStatus(date, OrderStatus.SUCCESS);
// 3. 构建Map,方便比对
Map<String, Order> orderMap = localOrders.stream()
.collect(Collectors.toMap(Order::getOrderNo, o -> o, (o1, o2) -> o1));
Set<String> channelOrderNos = channelTransactions.stream()
.map(ChannelTransaction::getOrderNo)
.collect(Collectors.toSet());
Set<String> localOrderNos = orderMap.keySet();
// 4. 找出单边账(渠道有,本地没有)
for (String orderNo : channelOrderNos) {
if (!localOrderNos.contains(orderNo)) {
// 长款:用户已付款,但本地订单未更新
log.error("发现长款,订单号: {}", orderNo);
discrepancyMapper.insert(new DiscrepancyRecord(orderNo, DiscrepancyType.OVER, date));
// 触发自动退款或人工介入
}
}
// 5. 找出短账(本地有,渠道没有)
for (String orderNo : localOrderNos) {
if (!channelOrderNos.contains(orderNo)) {
// 短款:本地显示成功,但渠道无记录
log.error("发现短款,订单号: {}", orderNo);
discrepancyMapper.insert(new DiscrepancyRecord(orderNo, DiscrepancyType.UNDER, date));
// 触发补单或人工核查
}
}
// 6. 比对金额是否一致
for (String orderNo : channelOrderNos) {
if (localOrderNos.contains(orderNo)) {
ChannelTransaction channelTx = channelTransactions.stream()
.filter(t -> t.getOrderNo().equals(orderNo))
.findFirst().orElse(null);
Order localOrder = orderMap.get(orderNo);
if (channelTx.getAmount().compareTo(localOrder.getAmount()) != 0) {
log.error("金额不一致,订单号: {}, 渠道: {}, 本地: {}",
orderNo, channelTx.getAmount(), localOrder.getAmount());
discrepancyMapper.insert(new DiscrepancyRecord(orderNo, DiscrepancyType.AMOUNT_ERROR, date));
}
}
}
}
}
关键点:
- 对账不是实时进行的,而是T+1(隔天)进行。
- 发现差异后,系统应自动触发退款或补单流程,而不是仅仅记录日志。
- 对账结果是财务审计的重要依据。
四、 高并发下的性能优化:如何支撑每秒万级交易?
上述架构保证了正确性,但在高并发场景下,性能可能成为瓶颈。以下是几个关键的优化
