电商大促MySQL高并发实战:锁优化与分库分表方案解析
记得第一次经历双11大促的那年,凌晨零点刚过,我们的监控大屏上流量曲线直冲云霄,MySQL的慢查询日志开始疯狂刷屏,连接数瞬间飙升到上限,运维团队紧急拉起的会议上,空气凝固得能拧出水来。那时候我才真正意识到,平时跑得好好的系统,在面对十倍甚至百倍于日常的并发时,会有多么不堪一击。
今天想跟你聊聊的,就是这套在高并发场景下,我们是如何一步步把MySQL的锁问题和分库分表方案啃下来的实战经历。
高并发场景下的锁竞争:一个真实的崩溃现场
大促期间最常见的性能杀手,就是锁竞争。我们先来看一个真实发生过的场景。
某次大促活动中,订单表的核心热点集中在两个操作:库存扣减和订单创建。业务逻辑看似简单——用户下单时先检查库存,库存够就扣减,然后插入订单记录。但在万级并发下,问题就来了。
当时我们的数据库连接池配置是200,高并发瞬间涌入时,连接数很快被打满。通过SHOW PROCESSLIST观察,发现大量请求都卡在LOCK WAIT状态。进一步用information_schema.innodb_locks查询,定位到问题根源:库存表的主键更新语句UPDATE stock SET count = count - #{num} WHERE product_id = #{productId}触发了行锁竞争,多个并发请求同时更新同一行或相邻行时,InnoDB的间隙锁机制导致锁等待队列越来越长,最终形成雪崩效应。
-- 问题代码示例:粗暴的行锁更新
UPDATE stock
SET count = count - #{purchaseNum},
update_time = NOW()
WHERE product_id = #{productId}
AND count >= #{purchaseNum}
这段代码看起来天衣无缝,逻辑也没问题,但在高并发下,每笔订单都试图获取同一行的行锁,请求排队等待,响应时间直线上升,数据库CPU被打满,甚至出现死锁。
锁优化的第一个思路:减少锁的粒度
我们首先想到的是缩小锁的范围。传统的行锁更新,锁住的是整行数据,并发高的时候竞争就激烈。于是我们把库存扣减逻辑做了拆分:
// 优化思路:使用分布式锁 + 乐观锁 双重保障
public class StockService {
// 1. 先通过乐观锁快速判断,减少行锁持有时间
public boolean tryReduceStock(Long productId, int num) {
// 乐观锁版本号机制
int affected = stockMapper.updateStockByVersion(
productId, num, System.currentTimeMillis()
);
return affected > 0;
}
}
-- 使用乐观锁替代悲观行锁
UPDATE stock
SET count = count - #{num},
version = version + 1,
update_time = NOW()
WHERE product_id = #{productId}
AND count >= #{num}
AND version = #{currentVersion}
乐观锁的思路是:更新时带上当前的版本号,如果数据库里的版本号和查询时一致,才执行更新。这样就不需要长时间持有行锁,大大减少了锁等待时间。当然,乐观锁在高冲突场景下会有大量失败重试,所以我们需要配合业务层的重试机制。
更进一步的优化:库存预扣减与异步最终一致
对于大促这种极端场景,光靠锁的优化还不够。我们后来采用了一套库存预扣减 + 异步最终一致的方案:
// 库存预扣减方案:先扣Redis,再异步同步MySQL
@Service
public class PreStockService {
@Autowired
private StringRedisTemplate redisTemplate;
@Autowired
private MQProducer mqProducer;
// 扣减Redis中的预扣减库存
public boolean preReduceStock(String productId, int num) {
String key = "stock:pre:" + productId;
Long current = redisTemplate.opsForValue().decrement(key, num);
if (current < 0) {
redisTemplate.opsForValue().increment(key, num);
return false;
}
// 扣减成功,发送MQ消息异步同步到MySQL
mqProducer.send("stock.sync", buildStockMessage(productId, num));
return true;
}
// 异步同步到MySQL,批量操作减少数据库压力
@KafkaListener(topics = "stock.sync")
public void syncStockToDB(StockMessage message) {
batchUpdateStock(message.getProductId(), message.getNum());
}
}
这套方案的核心思想是:把高并发的压力从MySQL转移到Redis,Redis的单线程模型处理并发请求的性能远超MySQL,然后通过MQ异步把数据同步到MySQL,实现最终一致。
分库分表:打破单表性能瓶颈的必由之路
锁优化只是治标,当数据量持续增长,单表百万、千万级的时候,查询性能本身就会成为瓶颈。这时候分库分表就不得不提了。
我们当时面临的问题是:订单表随着业务增长,单表数据量突破了5000万,即使加了索引,复杂查询和全表扫描的效率依然很低。更严重的是,写入集中在几个热点分区,导致数据倾斜。
分库分表的核心设计:路由键的选择
分库分表最关键的一步是选择分片键(Sharding Key)。选错了,后续查询会非常痛苦。
我们的订单表有几个常见的查询场景:
- 按订单ID查询订单详情
- 按用户ID查询用户的订单列表
- 按商家ID查询商家的订单
- 按时间范围查询订单统计
经过分析,我们发现按order_id查询是最频繁的场景,而order_id本身是可以生成具有规律性的雪花ID,于是我们选择了订单ID作为分片键。
// 雪花ID生成器,保证ID全局唯一且具有时间有序性
public class SnowflakeIdGenerator {
private long workerId;
private long datacenterId;
private long sequence = 0L;
private long lastTimestamp = -1L;
public synchronized long nextId() {
long timestamp = currentTimeMillis();
if (timestamp < lastTimestamp) {
throw new RuntimeException("时钟回拨异常");
}
if (timestamp == lastTimestamp) {
sequence = (sequence + 1) & 4095;
if (sequence == 0) {
timestamp = waitNextMillis(lastTimestamp);
}
} else {
sequence = 0L;
}
lastTimestamp = timestamp;
return ((timestamp - EPOCH) << 22)
| (datacenterId << 17)
| (workerId << 12)
| sequence;
}
}
分库分表的路由算法:
// 分库分表策略:基于order_id取模
public class OrderShardingStrategy implements ShardingStrategy<Long> {
private static final int DB_COUNT = 16;
private static final int TABLE_COUNT_PER_DB = 64;
@Override
public int routeDb(Long shardingKey) {
return Math.abs(shardingKey % DB_COUNT);
}
@Override
public int routeTable(Long shardingKey, int dbIndex) {
return Math.abs((shardingKey / DB_COUNT) % TABLE_COUNT_PER_DB);
}
// 生成实际的数据源和表名
public String getDataSourceName(Long orderId) {
int dbIndex = routeDb(orderId);
return "order_db_" + dbIndex;
}
public String getTableName(Long orderId) {
int dbIndex = routeDb(orderId);
int tableIndex = routeTable(orderId, dbIndex);
return "t_order_" + tableIndex;
}
}
分库分表后的查询难题:非分片键查询怎么办?
这是分库分表后最头疼的问题。比如我们要按用户ID查询订单,但分片键是订单ID,这就意味着需要遍历所有分片才能查到结果。
我们的解决方案是建立反向索引表:
-- 用户维度的反向索引表,用于按userId查询
CREATE TABLE t_user_order_index (
id BIGINT PRIMARY KEY,
user_id BIGINT NOT NULL,
order_id BIGINT NOT NULL,
create_time DATETIME,
INDEX idx_user_id (user_id, create_time)
);
-- 商家维度的反向索引表
CREATE TABLE t_merchant_order_index (
id BIGINT PRIMARY KEY,
merchant_id BIGINT NOT NULL,
order_id BIGINT NOT NULL,
create_time DATETIME,
INDEX idx_merchant_id (merchant_id)
);
这套索引表和我们主订单表通过双写+消息队列保持同步:
@Service
public class OrderWriteService {
@Autowired
private OrderMapper orderMapper;
@Autowired
private OrderIndexMapper indexMapper;
@Autowired
private MQProducer mqProducer;
public Long createOrder(OrderCreateRequest request) {
Long orderId = idGenerator.nextId();
// 1. 写入主订单表
Order order = buildOrder(orderId, request);
orderMapper.insert(order);
// 2. 写入反向索引表(同步,保证查询一致性)
OrderIndex index = buildOrderIndex(orderId, request.getUserId(),
request.getMerchantId());
indexMapper.insert(index);
// 3. 发送MQ消息(异步,用于数据补偿和最终一致性保障)
mqProducer.send("order.create", buildOrderMessage(orderId, request));
return orderId;
}
}
分页查询的优化
分库分表后的分页查询也是一个经典难题。传统的LIMIT offset, size在跨分片场景下性能极差。
// 基于游标的分页方案,避免深分页
public class CursorBasedPaging {
/**
* 游标分页查询,适合分库分表场景
* @param lastOrderId 上一页最后一条记录的orderId
* @param limit 每页大小
*/
public List<Order> queryWithCursor(Long lastOrderId, int limit) {
// 根据lastOrderId确定从哪个分片开始查询
int startDb = routeDb(lastOrderId);
int startTable = routeTable(lastOrderId, startDb);
List<Order> result = new ArrayList<>();
// 从当前分片开始遍历,直到收集够limit条记录
for (int db = startDb; db < DB_COUNT; db++) {
for (int table = (db == startDb) ? startTable : 0;
table < TABLE_COUNT_PER_DB; table++) {
String sql = "SELECT * FROM t_order_" + table
+ " WHERE order_id > ? ORDER BY order_id ASC LIMIT ?";
List<Order> batch = orderMapper.queryByCursor(sql, lastOrderId, limit);
result.addAll(batch);
if (result.size() >= limit) {
return result.subList(0, limit);
}
// 更新游标
if (!batch.isEmpty()) {
lastOrderId = batch.get(batch.size() - 1).getOrderId();
}
}
}
return result;
}
}
分布式事务:分库分表后的又一挑战
分库分表之后,跨分片的分布式事务成了新的问题。比如用户下单时,需要同时操作订单库、库存库、优惠券库,这在单个数据库内可以通过本地事务解决,但在分库分表后就需要分布式事务了。
我们最终采用的是TCC(Try-Confirm-Cancel)模式,配合本地消息表实现最终一致性:
/**
* TCC分布式事务示例:下单操作
*/
@Component
public class OrderTccTransaction {
@Autowired
private StockService stockService;
@Autowired
private CouponService couponService;
@Autowired
private OrderMapper orderMapper;
@Autowired
private LocalMessageMapper localMessageMapper;
/**
* Try阶段:预留资源
*/
public String tryCreateOrder(OrderRequest request) {
// 1. 尝试扣减库存(预留,不真正扣减)
boolean stockReserved = stockService.tryReserveStock(
request.getProductId(), request.getNum()
);
if (!stockReserved) {
throw new BusinessException("库存不足");
}
// 2. 尝试使用优惠券
boolean couponApplied = couponService.tryApplyCoupon(
request.getUserId(), request.getCouponId()
);
if (!couponApplied) {
stockService.cancelReserveStock(
request.getProductId(), request.getNum()
);
throw new BusinessException("优惠券不可用");
}
// 3. 创建订单(预下单状态)
Order order = buildPreOrder(request);
order.setStatus(OrderStatus.PENDING);
orderMapper.insert(order);
// 4. 记录本地消息,用于Confirm阶段
LocalMessage message = buildLocalMessage("order.confirm", order.getOrderId());
localMessageMapper.insert(message);
return order.getOrderId().toString();
}
/**
* Confirm阶段:确认提交
*/
@Transactional
public void confirmOrder(String orderId) {
// 1. 真正扣减库存
stockService.confirmReduceStock(orderId);
// 2. 更新订单状态为已下单
orderMapper.updateStatus(orderId, OrderStatus.CONFIRMED);
// 3. 标记本地消息已处理
localMessageMapper.markProcessed(orderId);
}
/**
* Cancel阶段:回滚
*/
@Transactional
public void cancelOrder(String orderId) {
// 1. 取消库存预留
stockService.cancelReserveStockByOrderId(orderId);
// 2. 取消优惠券
couponService.cancelApplyCoupon(orderId);
// 3. 删除预下单
orderMapper.deletePreOrder(orderId);
// 4. 标记本地消息已取消
localMessageMapper.markCancelled(orderId);
}
}
本地消息表的实现也很关键,它是保证最终一致性的核心:
-- 本地消息表,用于异步保证分布式事务的最终一致性
CREATE TABLE t_local_message (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
message_id VARCHAR(64) NOT NULL UNIQUE,
topic VARCHAR(128) NOT NULL,
content TEXT NOT NULL,
status TINYINT NOT NULL DEFAULT 0 COMMENT '0-待发送 1-已发送 2-已处理',
retry_count INT DEFAULT 0,
create_time DATETIME DEFAULT CURRENT_TIMESTAMP,
update_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
INDEX idx_status (status, create_time)
);
// 本地消息发送补偿任务
@Component
public class LocalMessageSender {
@Scheduled(fixedRate = 5000)
public void sendPendingMessages() {
// 查询待发送的本地消息
List<LocalMessage> pendingMessages = localMessageMapper
.selectByStatus(LocalMessageStatus.PENDING.getCode(), 100);
for (LocalMessage message : pendingMessages) {
try {
mqProducer.send(message.getTopic(), message.getContent());
localMessageMapper.markSent(message.getId());
} catch (Exception e) {
// 发送失败,增加重试次数
localMessageMapper.incrementRetryCount(message.getId());
}
}
}
}
监控与弹性扩容:大促期间的持续保障
方案再好,也需要在大促期间持续监控和动态调整。我们建立了一套完整的监控体系:
# Prometheus监控配置示例
scrape_configs:
- job_name: 'mysql'
static_configs:
- targets: ['mysql-exporter:9104']
metrics_path: /metrics
- job_name: 'application'
static_configs:
- targets: ['app-server:8080']
// 应用层面的性能监控
@Component
public class PerformanceMonitor {
@Autowired
private MeterRegistry meterRegistry;
/**
* 监控SQL执行耗时
*/
public <T> T monitorExecution(String metricName, Callable<T> task) {
Timer.Sample sample = Timer.start(meterRegistry);
try {
T result = task.call();
sample.stop(Timer.builder(metricName)
.tag("status", "success")
.register(meterRegistry));
return result;
} catch (Exception e) {
sample.stop(Timer.builder(metricName)
.tag("status", "error")
.register(meterRegistry));
throw e;
}
}
/**
* 实时监控数据库连接池状态
*/
public void reportDataSourceStatus(DataSource dataSource) {
HikariDataSource hikari = (HikariDataSource) dataSource;
Gauge.builder("db.pool.active", hikari,
ds -> ds.getHikariPoolMXBean().getActiveConnections())
.tag("pool", "order-db")
.register(meterRegistry);
Gauge.builder("db.pool.idle", hikari,
ds -> ds.getHikariPoolMXBean().getIdleConnections())
.tag("pool", "order-db")
.register(meterRegistry);
}
}
大促期间最关键的还有动态扩容能力。我们采用了数据库中间件配合自动扩缩容的方案:
/**
* 动态分片配置,支持运行时调整
*/
@Configuration
public class ShardingDynamicConfig {
@Autowired
private ShardingSphereDataSource shardingDataSource;
/**
* 大促期间动态增加分片
*/
public void scaleOut(int newDbCount, int newTableCountPerDb) {
// 重新计算分片规则
ShardingRuleConfiguration newRule = buildNewShardingRule(
newDbCount, newTableCountPerDb
);
// 动态更新分片配置
shardingDataSource.getConfigurationManager()
.updateRuleConfiguration(newRule);
// 触发数据迁移(增量数据通过CDC同步)
dataMigrationService.startIncrementalMigration(
newDbCount, newTableCountPerDb
);
}
/**
* 大促结束后缩容
*/
public void scaleIn(int targetDbCount, int targetTableCountPerDb) {
// 1. 停止新数据写入到扩容的分片
// 2. 等待CDC同步完成
// 3. 合并分片
// 4. 更新分片规则
}
}
压测验证:用数据说话
方案上线前,压测是必不可少的环节。我们使用了专业的压测工具对优化后的系统进行全面测试:
/**
* 高并发压测场景模拟
*/
@SpringBootTest
class StressTest {
@Autowired
private OrderService orderService;
@Test
void testOrderCreationUnderHighConcurrency() throws Exception {
int threadCount = 500;
int requestCount = 10000;
// 使用CountDownLatch同步并发
CountDownLatch latch = new CountDownLatch(threadCount);
AtomicInteger successCount = new AtomicInteger(0);
AtomicInteger failCount = new AtomicInteger(0);
long startTime = System.currentTimeMillis();
// 启动并发线程
for (int i = 0; i < threadCount; i++) {
final int threadId = i;
new Thread(() -> {
try {
for (int j = 0; j < requestCount / threadCount; j++) {
try {
Long orderId = orderService.createOrder(buildOrderRequest(threadId, j));
if (orderId != null) {
successCount.incrementAndGet();
} else {
failCount.incrementAndGet();
}
} catch (Exception e) {
failCount.incrementAndGet();
}
}
} finally {
latch.countDown();
}
}).start();
}
latch.await();
long endTime = System.currentTimeMillis();
System.out.println("========================================");
System.out.println("压测结果:");
System.out.println("总请求数: " + requestCount);
System.out.println("成功数: " + successCount.get());
System.out.println("失败数: " + failCount.get());
System.out.println("成功率: " + (successCount.get() * 100.0 / requestCount) + "%");
System.out.println("总耗时: " + (endTime - startTime) + "ms");
System.out.println("QPS: " + (requestCount * 1000.0 / (endTime - startTime)));
System.out.println("平均响应时间: " + (endTime - startTime) * 1.0 / requestCount + "ms");
System.out.println("========================================");
}
}
压测的结果让我们对系统能力有了清晰的认识:优化后,订单创建的吞吐量从原来的每秒500单提升到了每秒3000单以上,平均响应时间从200ms降低到了50ms以内,数据库连接数峰值从800降到了150左右。
一些实战中的经验教训
回过头来看这段经历,有几个点特别值得分享:
第一,不要等到问题出现了再优化。 很多团队都是大促当天被流量打趴下了才开始想办法,这时候再优化已经来不及了。应该在平时就做充分的压测,预留足够的性能余量。
第二,分库分表要尽早规划。 我们当时是数据量上来了才做分库分表,迁移过程非常痛苦,还出了不少线上问题。如果一开始就做好分片设计,后续会顺畅很多。
第三,监控比优化更重要。 有了完善的监控,才能快速定位问题。我们后来建立了从应用层到数据库层的完整监控链路,包括慢查询、锁等待、连接池、CPU、内存等关键指标,出现问题时能快速定位。
第四,降级预案必不可少。 大促期间,有些功能是可以降级的。比如订单详情中的推荐商品、用户评价等,这些非核心功能在流量高峰时可以考虑降级,把资源留给核心交易链路。
第五,数据一致性要分场景对待。 不是所有场景都需要强一致性。像库存扣减、订单创建这种核心场景,需要保证一致性;但像订单统计、用户行为分析这类场景,最终一致性就够了,这样可以大幅降低系统复杂度。
这套方案经过多次大促的检验,已经比较成熟了。但高并发领域没有银弹,每次大促都会遇到新的问题,需要持续优化和改进。希望这些经验能给你一些参考,如果你们也在面临类似的高并发挑战,欢迎一起交流探讨。
