电商大促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、内存等关键指标,出现问题时能快速定位。

第四,降级预案必不可少。 大促期间,有些功能是可以降级的。比如订单详情中的推荐商品、用户评价等,这些非核心功能在流量高峰时可以考虑降级,把资源留给核心交易链路。

第五,数据一致性要分场景对待。 不是所有场景都需要强一致性。像库存扣减、订单创建这种核心场景,需要保证一致性;但像订单统计、用户行为分析这类场景,最终一致性就够了,这样可以大幅降低系统复杂度。

这套方案经过多次大促的检验,已经比较成熟了。但高并发领域没有银弹,每次大促都会遇到新的问题,需要持续优化和改进。希望这些经验能给你一些参考,如果你们也在面临类似的高并发挑战,欢迎一起交流探讨。