某电商大促MySQL宕机后 架构师总结高并发优化方案 连接池读写分离分库分表实战经验分享
一、那晚发生了什么
说出来你可能不信,那天是某个电商平台的大促首日,下午三点,GMV正在爬坡。结果监控大屏突然红灯一片——MySQL主库连接数打满,响应时间飙到30秒,整个订单系统直接瘫痪。
事后复盘,架构组把事故原因总结成一句话:以为业务量增长是线性的,但实际上高并发下的资源竞争是非线性的。
下面这篇文章,是我参与那次事故复盘后,带着团队重新梳理高并发优化方案的心路历程。没有空洞的理论,全是踩过的坑和真正落地的代码。
二、连接池:最先被忽视的”小水管”
2.1 事故前的配置有多简陋
事故发生前,我们用的是Spring Boot默认的HikariCP配置:
spring:
datasource:
url: jdbc:mysql://192.168.1.100:3306/ecommerce
username: root
password: "123456"
hikari:
# 默认连接池配置,maxPoolSize = 10,minIdle = 10
maximum-pool-size: 10
minimum-idle: 10
idle-timeout: 600000
connection-timeout: 30000
看起来没问题对吧? 10个连接,应付日常流量绰绰有余。但大促时,每秒几千笔订单同时涌入,每个请求都需要一个数据库连接,10个连接瞬间被打满,后续请求全部排队等待,最后直接OOM。
2.2 我们是如何重新设计连接池的
第一步:动态计算合理的连接数
连接池大小不是越大越好,也不是越小越好。有一个经验公式:
最佳连接数 = CPU核数 × 2 + 磁盘数
我们的服务器是8核16G,磁盘是2个SSD,所以:
最佳连接数 = 8 × 2 + 2 = 18
但考虑到大促期间的峰值流量,我们决定给到 25 个连接。
第二步:完整的HikariCP优化配置
spring:
datasource:
url: jdbc:mysql://192.168.1.100:3306/ecommerce
username: root
password: "123456"
driver-class-name: com.mysql.cj.jdbc.Driver
hikari:
# 连接池核心参数
maximum-pool-size: 25 # 最大连接数
minimum-idle: 10 # 最小空闲连接数
idle-timeout: 600000 # 空闲连接超时时间(10分钟)
max-lifetime: 1800000 # 连接最大生命周期(30分钟)
connection-timeout: 5000 # 获取连接超时时间(5秒)
leak-detection-threshold: 60000 # 连接泄漏检测阈值(60秒)
# 连接测试配置(确保连接有效)
connection-test-query: SELECT 1
validation-timeout: 3000
第三步:用Java代码动态管理连接池
光靠配置还不够,我们需要在代码层面做更精细的控制:
@Component
public class DynamicDataSourceConfig {
@Autowired
private HikariDataSource dataSource;
/**
* 根据当前压力动态调整连接池大小
*/
@Scheduled(fixedRate = 30000) // 每30秒检查一次
public void adjustPoolSize() {
int currentPoolSize = dataSource.getHikariPoolMXBean().getActiveConnections();
int maxSize = dataSource.getMaximumPoolSize();
// 当活跃连接数超过80%时,适当增加连接
if (currentPoolSize > maxSize * 0.8) {
dataSource.setMaximumPoolSize(maxSize + 5);
log.warn("检测到高负载,连接池从{}扩大到{}",
maxSize - 5, maxSize + 5);
}
// 当活跃连接数低于20%时,适当缩减
else if (currentPoolSize < maxSize * 0.2) {
dataSource.setMaximumPoolSize(Math.max(maxSize - 5, 10));
log.info("负载降低,连接池从{}缩减到{}",
maxSize + 5, Math.max(maxSize - 5, 10));
}
}
/**
* 连接池健康监控
*/
@Scheduled(fixedRate = 10000)
public void monitorPoolHealth() {
HikariPoolMXBean poolBean = dataSource.getHikariPoolMXBean();
log.info("连接池状态 - " +
"活跃连接: {}, 空闲连接: {}, 总连接数: {}, 等待连接数: {}",
poolBean.getActiveConnections(),
poolBean.getIdleConnections(),
poolBean.getTotalConnections(),
poolBean.getThreadsAwaitingConnection());
// 当等待连接数超过阈值时,发送告警
if (poolBean.getThreadsAwaitingConnection() > 10) {
sendAlert("连接池压力过大,等待连接的线程数:"
+ poolBean.getThreadsAwaitingConnection());
}
}
}
第四步:处理连接泄漏问题
连接泄漏是连接池最常见的坑之一。比如下面这种代码:
// 错误示范:忘记关闭连接
Connection conn = dataSource.getConnection();
try {
// 处理业务...
// 如果这里抛异常,连接永远不会被释放!
} catch (Exception e) {
e.printStackTrace();
}
// conn没有关闭!
正确写法:
// 使用try-with-resources自动管理连接
try (Connection conn = dataSource.getConnection();
PreparedStatement stmt = conn.prepareStatement("SELECT * FROM orders")) {
try (ResultSet rs = stmt.executeQuery()) {
while (rs.next()) {
// 处理结果...
}
}
} catch (SQLException e) {
log.error("数据库查询失败", e);
}
// 所有资源自动关闭
三、读写分离:把压力分散到多个节点
3.1 为什么需要读写分离
MySQL单节点的瓶颈在于:写操作会锁表,读操作会阻塞写操作。
大促期间,订单写入量激增,同时用户查询订单状态、商品详情等读操作也在疯狂增加。如果所有请求都打到主库,主库会同时承担读写压力,最终不堪重负。
读写分离的核心思路:主库只负责写,从库负责读,通过主从复制保证数据一致性。
3.2 完整的读写分离架构
┌─────────────┐
│ 应用层 │
│ (Spring Boot)│
└──────┬──────┘
│
┌─────────────────┼─────────────────┐
│ │ │
┌─────▼─────┐ ┌─────▼─────┐ ┌─────▼─────┐
│ 主库MySQL │ │ 从库MySQL │ │ 从库MySQL │
│ (写操作) │ │ (读操作) │ │ (读操作) │
└─────┬─────┘ └───────────┘ └───────────┘
│
┌─────▼─────┐
│ MySQL主从 │
│ 复制 │
└───────────┘
3.3 基于Spring的动态路由数据源
@Configuration
public class DataSourceConfig {
/**
* 主数据源 - 负责写操作
*/
@Bean
@Primary
@ConfigurationProperties("spring.datasource.master")
public DataSource masterDataSource() {
return DataSourceBuilder.create().type(HikariDataSource.class).build();
}
/**
* 从数据源 - 负责读操作(可以有多个)
*/
@Bean
@ConfigurationProperties("spring.datasource.slave1")
public DataSource slave1DataSource() {
return DataSourceBuilder.create().type(HikariDataSource.class).build();
}
/**
* 从数据源2
*/
@Bean
@ConfigurationProperties("spring.datasource.slave2")
public DataSource slave2DataSource() {
return DataSourceBuilder.create().type(HikariDataSource.class).build();
}
/**
* 动态数据源路由
*/
@Bean
public DataSource routingDataSource(
DataSource masterDataSource,
DataSource slave1DataSource,
DataSource slave2DataSource) {
Map<Object, Object> targetDataSources = new HashMap<>();
targetDataSources.put(DataSourceType.MASTER, masterDataSource);
targetDataSources.put(DataSourceType.SLAVE, slave1DataSource);
targetDataSources.put(DataSourceType.SLAVE, slave2DataSource); // 同一类型会合并
RoutingDataSource routingDataSource = new RoutingDataSource();
routingDataSource.setTargetDataSources(targetDataSources);
routingDataSource.setDefaultTargetDataSource(masterDataSource);
return routingDataSource;
}
/**
* 基于注解的数据源选择器
*/
public static class RoutingDataSource extends AbstractRoutingDataSource {
@Override
protected Object determineCurrentLookupKey() {
return DataSourceContextHolder.getDataSourceType();
}
}
/**
* 数据源类型枚举
*/
public enum DataSourceType {
MASTER, SLAVE
}
}
3.4 数据源切换上下文
public class DataSourceContextHolder {
private static final ThreadLocal<DataSourceType> CONTEXT = new ThreadLocal<>();
public static void setDataSourceType(DataSourceType type) {
CONTEXT.set(type);
}
public static DataSourceType getDataSourceType() {
return CONTEXT.get();
}
public static void clearDataSourceType() {
CONTEXT.remove();
}
}
3.5 AOP自动切换数据源
@Aspect
@Component
public class DataSourceAop {
/**
* 拦截所有带有@Master注解的方法,使用主库
*/
@Before("@annotation(master)")
public void setMasterDataSource(JoinPoint point, Master master) {
DataSourceContextHolder.setDataSourceType(DataSourceType.MASTER);
}
/**
* 拦截所有Service方法,默认使用从库(读操作)
*/
@Pointcut("execution(* com.ecommerce.service.*.*(..))")
public void serviceLayer() {}
@Around("serviceLayer()")
public Object aroundService(ProceedingJoinPoint point) throws Throwable {
try {
// 判断是否是写操作(通过方法名判断)
String methodName = point.getSignature().getName();
if (isWriteOperation(methodName)) {
DataSourceContextHolder.setDataSourceType(DataSourceType.MASTER);
} else {
DataSourceContextHolder.setDataSourceType(DataSourceType.SLAVE);
}
return point.proceed();
} finally {
DataSourceContextHolder.clearDataSourceType();
}
}
private boolean isWriteOperation(String methodName) {
// 常见的写操作关键词
String[] writeKeywords = {"insert", "update", "delete", "add", "create", "save", "remove"};
return Arrays.stream(writeKeywords)
.anyMatch(keyword -> methodName.toLowerCase().startsWith(keyword));
}
}
3.6 多从库负载均衡
当有多个从库时,如何实现负载均衡?我们可以自己实现一个简单的轮询策略:
@Component
public class LoadBalanceDataSource extends AbstractRoutingDataSource {
private final List<DataSource> slaveDataSources;
private final AtomicInteger counter = new AtomicInteger(0);
public LoadBalanceDataSource(DataSource masterDataSource, List<DataSource> slaveDataSources) {
this.slaveDataSources = slaveDataSources;
Map<Object, Object> targetDataSources = new HashMap<>();
targetDataSources.put(DataSourceType.MASTER, masterDataSource);
targetDataSources.put(DataSourceType.SLAVE, createLoadBalanceWrapper());
setTargetDataSources(targetDataSources);
afterPropertiesSet();
}
/**
* 创建负载均衡包装器
*/
private DataSource createLoadBalanceWrapper() {
return new LoadBalanceDataSourceWrapper(slaveDataSources);
}
@Override
protected Object determineCurrentLookupKey() {
return DataSourceContextHolder.getDataSourceType();
}
}
/**
* 负载均衡数据源包装器
*/
public class LoadBalanceDataSourceWrapper implements DataSource {
private final List<DataSource> dataSources;
private final ReadWriteLock lock = new ReentrantReadWriteLock();
public LoadBalanceDataSourceWrapper(List<DataSource> dataSources) {
this.dataSources = dataSources;
}
@Override
public Connection getConnection() throws SQLException {
// 轮询选择从库
int index = counter.getAndIncrement() % dataSources.size();
if (counter.get() < 0) { // 防止溢出
counter.set(0);
index = 0;
}
return dataSources.get(index).getConnection();
}
// 其他方法代理...
}
四、分库分表:终极武器,但要用对地方
4.1 什么时候需要分库分表?
先说一个误区:不是所有系统都需要分库分表。
我们的经验是:
- 单表数据量超过 500万 行时,查询性能开始明显下降
- 单库QPS超过 5000 时,主库压力开始显现
- 业务增长预期在半年内会导致数据量翻倍
满足以上任一条件,就需要考虑分库分表了。
4.2 我们的分库分表策略
第一步:确定分片键
分片键的选择至关重要。我们的订单表以 user_id 作为分片键:
-- 原始表结构
CREATE TABLE `orders` (
`id` bigint(20) NOT NULL AUTO_INCREMENT,
`user_id` bigint(20) NOT NULL COMMENT '用户ID',
`order_no` varchar(32) NOT NULL COMMENT '订单编号',
`amount` decimal(10,2) NOT NULL COMMENT '订单金额',
`status` tinyint(4) NOT NULL COMMENT '订单状态',
`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_order_no` (`order_no`),
KEY `idx_user_id` (`user_id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='订单表';
第二步:使用ShardingSphere进行分表
ShardingSphere是目前最流行的分库分表框架之一。以下是完整配置:
# application-sharding.yml
spring:
shardingsphere:
# 数据源配置
datasource:
names: ds0,ds1,ds2,ds3
ds0:
type: com.zaxxer.hikari.HikariDataSource
driver-class-name: com.mysql.cj.jdbc.Driver
jdbc-url: jdbc:mysql://192.168.1.100:3306/ecommerce_ds0
username: root
password: "123456"
hikari:
maximum-pool-size: 20
ds1:
type: com.zaxxer.hikari.HikariDataSource
driver-class-name: com.mysql.cj.jdbc.Driver
jdbc-url: jdbc:mysql://192.168.1.101:3306/ecommerce_ds1
username: root
password: "123456"
hikari:
maximum-pool-size: 20
ds2:
type: com.zaxxer.hikari.HikariDataSource
driver-class-name: com.mysql.cj.jdbc.Driver
jdbc-url: jdbc:mysql://192.168.1.102:3306/ecommerce_ds2
username: root
password: "123456"
hikari:
maximum-pool-size: 20
ds3:
type: com.zaxxer.hikari.HikariDataSource
driver-class-name: com.mysql.cj.jdbc.Driver
jdbc-url: jdbc:mysql://192.168.1.103:3306/ecommerce_ds3
username: root
password: "123456"
hikari:
maximum-pool-size: 20
# 分片规则配置
rules:
sharding:
tables:
orders:
actual-data-nodes: ds$->{0..3}.orders_$->{0..3}
# 数据库分片策略
database-strategy:
standard:
sharding-column: user_id
sharding-algorithm-name: database-inline
# 表分片策略
table-strategy:
standard:
sharding-column: user_id
sharding-algorithm-name: table-inline
# 主键生成策略
key-generate-strategy:
column: id
key-generator-name: snowflake
sharding-algorithms:
database-inline:
type: INLINE
props:
algorithm-expression: ds$->{user_id % 4}
table-inline:
type: INLINE
props:
algorithm-expression: orders_$->{user_id % 4}
key-generators:
snowflake:
type: SNOWFLAKE
props:
worker-id: 123 # 每台机器不同
# 其他配置
props:
sql-show: true # 打印SQL,便于调试
max-connections-size-per-query: 1
check-table-does-not-exist: false
第三步:分片键的选择原则
好分片键的标准:
- 高频查询字段:订单查询通常按用户ID查询
- 数据分布均匀:用户ID是数字,取模后分布均匀
- 避免跨分片查询:尽量避免按订单号、时间等字段查询
不好的分片键:
- 用户手机号:涉及隐私,且可能存在分布不均
- 订单金额:数值范围有限,容易聚集
- 创建时间:时间趋势明显,导致数据倾斜
4.3 跨分片查询的处理方案
分库分表后,最头疼的问题是跨分片查询。比如按订单号查询:
-- 这个查询需要全表扫描所有分片,性能极差!
SELECT * FROM orders WHERE order_no = '20241201001';
解决方案:引入ES搜索引擎
@Component
public class OrderEsSyncService {
@Autowired
private ElasticsearchRestTemplate elasticsearchTemplate;
@Autowired
private OrderMapper orderMapper;
/**
* 监听订单变更,同步到ES
*/
@Async
public void syncOrderToEs(Long orderId) {
Order order = orderMapper.selectById(orderId);
if (order != null) {
OrderEsDoc esDoc = convertToEsDoc(order);
elasticsearchTemplate.save(esDoc);
}
}
/**
* 在ES中按订单号查询
*/
public Order searchOrderByOrderNo(String orderNo) {
QueryBuilder queryBuilder = QueryBuilders.termQuery("orderNo", orderNo);
SearchSourceBuilder sourceBuilder = new SearchSourceBuilder();
sourceBuilder.query(queryBuilder).size(1);
NativeQuery nativeQuery = new NativeQuery(sourceBuilder, OrderEsDoc.class);
SearchHits<OrderEsDoc> hits = elasticsearchTemplate.search(nativeQuery);
if (hits.hasSearchHits()) {
OrderEsDoc esDoc = hits.getContent();
return convertToOrder(esDoc);
}
return null;
}
}
4.4 分库分表后的数据迁移
如果现有系统没有分库分表,如何平滑迁移?
方案:双写 + 数据回放
@Component
public class DualWriteService {
@Autowired
private OrderMapper orderMapper; // 原表
@Autowired
private OrderShardingMapper shardingOrderMapper; // 分表后
/**
* 双写:同时写入原表和分表
*/
@Transactional
public void dualWriteOrder(Order order) {
// 写入原表
orderMapper.insert(order);
// 写入分表(根据user_id路由到不同分片)
shardingOrderMapper.insert(order);
}
/**
* 数据回放:将历史数据迁移到分表
*/
public void migrateHistoryData() {
long total = orderMapper.count();
long batchSize = 1000;
long offset = 0;
while (offset < total) {
List<Order> orders = orderMapper.selectListByOffset(offset, batchSize);
// 按user_id路由到不同分片
Map<Long, List<Order>> groupByUserId = orders.stream()
.collect(Collectors.groupingBy(Order::getUserId));
groupByUserId.forEach((userId, orderList) -> {
int shardIndex = (int)(userId % 4);
shardingOrderMapper.batchInsert(shardIndex, orderList);
});
offset += batchSize;
log.info("已迁移 {}/{} 条记录", offset, total);
}
}
}
五、完整的架构图
经过上述优化,我们的高并发架构如下:
┌─────────────────────────────────────────────────────────────────────────┐
│ 客户端请求 │
└───────────────────────────────┬─────────────────────────────────────────┘
│
┌───────────▼───────────┐
│ 负载均衡层 │
│ (Nginx / SLB) │
└───────────┬───────────┘
│
┌─────────────────┼─────────────────┐
│ │ │
┌─────────▼──────┐ ┌────────▼──────┐ ┌────────▼──────┐
│ 应用服务器A │ │ 应用服务器B │ │ 应用服务器C │
│ (Spring Boot) │ │ (Spring Boot) │ │ (Spring Boot) │
└─────────┬──────┘ └────────┬──────┘ └────────┬──────┘
│ │ │
└─────────────────┼─────────────────┘
│
┌───────────▼───────────┐
│ 缓存层 (Redis) │
│ - 热点数据缓存 │
│ - 分布式锁 │
│ - 秒杀库存预扣 │
└───────────┬───────────┘
│
┌─────────────────┼─────────────────┐
│ │ │
┌─────────▼──────┐ ┌────────▼──────┐ ┌────────▼──────┐
│ 主库 (写) │ │ 从库 (读) │ │ 从库 (读) │
│ MySQL Master │ │ MySQL Slave │ │ MySQL Slave │
└─────────┬──────┘ └────────┬──────┘ └────────┬──────┘
│ │ │
│ 主从复制 │ │
│ │ │
┌─────────▼─────────────────▼─────────────────▼──────┐
│ 分库分表层 (ShardingSphere) │
│ ds0: orders_0, orders_1 │
│ ds1: orders_0, orders_1 │
│ ds2: orders_0, orders_1 │
│ ds3: orders_0, orders_1 │
└─────────────────────────┬─────────────────────────┘
│
┌───────────────┼───────────────┐
│ │ │
┌─────────▼──────┐ ┌──────▼──────┐ ┌──────▼──────┐
│ Elasticsearch │ │ MongoDB │ │ ClickHouse │
│ (订单查询) │ │ (日志存储) │ │ (数据分析) │
└─────────────────┘ └─────────────┘ └─────────────┘
六、一些实战中的”血泪教训”
6.1 连接池不是越大越好
我们曾经把连接池设到50,结果反而更慢。原因是:
- 连接太多,服务器内存压力大
- 连接切换的上下文开销增加
- 数据库本身也承受不住这么多连接
建议:连接池大小 = CPU核数 × 2 + 磁盘数,这是一个比较合理的起点。
6.2 读写分离有延迟
MySQL主从复制是异步的,从库数据可能有几秒延迟。这在以下场景会有问题:
// 问题代码:刚写完订单,马上查询,可能查到旧数据
orderService.insert(order);
Order queryResult = orderService.selectById(order.getId());
// queryResult可能是旧数据!
解决方案:
// 方案1:强制走主库
@Master
public Order insertAndGetOrder(Order order) {
orderMapper.insert(order);
return orderMapper.selectById(order.getId()); // 强制从主库查
}
// 方案2:使用Redis缓存保证一致性
@Transactional
public Order insertOrder(Order order) {
orderMapper.insert(order);
redisTemplate.opsForValue().set("order:" + order.getId(), order, 30, TimeUnit.SECONDS);
return order;
}
6.3 分库分表后的ID生成
不要使用MySQL自增ID,因为分库后ID会冲突。我们使用雪花算法:
@Component
public class SnowflakeIdGenerator {
private final AtomicLong workerId;
private final AtomicLong datacenterId;
private final AtomicLong sequence = new AtomicLong(0);
private final long twepoch = 1288834974657L; // Twitter的纪元时间
private final long workerIdBits = 5L;
private final long datacenterIdBits = 5L;
private final long maxWorkerId = -1L ^ (-1L << workerIdBits);
private final long maxDatacenterId = -1L ^ (-1L << datacenterIdBits);
private final long sequenceBits = 12L;
private final long workerIdShift = sequenceBits;
private final long datacenterIdShift = sequenceBits + workerIdBits;
private final long timestampLeftShift = sequenceBits + workerIdBits + datacenterIdBits;
private final long sequenceMask = -1L ^ (-1L << sequenceBits);
private long lastTimestamp = -1L;
public SnowflakeIdGenerator(long workerId, long datacenterId) {
if (workerId > maxWorkerId || workerId < 0) {
throw new IllegalArgumentException(
String.format("worker Id can't be greater than %d or less than 0", maxWorkerId));
}
if (datacenterId > maxDatacenterId || datacenterId < 0) {
throw new IllegalArgumentException(
String.format("datacenter Id can't be greater than %d or less than 0", maxDatacenterId));
}
this.workerId = new AtomicLong(workerId);
this.datacenterId = new AtomicLong(datacenterId);
}
public synchronized long nextId() {
long timestamp = timeGen();
if (timestamp < lastTimestamp) {
// 时钟回拨,等待或抛出异常
throw new RuntimeException("Clock moved backwards. Refusing to generate id");
}
if (lastTimestamp == timestamp) {
sequence.set((sequence.get() + 1) & sequenceMask);
if (sequence.get() == 0) {
timestamp = tilNextMillis(lastTimestamp);
}
} else {
sequence.set(0);
}
lastTimestamp = timestamp;
return ((timestamp - twepoch) << timestampLeftShift)
| (datacenterId.get() << datacenterIdShift)
| (workerId.get() << workerIdShift)
| sequence.get();
}
private long tilNextMillis(long lastTimestamp) {
long timestamp = timeGen();
while (timestamp <= lastTimestamp) {
timestamp = timeGen();
}
return timestamp;
}
private long timeGen() {
return System.currentTimeMillis();
}
}
6.4 监控告警必不可少
没有监控的优化都是空谈。我们引入了以下监控:
@Component
public class DatabaseMonitor {
@Autowired
private MeterRegistry meterRegistry;
/**
* 记录数据库查询耗时
*/
@Around("execution(* com.ecommerce.mapper.*.*(..))")
public Object monitorQuery(ProceedingJoinPoint point) throws Throwable {
StopWatch stopWatch = new StopWatch();
stopWatch.start();
try {
return point.proceed();
} finally {
stopWatch.stop();
long duration = stopWatch.getLastTaskTimeMillis();
// 记录到监控平台
meterRegistry.timer("db.query.duration")
.record(duration, TimeUnit.MILLISECONDS);
// 慢查询告警
if (duration > 1000) {
log.warn("慢查询告警: method={}, duration={}ms",
point.getSignature().getName(), duration);
}
}
}
}
七、总结
那次宕机事故后,我们总结出的核心经验:
- 连接池要合理配置,不要依赖默认值,也不要盲目加大
- 读写分离能有效分散读压力,但要注意数据延迟问题
- 分库分表是终极方案,但要慎重选择分片键,并做好数据迁移规划
- 监控告警是保障,没有监控的优化等于盲人摸象
高并发优化不是一蹴而就的,而是一个持续迭代的过程。希望这篇文章能帮到你。