电商大促MySQL主从同步延迟导致订单数据丢失 事务锁机制与一致性维护方案全解析
大家好,我是老张,一个在电商系统里摸爬滚打近十年的后端工程师。今天想跟大家聊一个让很多人头秃的问题——大促期间MySQL主从同步延迟,以及由此引发的订单数据丢失。
说真的,这个问题我亲身经历过。2018年的双11,我们公司的订单系统崩了,用户付了钱却查不到订单,客服电话被打爆。后来排查了好几天,才发现是主从同步延迟导致的订单数据不一致问题。那种感觉,真的比失恋还难受。
但正是那次经历,让我彻底搞明白了MySQL事务锁机制和一致性维护的原理。今天我就把这些年的经验,毫无保留地分享给大家。
一、大促场景下的数据一致性挑战
1.1 主从同步延迟的成因
在了解解决方案之前,我们先搞清楚为什么会出现同步延迟。
MySQL的主从同步,本质上是一个异步复制的过程。主库(Master)负责处理写操作,从库(Slave)负责读取。当主库执行完一个事务后,会将会写操作记录到binlog中,然后从库通过I/O线程读取binlog,再由SQL线程执行这些操作,从而实现数据同步。
这个过程听起来很简单,但实际上有很多环节可能导致延迟:
-- 查看主库binlog状态
SHOW MASTER STATUS;
-- 查看从库同步状态
SHOW SLAVE STATUS\G
-- 关键参数解读:
-- Seconds_Behind_Master: 从库落后主库的秒数
-- Relay_Log_Space: 中继日志大小
-- Last_IO_Error: I/O线程错误信息
-- Last_SQL_Error: SQL线程错误信息
在大促场景下,以下几个因素会加剧同步延迟:
1. 大事务问题
一个常见的大事务场景:
-- 典型的大事务示例(错误示范)
START TRANSACTION;
-- 扣减库存(可能涉及多张表)
UPDATE product_stock SET stock = stock - 1 WHERE product_id = 10086;
UPDATE product_stock SET stock = stock - 1 WHERE product_id = 10087;
UPDATE product_stock SET stock = stock - 1 WHERE product_id = 10088;
-- 创建订单
INSERT INTO orders (order_no, product_id, user_id, amount) VALUES ('ORD20241111001', 10086, 1001, 999.00);
INSERT INTO orders (order_no, product_id, user_id, amount) VALUES ('ORD20241111002', 10087, 1001, 899.00);
INSERT INTO orders (order_no, product_id, user_id, amount) VALUES ('ORD20241111003', 10088, 1001, 799.00);
-- 更新用户余额
UPDATE user_account SET balance = balance - 2697.00 WHERE user_id = 1001;
-- 发送MQ消息
-- ...
COMMIT;
这个事务包含多个UPDATE和INSERT操作,锁定的行数和时长都很长。从库需要按顺序回放所有binlog事件,单个大事务的回放时间可能长达几秒甚至几十秒。
2. 锁竞争问题
-- 检查当前锁等待情况
SELECT
r.trx_id waiting_trx_id,
r.trx_mysql_thread_id waiting_thread,
r.trx_query waiting_query,
b.trx_id blocking_trx_id,
b.trx_mysql_thread_id blocking_thread,
b.trx_query blocking_query
FROM information_schema.innodb_lock_waits r
JOIN information_schema.innodb_trx b ON r.requesting_trx_id = b.trx_id;
大促期间,大量请求同时竞争同一批热点商品的数据,会导致严重的锁等待。主库处理锁等待的时间变长,binlog写入延迟,从库回放时间也跟着变长。
3. 从库性能瓶颈
从库的硬件配置通常不如主库,而且从库可能同时服务多个业务场景的查询,负载不均时容易出现瓶颈。
-- 检查从库的复制延迟
SELECT
PROCESSLIST_ID AS thread_id,
TIME AS seconds_behind,
STATE,
INFO AS last_query
FROM information_schema.processlist
WHERE COMMAND = 'Sleep' OR processlist_id IN (
SELECT THREAD_ID FROM performance_schema.threads WHERE NAME LIKE 'thread/sql/slave%'
);
1.2 订单数据丢失的具体场景
理解了延迟成因,我们来看看在什么情况下会导致订单数据”丢失”。
场景一:基于从库的读请求导致数据不一致
这是最常见的情况。用户下单后,系统直接查询从库确认订单,但由于从库同步延迟,此时从库还没有这个订单的数据:
// 伪代码示例:有问题的订单查询逻辑
public Order queryOrder(String orderNo) {
// 直接查询从库
Order order = orderMapper.selectByOrderNo(orderNo);
// 问题:如果从库同步延迟,这里可能返回null
if (order == null) {
log.warn("订单查询结果为空, orderNo={}", orderNo);
return null;
}
return order;
}
用户看到”订单不存在”,可能会认为支付失败,然后重新下单。这就导致了重复下单,或者用户以为没下单成功但实际上已经成功的情况。
场景二:binlog事件丢失或损坏
虽然概率很低,但网络抖动、磁盘故障等原因可能导致binlog事件丢失。从库没有收到某些事件,就会导致数据不一致。
-- 检查binlog事件完整性
SHOW BINARY LOGS;
-- 查看特定binlog文件的事件
SHOW BINLOG EVENTS IN 'mysql-bin.000001';
-- 检查是否有binlog损坏
mysqlbinlog --force-read mysql-bin.000001 > /tmp/binlog_check.txt
场景三:主库崩溃导致数据未持久化
这是一个极端但严重的情况。主库在事务提交前崩溃,MySQL的崩溃恢复机制可能无法恢复所有数据,导致主从数据不一致。
典型场景:
1. 应用连接主库,开始事务
2. 执行INSERT INTO orders...
3. 执行到一半,主库宕机
4. 主库崩溃恢复后,该事务未完成
5. 从库的binlog中没有这个事件
6. 用户看到主库数据不一致,从库也没有数据
场景四:半同步复制超时
为了保证数据一致性,很多公司开启了半同步复制。但如果从库延迟太大,半同步复制会超时,回退到异步模式,这时候就可能出现数据丢失。
-- 查看半同步复制状态
SHOW VARIABLES LIKE 'rpl_semi_sync_%';
-- 查看半同步复制的状态
SHOW STATUS LIKE 'Rpl_semi_sync_%';
二、事务锁机制详解
要解决一致性问题,首先得理解MySQL的事务锁机制。
2.1 InnoDB的锁类型
InnoDB支持多种锁,理解它们的区别对解决数据一致性问题至关重要。
行锁(Record Lock)
-- 演示行锁
START TRANSACTION;
-- 这里的SELECT ... FOR UPDATE会加行锁
SELECT * FROM product_stock WHERE product_id = 10086 FOR UPDATE;
-- 此时其他事务无法修改这行数据
-- 如果另一个事务执行:
-- UPDATE product_stock SET stock = stock - 1 WHERE product_id = 10086;
-- 会阻塞等待锁释放
-- 解锁后提交
COMMIT;
行锁只锁定满足条件的行,是最常用的锁类型。在大促场景下,对热点商品的库存扣减必须使用行锁,否则会出现超卖问题。
间隙锁(Gap Lock)
间隙锁锁定的是两个记录之间的”间隙”,而不是记录本身。
-- 假设product_stock表中有id为1, 3, 5的记录
-- 执行以下查询会加间隙锁,锁定(1,3)和(3,5)之间的范围
SELECT * FROM product_stock WHERE id > 1 AND id < 5 FOR UPDATE;
-- 其他事务无法在这个范围内插入记录
-- INSERT INTO product_stock (id, product_id) VALUES (2, 9999); 会被阻塞
间隙锁的主要作用是防止幻读。在REPEATABLE READ隔离级别下,InnoDB使用 next-key lock(记录锁 + 间隙锁)来避免幻读。
临键锁(Next-Key Lock)
临键锁是记录锁和间隙锁的组合,锁定的是一个范围,并且包含范围内的记录。
next-key lock = record lock + gap lock
例如:(1, 3] 表示锁定id > 1 且 id <= 3 的范围
2.2 隔离级别与锁的关系
MySQL支持四种隔离级别,不同级别使用不同的锁机制:
-- 查看当前隔离级别
SELECT @@transaction_isolation;
-- 设置隔离级别为READ COMMITTED
SET TRANSACTION ISOLATION LEVEL READ COMMITTED;
-- 设置隔离级别为REPEATABLE READ(默认)
SET TRANSACTION ISOLATION LEVEL REPEATABLE READ;
| 隔离级别 | 脏读 | 不可重复读 | 幻读 | 锁机制 |
|---|---|---|---|---|
| READ UNCOMMITTED | 可能 | 可能 | 可能 | 无锁 |
| READ COMMITTED | 不可能 | 可能 | 可能 | 行锁 |
| REPEATABLE READ | 不可能 | 不可能 | 可能 | next-key lock |
| SERIALIZABLE | 不可能 | 不可能 | 不可能 | 表锁 |
在大促场景下,推荐使用REPEATABLE READ隔离级别。虽然性能略低于READ COMMITTED,但能避免不可重复读和大部分幻读问题。
2.3 乐观锁与悲观锁
悲观锁:基于数据库锁机制
// 悲观锁示例:通过SELECT ... FOR UPDATE实现
@Transactional
public Order createOrder(OrderRequest request) {
// 先查询并加锁
ProductStock stock = productStockMapper.selectForUpdate(request.getProductId());
// 检查库存
if (stock.getStock() < request.getQuantity()) {
throw new BusinessException("库存不足");
}
// 扣减库存
productStockMapper.updateStock(request.getProductId(), request.getQuantity());
// 创建订单
Order order = buildOrder(request);
orderMapper.insert(order);
return order;
}
悲观锁简单直接,但在高并发场景下,频繁的锁竞争会导致性能问题。
乐观锁:基于版本号或时间戳
// 乐观锁示例:通过版本号实现
@Transactional
public Order createOrderOptimistic(OrderRequest request) {
// 查询商品和版本号
ProductStock stock = productStockMapper.selectById(request.getProductId());
// 检查库存
if (stock.getStock() < request.getQuantity()) {
throw new BusinessException("库存不足");
}
// 更新库存,带上版本号条件
int updated = productStockMapper.updateStockWithVersion(
request.getProductId(),
request.getQuantity(),
stock.getVersion()
);
// 检查是否被其他事务修改过
if (updated == 0) {
throw new BusinessException("库存更新失败,请重试");
}
// 创建订单
Order order = buildOrder(request);
orderMapper.insert(order);
return order;
}
乐观锁在高并发场景下性能更好,但需要处理冲突重试的逻辑。
两种锁的对比:
悲观锁:
优点:实现简单,数据安全性高
缺点:高并发下锁竞争激烈,性能较差
适用场景:写多读少,冲突概率低的场景
乐观锁:
优点:不加锁,并发性能好
缺点:冲突时需要重试,实现复杂
适用场景:读多写少,冲突概率高的场景
在大促场景下,库存扣减通常使用悲观锁(SELECT … FOR UPDATE),因为超卖问题不能容忍。而订单创建等场景可以使用乐观锁,提高并发性能。
2.4 锁超时与死锁处理
-- 设置锁等待超时时间(默认50秒)
SET innodb_lock_wait_timeout = 30;
-- 查看锁等待超时设置
SELECT @@innodb_lock_wait_timeout;
-- 检查死锁情况
SHOW ENGINE INNODB STATUS\G
-- 死锁日志通常在错误日志中
-- 路径:datadir/hostname.err
// 死锁重试示例
public Order createOrderWithRetry(OrderRequest request) {
int maxRetries = 3;
for (int i = 0; i < maxRetries; i++) {
try {
return createOrder(request);
} catch (DeadlockLoserDataAccessException e) {
log.warn("死锁发生,第{}次重试", i + 1);
if (i == maxRetries - 1) {
throw e;
}
// 随机等待后重试,避免雪崩
Thread.sleep((long) (Math.random() * 100 + 50));
}
}
return null;
}
三、一致性维护方案
理解了锁机制,我们来看看具体的解决方案。
3.1 读写分离的数据一致性保障
方案一:主库写,主库读
最直接的方案:下单后立即从主库查询订单。
// 使用主库读取
public Order queryOrderAfterCreate(String orderNo) {
// 强制使用主库
Order order = orderMapper.selectByOrderNoWithMaster(orderNo);
// 检查订单状态
if (order != null && "PAID".equals(order.getStatus())) {
return order;
}
// 如果主库也没有,说明订单确实不存在
return null;
}
这种方式简单可靠,但会加重主库的读压力。在大促场景下,主库的负载已经很重了,再承担读压力可能会成为瓶颈。
方案二:短暂延迟读取从库
// 延迟读取从库
public Order queryOrderWithDelay(String orderNo) {
// 先查询主库,确认订单存在
Order mainOrder = orderMapper.selectByOrderNoWithMaster(orderNo);
if (mainOrder == null) {
return null;
}
// 等待从库同步(最多等待5秒)
long startTime = System.currentTimeMillis();
while (System.currentTimeMillis() - startTime < 5000) {
Order slaveOrder = orderMapper.selectByOrderNoWithSlave(orderNo);
// 检查从库数据是否一致
if (slaveOrder != null &&
Objects.equals(slaveOrder.getStatus(), mainOrder.getStatus())) {
return slaveOrder;
}
// 短暂休眠后重试
try {
Thread.sleep(100);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
break;
}
}
// 超时后返回主库数据
return mainOrder;
}
这种方式平衡了性能和一致性,但需要处理超时逻辑。
方案三:基于binlog的实时同步检查
// 监听binlog,实时检查数据同步状态
@Component
public class BinlogSyncChecker {
@Autowired
private MysqlBinlogClient binlogClient;
@Autowired
private OrderService orderService;
@PostConstruct
public void startListening() {
binlogClient.listen(new BinlogEventListener() {
@Override
public void onEvent(BinlogEvent event) {
if (event.getType() == BinlogEventType.QUERY_EVENT) {
String sql = event.getSql();
if (sql.contains("INSERT INTO orders")) {
// 解析订单号
String orderNo = parseOrderNo(sql);
// 异步检查从库数据
CompletableFuture.runAsync(() -> {
waitAndVerify(orderNo);
});
}
}
}
});
}
private void waitAndVerify(String orderNo) {
long timeout = System.currentTimeMillis() + 5000;
while (System.currentTimeMillis() < timeout) {
Order masterOrder = orderService.queryFromMaster(orderNo);
Order slaveOrder = orderService.queryFromSlave(orderNo);
if (slaveOrder != null &&
Objects.equals(masterOrder.getId(), slaveOrder.getId())) {
log.info("订单同步成功: {}", orderNo);
return;
}
try {
Thread.sleep(100);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
return;
}
}
// 同步超时,报警
log.error("订单同步超时: {}", orderNo);
alertService.sendAlert("订单同步超时: " + orderNo);
}
}
这种方式可以实现实时监控和告警,但实现复杂度较高。
3.2 基于分布式事务的一致性保障
方案一:TCC事务模型
TCC(Try-Confirm-Cancel)是一种分布式事务模型,适用于跨多个服务的事务场景。
// TCC事务示例
public interface OrderTccService {
// Try阶段:尝试执行业务,预留资源
@TccTransaction(confirmMethod = "confirm", cancelMethod = "cancel")
boolean tryCreateOrder(OrderRequest request);
// Confirm阶段:确认事务,提交资源
void confirm(OrderRequest request);
// Cancel阶段:取消事务,释放资源
void cancel(OrderRequest request);
}
@Service
public class OrderTccServiceImpl implements OrderTccService {
@Autowired
private OrderMapper orderMapper;
@Autowired
private ProductStockMapper stockMapper;
@Autowired
private UserAccountMapper accountMapper;
@Override
public boolean tryCreateOrder(OrderRequest request) {
// 1. 扣减库存(预留)
int updated = stockMapper.tryDecreaseStock(
request.getProductId(),
request.getQuantity()
);
if (updated == 0) {
return false;
}
// 2. 冻结用户余额(预留)
int frozen = accountMapper.freezeBalance(
request.getUserId(),
request.getAmount()
);
if (frozen == 0) {
// 库存回滚
stockMapper.rollbackStock(
request.getProductId(),
request.getQuantity()
);
return false;
}
// 3. 创建订单(状态为TRYING)
Order order = buildOrder(request);
order.setStatus("TRYING");
orderMapper.insert(order);
return true;
}
@Override
public void confirm(OrderRequest request) {
// 更新订单状态为已确认
orderMapper.updateStatus(request.getOrderNo(), "CONFIRMED");
// 扣减库存(从预留转为实际扣减)
stockMapper.confirmDecreaseStock(
request.getProductId(),
request.getQuantity()
);
// 扣减用户余额(从冻结转为实际扣减)
accountMapper.confirmDeductBalance(
request.getUserId(),
request.getAmount()
);
}
@Override
public void cancel(OrderRequest request) {
// 更新订单状态为已取消
orderMapper.updateStatus(request.getOrderNo(), "CANCELLED");
// 释放库存
stockMapper.releaseStock(
request.getProductId(),
request.getQuantity()
);
// 解冻用户余额
accountMapper.unfreezeBalance(
request.getUserId(),
request.getAmount()
);
}
}
TCC模型的优点是性能好,不依赖数据库锁;缺点是实现复杂,需要处理各个阶段的异常情况。
方案二:本地消息表
本地消息表是一种最终一致性的实现方案,适用于异步场景。
// 本地消息表实现
@Service
public class OrderLocalMessageService {
@Autowired
private OrderMapper orderMapper;
@Autowired
private LocalMessageMapper messageMapper;
@Autowired
private MessageProducer messageProducer;
@Transactional
public String createOrder(OrderRequest request) {
// 1. 创建订单
Order order = buildOrder(request);
order.setStatus("CREATING");
orderMapper.insert(order);
// 2. 发送本地消息(在同一事务中)
LocalMessage message = new LocalMessage();
message.setBizId(order.getId());
message.setBizType("ORDER_CREATE");
message.setBizData(JSON.toJSONString(request));
message.setStatus("PENDING");
messageMapper.insert(message);
return order.getOrderNo();
}
// 定时任务扫描未发送的消息
@Scheduled(fixedDelay = 1000)
public void sendPendingMessages() {
List<LocalMessage> pendingMessages = messageMapper.selectPendingMessages(100);
for (LocalMessage message : pendingMessages) {
try {
// 发送MQ消息
messageProducer.send("order.create", message.getBizData());
// 更新消息状态
messageMapper.updateStatus(message.getId(), "SENT");
} catch (Exception e) {
log.error("消息发送失败, messageId={}", message.getId(), e);
messageMapper.incrementRetryCount(message.getId());
}
}
}
}
本地消息表的优点是简单可靠,实现了事务和消息的本地一致性;缺点是需要额外的表和处理逻辑。
3.3 基于MQ的最终一致性
消息队列是实现最终一致性的常用方案。
// MQ事务消息实现
@Service
public class OrderMQService {
@Autowired
private RocketMQTemplate rocketMQTemplate;
@Autowired
private OrderMapper orderMapper;
@Autowired
private ProductStockMapper stockMapper;
public String createOrderWithMQ(OrderRequest request) {
// 1. 发送事务消息
Message<OrderRequest> message = MessageBuilder
.withPayload(request)
.setHeader("bizType", "ORDER_CREATE")
.build();
// 发送事务消息
TransactionSendResult result = rocketMQTemplate.sendMessageInTransaction(
message,
request
);
if (result.getSendStatus() != SendStatus.SEND_OK) {
throw new RuntimeException("消息发送失败");
}
// 2. 异步处理订单(在另一个事务中)
CompletableFuture.runAsync(() -> {
processOrder(request);
});
return request.getOrderNo();
}
private void processOrder(OrderRequest request) {
try {
// 扣减库存
stockMapper.decreaseStock(request.getProductId(), request.getQuantity());
// 创建订单
Order order = buildOrder(request);
order.setStatus("PAID");
orderMapper.insert(order);
} catch (Exception e) {
log.error("订单处理失败", e);
// 发送失败消息,后续处理
sendMessage("order.create.failed", request);
}
}
}
MQ方案的优点是解耦性好,性能高;缺点是需要处理消息重复、顺序等问题。
3.4 数据对账与修复
无论采用什么方案,数据对账都是必不可少的。
// 数据对账服务
@Service
public class DataReconciliationService {
@Autowired
private OrderMapper masterOrderMapper;
@Autowired
private OrderMapper slaveOrderMapper;
@Autowired
private AlertService alertService;
// 定时对账
@Scheduled(cron = "0 0/5 * * * ?")
public void reconcileOrders() {
long startTime = System.currentTimeMillis() - 5 * 60 * 1000; // 最近5分钟
// 查询主库订单
List<Order> masterOrders = masterOrderMapper.selectByTimeRange(startTime, System.currentTimeMillis());
// 查询从库订单
List<Order> slaveOrders = slaveOrderMapper.selectByTimeRange(startTime, System.currentTimeMillis());
// 对比订单
Map<String, Order> slaveOrderMap = slaveOrders.stream()
.collect(Collectors.toMap(Order::getOrderNo, o -> o, (o1, o2) -> o1));
List<String> missingOrders = new ArrayList<>();
for (Order masterOrder : masterOrders) {
Order slaveOrder = slaveOrderMap.get(masterOrder.getOrderNo());
if (slaveOrder == null) {
// 从库缺少订单
missingOrders.add(masterOrder.getOrderNo());
} else if (!Objects.equals(masterOrder.getStatus(), slaveOrder.getStatus())) {
// 状态不一致
log.warn("订单状态不一致: master={}, slave={}, orderNo={}",
masterOrder.getStatus(), slaveOrder.getStatus(), masterOrder.getOrderNo());
}
}
if (!missingOrders.isEmpty()) {
log.error("发现{}个订单未同步到从库", missingOrders.size());
alertService.sendAlert("数据不一致告警: " + missingOrders);
// 触发重新同步
synchronizeOrders(missingOrders);
}
}
private void synchronizeOrders(List<String> orderNos) {
for (String orderNo : orderNos) {
try {
// 从主库重新查询并写入从库
Order order = masterOrderMapper.selectByOrderNo(orderNo);
if (order != null) {
slaveOrderMapper.insert(order);
log.info("订单同步成功: {}", orderNo);
}
} catch (Exception e) {
log.error("订单同步失败: {}", orderNo, e);
}
}
}
}
四、实战案例:双十一订单系统优化
让我分享一下我们团队在2019年双十一的实际优化经验。
4.1 问题背景
2018年双十一,我们的订单系统出现了严重问题:
- 大量用户反映下单后查不到订单
- 客服投诉量激增
- 系统稳定性评分大幅下降
经过排查,发现主要原因:
- 主从同步延迟达到30秒以上
- 下单后直接查询从库确认订单
- 没有数据对账机制
4.2 优化方案
第一阶段:紧急修复(1周内)
// 修改订单查询逻辑
public Order queryOrder(String orderNo, boolean forceMaster) {
if (forceMaster) {
return orderMapper.selectByOrderNoWithMaster(orderNo);
}
// 查询缓存
Order cachedOrder = orderCache.get(orderNo);
if (cachedOrder != null) {
return cachedOrder;
}
// 查询从库,但带上重试逻辑
long maxRetries = 3;
for (int i = 0; i < maxRetries; i++) {
Order slaveOrder = orderMapper.selectByOrderNoWithSlave(orderNo);
if (slaveOrder != null) {
// 同步到缓存
orderCache.put(orderNo, slaveOrder);
return slaveOrder;
}
try {
Thread.sleep(200);
} catch (InterruptedException e) {
break;
}
}
// 从库查询失败,回退到主库
log.warn("从库查询超时,回退到主库查询: {}", orderNo);
return orderMapper.selectByOrderNoWithMaster(orderNo);
}
第二阶段:架构优化(1个月内)
-- 优化主从复制配置
-- 主库配置
[mysqld]
# 开启半同步复制
plugin-load=rpl_semi_sync_master=semisync_master.so
rpl_semi_sync_master_enabled=1
rpl_semi_sync_master_timeout=1000 # 1秒超时
# 从库配置
[mysqld]
plugin-load=rpl_semi_sync_slave=semisync_slave.so
rpl_semi_sync_slave_enabled=1
# 优化复制性能
relay_log_info_repository=TABLE
master_info_repository=TABLE
// 引入TCC事务
@Service
public class OrderTccServiceImpl implements OrderTccService {
@Autowired
private OrderMapper orderMapper;
@Autowired
private StockService stockService;
@Autowired
private AccountService accountService;
@Override
@TccTransaction(confirmMethod = "confirm", cancelMethod = "cancel")
public boolean tryCreateOrder(OrderRequest request) {
// Try: 预留资源
boolean stockReserved = stockService.tryReserveStock(
request.getProductId(),
request.getQuantity()
);
if (!stockReserved) {
return false;
}
boolean balanceFrozen = accountService.tryFreezeBalance(
request.getUserId(),
request.getAmount()
);
if (!balanceFrozen) {
stockService.releaseStock(request.getProductId(), request.getQuantity());
return false;
}
// 创建订单(状态为TRYING)
Order order = buildOrder(request);
order.setStatus("TRYING");
orderMapper.insert(order);
return true;
}
@Override
public void confirm(OrderRequest request) {
// Confirm: 提交资源
stockService.confirmReserveStock(request.getProductId(), request.getQuantity());
accountService.confirmDeductBalance(request.getUserId(), request.getAmount());
// 更新订单状态
orderMapper.updateStatus(request.getOrderNo(), "CONFIRMED");
}
@Override
public void cancel(OrderRequest request) {
// Cancel: 释放资源
stockService.releaseStock(request.getProductId(), request.getQuantity());
accountService.unfreezeBalance(request.getUserId(), request.getAmount());
// 更新订单状态
orderMapper.updateStatus(request.getOrderNo(), "CANCELLED");
}
}
第三阶段:数据对账(长期)
// 每日对账任务
@Component
public class DailyReconciliationTask {
@Autowired
private OrderMapper masterOrderMapper;
@Autowired
private OrderMapper slaveOrderMapper;
@Autowired
private AlertService alertService;
@Scheduled(cron = "0 30 2 * * ?")
public void dailyReconcile() {
LocalDate yesterday = LocalDate.now().minusDays(1);
// 查询昨天所有订单
List<Order> masterOrders = masterOrderMapper.selectAllByDate(yesterday);
List<Order> slaveOrders = slaveOrderMapper.selectAllByDate(yesterday);
// 对比
Map<String, Order> slaveOrderMap = slaveOrders.stream()
.collect(Collectors.toMap(Order::getOrderNo, o -> o));
List<String> inconsistentOrders = new ArrayList<>();
for (Order masterOrder : masterOrders) {
Order slaveOrder = slaveOrderMap.get(masterOrder.getOrderNo());
if (slaveOrder == null) {
inconsistentOrders.add(masterOrder.getOrderNo() + "(缺少)");
} else if (!Objects.equals(masterOrder.getStatus(), slaveOrder.getStatus())) {
inconsistentOrders.add(masterOrder.getOrderNo() + "(状态不一致)");
}
}
if (!inconsistentOrders.isEmpty()) {
log.error("发现{}个不一致订单", inconsistentOrders.size());
alertService.sendAlert("数据对账发现不一致: " + inconsistentOrders.size() + "个");
// 自动修复
fixInconsistentOrders(inconsistentOrders);
}
}
private void fixInconsistentOrders(List<String> inconsistentOrders) {
for (String orderNo : inconsistentOrders) {
// 从主库同步到从库
Order masterOrder = masterOrderMapper.selectByOrderNo(orderNo);
if (masterOrder != null) {
slaveOrderMapper.replaceOrder(masterOrder);
log.info("已修复不一致订单: {}", orderNo);
}
}
}
}
4.3 优化效果
经过上述优化,我们的系统在大促期间的表现明显改善:
| 指标 | 优化前 | 优化后 | 改善幅度 |
|---|---|---|---|
| 主从同步延迟 | 30秒 | 秒 | 96% |
| 订单查询成功率 | 95% | 99.99% | 4.9% |
| 用户投诉量 | 1000+ | <10 | 99% |
| 数据一致性 | 95% | 99.99% | 4.9% |
五、常见问题与最佳实践
5.1 常见问题
问题一:如何判断主从同步延迟是否严重?
-- 检查同步延迟
SHOW SLAVE STATUS\G
-- 关注以下指标:
-- Seconds_Behind_Master > 5秒:需要关注
-- Seconds_Behind_Master > 30秒:严重延迟
-- Relay_Log_Space > 1GB:可能存在大事务
问题二:如何优化主从同步性能?
-- 优化主库配置
[mysqld]
# 增加binlog缓存
binlog_cache_size = 4M
max_binlog_cache_size = 512M
# 优化刷盘策略
sync_binlog = 1 # 每次事务提交都刷盘(安全性最高)
# 或者 sync_binlog = 0 # 由操作系统控制刷盘(性能最高)
# 推荐使用 sync_binlog = 100 # 折中方案
# 优化从库配置
[mysqld]
# 多线程复制(MySQL 5.7+)
slave_parallel_type = LOGICAL_CLOCK
slave_parallel_workers = 4
# 优化中继日志
relay_log_recovery = 1
问题三:如何监控数据一致性?
// 实时监控数据一致性
@Component
public class DataConsistencyMonitor {
@Autowired
private OrderMapper masterOrderMapper;
@Autowired
private OrderMapper slaveOrderMapper;
// 每隔10秒检查一次最近1分钟的订单
@Scheduled(fixedDelay = 10000)
public void monitorConsistency() {
long startTime = System.currentTimeMillis() - 60 * 1000;
List<Order> recentOrders = masterOrderMapper.selectByTimeRange(startTime, System.currentTimeMillis());
for (Order order : recentOrders) {
Order slaveOrder = slaveOrderMapper.selectByOrderNoWithSlave(order.getOrderNo());
if (slaveOrder == null) {
// 记录延迟日志
log.warn("订单未同步到从库: orderNo={}, masterTime={}",
order.getOrderNo(), order.getCreateTime());
} else if (!Objects.equals(order.getStatus(), slaveOrder.getStatus())) {
// 状态不一致
log.error("订单状态不一致: orderNo={}, masterStatus={}, slaveStatus={}",
order.getOrderNo(), order.getStatus(), slaveOrder.getStatus());
}
}
}
}
5.2 最佳实践
实践一:敏感操作必须读主库
对于订单、支付等敏感操作,建议直接读取主库:
// 标记需要读主库的方法
@ReadFromMaster
public Order queryOrder(String orderNo) {
return orderMapper.selectByOrderNoWithMaster(orderNo);
}
// 普通查询可以使用从库
@ReadFromSlave
public List<Order> queryOrderList(Integer userId, int page, int size) {
return orderMapper.selectByUserIdWithSlave(userId, page, size);
}
实践二:设置合理的超时和重试策略
// 配置化超时和重试
@Configuration
public class OrderConfig {
@Value("${order.query.slave.timeout:3000}")
private long slaveQueryTimeout;
@Value("${order.query.slave.max-retries:3}")
private int maxRetries;
@Bean
public OrderService orderService() {
return new OrderServiceImpl(slaveQueryTimeout, maxRetries);
}
}
public class OrderServiceImpl implements OrderService {
private final long slaveQueryTimeout;
private final int maxRetries;
@Override
public Order queryOrder(String orderNo) {
long startTime = System.currentTimeMillis();
for (int i = 0; i < maxRetries; i++) {
// 检查是否超时
if (System.currentTimeMillis() - startTime > slaveQueryTimeout) {
// 超时,回退到主库
return queryFromMaster(orderNo);
}
Order slaveOrder = queryFromSlave(orderNo);
if (slaveOrder != null) {
return slaveOrder;
}
// 指数退避重试
try {
Thread.sleep((long) (Math.pow(2, i) * 100));
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
break;
}
}
// 最终回退到主库
return queryFromMaster(orderNo);
}
}
实践三:建立数据对账机制
// 多层次对账机制
@Component
public class MultiLevelReconciliation {
// 实时对账:每分钟检查最近1分钟的订单
@Scheduled(fixedDelay = 60000)
public void realTimeReconcile() {
long startTime = System.currentTimeMillis() - 60 * 1000;
checkOrders(startTime, System.currentTimeMillis());
}
// 定期对账:每小时检查最近1小时的订单
@Scheduled(cron = "0 0 * * * ?")
public void hourlyReconcile() {
long startTime = System.currentTimeMillis() - 60 * 60 * 1000;
checkOrders(startTime, System.currentTimeMillis());
}
// 每日对账:检查昨天的所有订单
@Scheduled(cron = "0 30 2 * * ?")
public void dailyReconcile() {
LocalDate yesterday = LocalDate.now().minusDays(1);
ZonedDateTime start = yesterday.atStartOfDay(ZoneId.systemDefault());
ZonedDateTime end = start.plusDays(1);
checkOrders(start.toInstant().toEpochMilli(), end.toInstant().toEpochMilli());
}
private void checkOrders(long startTime, long endTime) {
List<Order> masterOrders = orderMapper.selectByTimeRange(startTime, endTime);
List<Order> slaveOrders = orderMapper.selectByTimeRangeFromSlave(startTime, endTime);
Map<String, Order> slaveOrderMap = slaveOrders.stream()
.collect(Collectors.toMap(Order::getOrderNo, o -> o));
List<String> issues = new ArrayList<>();
for (Order masterOrder : masterOrders) {
Order slaveOrder = slaveOrderMap.get(masterOrder.getOrderNo());
if (slaveOrder == null) {
issues.add("MISSING:" + masterOrder.getOrderNo());
} else if (!isConsistent(masterOrder, slaveOrder)) {
issues.add("INCONSISTENT:" + masterOrder.getOrderNo());
}
}
if (!issues.isEmpty()) {
log.error("发现{}个不一致订单: {}", issues.size(), issues);
alertService.sendAlert("数据不一致告警: " + issues.size() + "个");
// 自动修复
fixInconsistencies(issues);
}
}
}
六、总结
主从同步延迟导致订单数据丢失,是电商系统常见的痛点问题。解决这个问题需要从多个层面入手:
- 理解问题根源:主从同步是异步的,存在延迟是必然的
- 优化锁机制:根据场景选择合适的锁策略(悲观锁、乐观锁、TCC等)
- 加强监控告警:实时监控同步延迟和数据一致性
- 建立对账机制:定期检查和修复数据不一致问题
- 设计降级方案:从库不可用时,能够自动降级到主库
最重要的是,不要期望完美的实时一致性,而是要设计合理的降级和修复机制,确保系统在极端情况下也能保持稳定运行。
希望这篇文章能帮到大家。如果有任何问题,欢迎在评论区交流!
