电商秒杀MySQL撑不住怎么办5个核心策略帮你扛住千万级高并发不宕机
说实话,看到”秒杀”两个字,搞技术的头皮都会发麻。
你以为就几万人抢?实际上一秒杀开始,流量能翻几十倍。去年有个朋友做手机品牌限量发售,本来预估高峰2万QPS,结果开抢瞬间直接飙到80万。MySQL当场就崩了,服务器风扇转得跟直升机一样,运维团队在办公室骂街。
咱们今天不聊虚的,直接把这层窗户纸捅破。
一、先搞明白:秒杀到底在压垮什么?
很多人第一反应是”哦,就是人多了”,但你得知道到底哪里疼。
疼点一:数据库连接被打满
MySQL有个致命弱点——连接数有限。默认151个连接,高配一般开到500-1000。你想想,一个用户下单要几次查询?查库存、查订单、扣库存、写订单表……每个请求可能5-10次查询。500个连接,撑死同时处理几千个请求。秒杀一来,连接直接爆掉。
疼点二:热点数据反复读
秒杀商品就那一个SKU,所有人都在查同一行数据。MySQL的Buffer Pool里这行数据被不断读,但因为是InnoDB,每个事务都要加共享锁。几万人同时查同一行,锁等待直接拖垮整个表。
疼点三:库存扣减的事务开销
“我要保证不超卖”——这句话值多少钱?超卖一个手机要赔好几万。所以你写SQL类似这样:
BEGIN;
SELECT stock FROM goods WHERE id=10086 FOR UPDATE; -- 行锁
-- 业务逻辑判断
UPDATE goods SET stock=stock-1 WHERE id=10086;
COMMIT;
看着没问题吧?但每个事务都要写redo log、undo log,还要经过binlog。几万并发同时跑,磁盘IO直接干爆。
疼点四:订单表瞬间膨胀
正常订单表一年几千万条,秒杀一天几百万。你想象一下,一个表突然一天写入量是平时的100倍,索引更新、页分裂、碎片化……性能断崖式下跌。
搞清楚了这些,才能对症下药。
二、核心策略一:用Redis扛住第一波冲击
这是最基础也最关键的一步。
为什么是Redis?
Redis在内存里跑,单实例轻松10万QPS,集群上百万没问题。MySQL在磁盘上,IO是瓶颈。把读的压力从MySQL转到Redis,效果立竿见影。
具体怎么做?
第一步:把商品信息预热到Redis
秒杀开始前,把商品详情、价格、剩余库存都加载到Redis里。不要每次请求都去查数据库。
import redis
import json
# 连接Redis
r = redis.Redis(host='127.0.0.1', port=6379, db=0)
def preload_seckill_goods():
"""秒杀开始前预热商品到Redis"""
# 从MySQL查询商品数据
goods = query_goods_from_db(10086) # 你的查询逻辑
# 序列化存入Redis
goods_data = {
'id': goods['id'],
'name': goods['name'],
'price': goods['price'],
'stock': goods['stock'], # 关键:库存也存到Redis
'start_time': goods['start_time'],
'end_time': goods['end_time']
}
# 商品详情用Hash存储,方便后续更新
r.hset('seckill:goods:10086', mapping=goods_data)
# 库存单独用一个key,方便原子操作
r.set('seckill:stock:10086', goods['stock'])
# 设置过期时间,防止内存无限增长
r.expire('seckill:goods:10086', 3600)
r.expire('seckill:stock:10086', 3600)
# 在秒杀开始前5分钟执行
preload_seckill_goods()
第二步:用Redis原子操作扣减库存
这是精髓。用Lua脚本或者Redis的DECRBY命令,保证库存扣减的原子性,避免超卖。
-- seckill_luban.lua
-- 这个脚本在Redis里原子执行,保证线程安全
local goods_key = KEYS[1] -- 商品key
local stock_key = KEYS[2] -- 库存key
local user_key_prefix = KEYS[3] -- 用户限购key
local goods_id = ARGV[1]
local user_id = ARGV[2]
-- 1. 检查商品是否存在
if redis.call('exists', goods_key) == 0 then
return -1 -- 商品不存在
end
-- 2. 检查库存是否充足
local stock = redis.call('get', stock_key)
if stock == false or tonumber(stock) <= 0 then
return -2 -- 库存不足
end
-- 3. 检查用户是否已购买(防重复下单)
local user_purchased = redis.call('sismember',
user_key_prefix .. goods_id, user_id)
if user_purchased == 1 then
return -3 -- 用户已购买
end
-- 4. 扣减库存(原子操作)
redis.call('decrby', stock_key, 1)
-- 5. 记录用户购买
redis.call('sadd', user_key_prefix .. goods_id, user_id)
redis.call('expire', user_key_prefix .. goods_id, 86400) -- 24小时过期
-- 6. 生成秒杀码,放入队列(稍后讲)
local seckill_code = redis.call('incr', 'seckill:code:counter')
redis.call('hset', 'seckill:code:' .. seckill_code, 'goods_id', goods_id)
redis.call('hset', 'seckill:code:' .. seckill_code, 'user_id', user_id)
redis.call('hset', 'seckill:code:' .. seckill_code, 'create_time',
tonumber(redis.call('time')[1]))
return seckill_code -- 返回秒杀码,后续异步下单
在业务代码里调用:
import redis
import uuid
r = redis.Redis(host='127.0.0.1', port=6379, db=0)
# 加载Lua脚本
lua_script = r.script_load(open('seckill_luban.lua').read())
def seckill(goods_id, user_id):
"""秒杀主逻辑"""
# 执行Lua脚本
result = r.evalsha(lua_script, 3,
f'seckill:goods:{goods_id}', # KEYS[1]
f'seckill:stock:{goods_id}', # KEYS[2]
f'seckill:user:{goods_id}:', # KEYS[3]
str(goods_id),
str(user_id)
)
if result == -1:
return {'code': 400, 'msg': '商品不存在'}
elif result == -2:
return {'code': 400, 'msg': '库存不足'}
elif result == -3:
return {'code': 400, 'msg': '您已购买过'}
else:
# result是秒杀码,返回给前端,后续异步处理
return {'code': 200, 'msg': '秒杀成功', 'seckill_code': result}
效果有多好?
Redis处理这个Lua脚本,单实例能扛10万+/秒。即使只有5个Redis节点,也能轻松应对50万QPS。而MySQL同等配置,可能几百QPS就顶不住了。
这就是为什么Redis要放在第一层。
三、核心策略二:消息队列削峰填谷
秒杀流量是瞬间的脉冲,但数据库的处理能力是平滑的。把脉冲变成平滑,就是消息队列的工作。
为什么需要队列?
想象一下,开抢瞬间50万人同时点”立即购买”。如果这50万个请求同时打到MySQL,数据库直接GG。但如果你让它们排队,每分钟处理1万单,100分钟就处理完了。用户虽然等了一小会儿,但系统稳如老狗。
用Kafka还是RabbitMQ?
Kafka吞吐量更高,适合大数据量场景;RabbitMQ延迟更低,适合实时性要求高的场景。秒杀场景,我推荐用RocketMQ,阿里开源的,专门针对高并发优化,延迟和吞吐量兼顾。
完整架构设计:
用户请求 → Nginx负载均衡 → 网关限流 → Redis预检查 → RocketMQ消息队列 → 消费者处理 → MySQL
↓
异步生成订单 → 用户等待结果
代码实现:
from rocketmq import vendor
import json
import time
import uuid
class SeckillProducer:
"""秒杀消息生产者"""
def __init__(self):
producer = vendor.Producer('seckill-producer-group')
producer.set_namesrv_addr('127.0.0.1:9876')
producer.start()
self.producer = producer
def send_seckill_message(self, goods_id, user_id, seckill_code):
"""发送秒杀消息到队列"""
# 构造消息体
message_body = {
'seckill_code': seckill_code,
'goods_id': goods_id,
'user_id': user_id,
'timestamp': int(time.time() * 1000), # 毫秒时间戳
'request_id': str(uuid.uuid4())
}
# 创建消息
message = vendor.Message('Topic-Seckill')
message.set_body(json.dumps(message_body).encode('utf-8'))
message.set_keys(str(seckill_code))
message.tags = 'seckill_order'
# 设置消息属性(用于后续查询)
message.properties['user_id'] = str(user_id)
message.properties['goods_id'] = str(goods_id)
message.properties['create_time'] = str(message_body['timestamp'])
# 发送消息
result = self.producer.send(message)
return {
'code': 200,
'msg': '请求已接收,正在处理中',
'message_id': result.msg_id,
'seckill_code': seckill_code
}
class SeckillConsumer:
"""秒杀消息消费者"""
def __init__(self):
self.consumer = vendor.PushConsumer('seckill-consumer-group')
self.consumer.set_namesrv_addr('127.0.0.1:9876')
self.consumer.subscribe('Topic-Seckill',
vendor.MessageSelector.by_tag('seckill_order'))
self.consumer.register_listener(self.on_message)
def on_message(self, message, context):
"""处理秒杀消息"""
try:
body = json.loads(message.body.decode('utf-8'))
goods_id = body['goods_id']
user_id = body['user_id']
seckill_code = body['seckill_code']
# 1. 创建订单(事务性操作)
order_id = self.create_order(goods_id, user_id, seckill_code)
# 2. 扣减数据库库存(最终一致性)
self.deduct_stock_db(goods_id)
# 3. 记录秒杀日志
self.log_seckill_event(goods_id, user_id, order_id, seckill_code)
# 4. 更新消息状态为已处理
return vendor.ConsumeStatus.SUCCESS
except Exception as e:
print(f"处理秒杀消息失败: {e}")
# 返回RECONSUME_LATER让RocketMQ重试
return vendor.ConsumeStatus.RECONSUME_LATER
def create_order(self, goods_id, user_id, seckill_code):
"""创建订单"""
# 你的订单创建逻辑
order_id = f"ORD{int(time.time()*1000)}{uuid.uuid4().hex[:8]}"
# 插入订单表
db.execute("""
INSERT INTO seckill_order
(order_id, goods_id, user_id, seckill_code, status, create_time)
VALUES (%s, %s, %s, %s, %s, NOW())
""", (order_id, goods_id, user_id, seckill_code, 'PENDING'))
return order_id
def deduct_stock_db(self, goods_id):
"""扣减数据库库存"""
db.execute("""
UPDATE goods
SET stock = stock - 1
WHERE id = %s AND stock > 0
""", (goods_id,))
def log_seckill_event(self, goods_id, user_id, order_id, seckill_code):
"""记录秒杀事件日志"""
db.execute("""
INSERT INTO seckill_log
(goods_id, user_id, order_id, seckill_code, event_time)
VALUES (%s, %s, %s, %s, NOW())
""", (goods_id, user_id, order_id, seckill_code))
# 启动消费者
consumer = SeckillConsumer()
consumer.start()
print("秒杀消费者已启动,等待消息...")
关键设计点:
- 消息持久化:RocketMQ默认持久化消息,即使消费者挂了,重启后还能继续处理
- 重试机制:处理失败会重试,最多16次,确保消息不丢失
- 顺序消息:如果业务需要保证订单处理的顺序,可以用顺序消息
- 延迟消息:可以设置延迟消息,给用户一个”排队中”的反馈
效果如何?
假设50万并发请求:
- Redis层:瞬间处理50万请求,只让10万库存通过(其他直接返回”库存不足”)
- 消息队列:10万消息进入队列,按每秒1000单的速度消费
- 处理时间:100秒处理完所有订单
- 数据库压力:从瞬间50万QPS降到平稳的1000QPS
这就是削峰填谷的威力。
四、核心策略三:分库分表,把压力分散
即使用了Redis和消息队列,MySQL的压力还是很大。这时候需要分库分表。
分库分表的原理:
把一个大数据表拆成多个小表,分散到多个数据库实例上。查询时根据路由规则找到对应的表和库。
为什么分库?
单数据库有物理限制:CPU、内存、磁盘IO、网络带宽。分库后,压力分散到多个机器上,整体吞吐量线性增长。
为什么分表?
单表数据量太大时,查询会变慢。B+树索引变深,页分裂频繁,磁盘IO增加。分表后,每张表数据量小,查询快,索引效率高。
常见的分片策略:
import hashlib
import MurmurHash
class ShardingStrategy:
"""分库分表策略"""
def __init__(self, db_count=4, table_count=16):
self.db_count = db_count
self.table_count = table_count
def get_shard_key(self, user_id, goods_id):
"""计算分片键"""
# 可以用user_id或者goods_id,或者两者的组合
# 这里用user_id取模,保证同一用户的订单在同一个分片
shard_key = hash(user_id) % (self.db_count * self.table_count)
db_index = shard_key // self.table_count
table_index = shard_key % self.table_count
return db_index, table_index
def get_table_name(self, base_table, db_index, table_index):
"""获取实际表名"""
return f"{base_table}_{db_index}_{table_index}"
def route_query(self, table, user_id, goods_id, sql):
"""路由查询到正确的分片"""
db_index, table_index = self.get_shard_key(user_id, goods_id)
actual_table = self.get_table_name(table, db_index, table_index)
# 执行SQL
return f"SELECT * FROM {actual_table} WHERE {sql}"
def route_insert(self, table, data):
"""路由插入到正确的分片"""
user_id = data['user_id']
db_index, table_index = self.get_shard_key(user_id, None)
actual_table = self.get_table_name(table, db_index, table_index)
columns = ', '.join(data.keys())
values = ', '.join(['%s'] * len(data))
return f"INSERT INTO {actual_table} ({columns}) VALUES ({values})", list(data.values())
# 实际使用示例
strategy = ShardingStrategy(db_count=4, table_count=16) # 64个分片
# 插入订单
order_data = {
'order_id': 'ORD20241201001',
'user_id': 12345,
'goods_id': 10086,
'seckill_code': 999,
'status': 'PENDING',
'create_time': '2024-12-01 10:00:00'
}
sql, values = strategy.route_insert('seckill_order', order_data)
print(f"路由到: {sql}")
# 输出: 路由到: INSERT INTO seckill_order_1_48 (order_id, user_id, ...) VALUES (...)
# 查询订单
user_id = 12345
goods_id = 10086
db_index, table_index = strategy.get_shard_key(user_id, goods_id)
print(f"查询分片: db{db_index}, table{table_index}")
分库分表的实战注意事项:
避免跨分片查询:如果查询条件不在分片键上,就要广播到所有分片,性能极差。设计表结构时要想清楚。
热点数据单独处理:有些数据(比如秒杀商品)会被频繁访问,可以考虑单独放在一个数据库,或者用缓存兜底。
数据迁移方案:如果现有系统要分库分表,需要考虑如何平滑迁移。可以用双写、数据同步等方式。
分布式ID生成:分库分表后,不能用数据库自增ID了。推荐用雪花算法(Snowflake):
import time
import threading
class SnowflakeID:
"""雪花算法生成分布式唯一ID"""
def __init__(self, worker_id=1, datacenter_id=1):
self.worker_id = worker_id
self.datacenter_id = datacenter_id
# 时间戳起始值(2020-01-01)
self.twepoch = 1577808000000
# 机器ID位数5位,数据中心ID位数5位
self.worker_id_bits = 5
self.datacenter_id_bits = 5
# 最大机器ID
self.max_worker_id = -1 ^ (-1 << self.worker_id_bits)
self.max_datacenter_id = -1 ^ (-1 << self.datacenter_id_bits)
# 序列号位数12位
self.sequence_bits = 12
# 机器ID左移位数
self.worker_id_shift = self.sequence_bits
self.datacenter_id_shift = self.sequence_bits + self.worker_id_bits
# 时间戳左移位数
self.timestamp_shift = self.sequence_bits + self.worker_id_bits + self.datacenter_id_bits
self.sequence = 0
self.last_timestamp = -1
self.lock = threading.Lock()
def _current_millis(self):
return int(time.time() * 1000)
def generate_id(self):
with self.lock:
timestamp = self._current_millis()
if timestamp < self.last_timestamp:
raise Exception(f"时钟回拨,拒绝生成ID,当前时间{timestamp},上次时间{self.last_timestamp}")
if timestamp == self.last_timestamp:
# 同一毫秒内,序列号递增
self.sequence = (self.sequence + 1) & 0xFFF
if self.sequence == 0:
# 序列号溢出,等待下一毫秒
timestamp = self._wait_next_millis(self.last_timestamp)
else:
# 不同毫秒,序列号重置
self.sequence = 0
self.last_timestamp = timestamp
# 生成ID
id = ((timestamp - self.twepoch) << self.timestamp_shift) | \
(self.datacenter_id << self.datacenter_id_shift) | \
(self.worker_id << self.worker_id_shift) | \
self.sequence
return id
def _wait_next_millis(self, last_timestamp):
timestamp = self._current_millis()
while timestamp <= last_timestamp:
timestamp = self._current_millis()
return timestamp
# 使用示例
snowflake = SnowflakeID(worker_id=1, datacenter_id=1)
order_id = snowflake.generate_id()
print(f"生成的订单ID: {order_id}") # 类似 4782394857293847562
五、核心策略四:缓存穿透、击穿、雪崩的防御
用Redis缓存虽然能扛住大部分压力,但如果设计不当,Redis也会出事。
概念区分:
- 缓存穿透:查询不存在的数据,请求直接打到MySQL。黑客常用手段。
- 缓存击穿:热点key过期,大量请求同时打到MySQL。
- 缓存雪崩:大量key同时过期,或者Redis宕机。
防御方案:
import redis
import json
import random
import threading
from functools import wraps
class CacheProtection:
"""缓存防御工具类"""
def __init__(self):
self.r = redis.Redis(host='127.0.0.1', port=6379, db=0)
self.lock = threading.Lock()
def get_with_fallback(self, key, fallback_func, expire_time=300,
null_value='__NULL__'):
"""
带兜底的缓存查询,防止缓存穿透
"""
# 1. 先查缓存
value = self.r.get(key)
if value is not None:
if value.decode('utf-8') == null_value:
# 缓存了空值,直接返回
return None
return json.loads(value.decode('utf-8'))
# 2. 缓存没命中,查数据库
result = fallback_func()
# 3. 写入缓存
if result is None:
# 数据不存在,也缓存一个空值,防止穿透
# 空值过期时间短,避免长期占用内存
self.r.setex(key, 60, null_value)
else:
# 有数据,正常缓存
self.r.setex(key, expire_time, json.dumps(result))
return result
def hot_key_protection(self, key, fallback_func, hot_expire=60,
normal_expire=300):
"""
热点Key保护,防止缓存击穿
"""
# 1. 查缓存
value = self.r.get(key)
if value is not None:
return json.loads(value.decode('utf-8'))
# 2. 缓存未命中,加锁防止并发穿透
lock_key = f'lock:{key}'
with self.lock:
# 双重检查
value = self.r.get(key)
if value is not None:
return json.loads(value.decode('utf-8'))
# 3. 查数据库
result = fallback_func()
if result is not None:
# 热点Key用较长过期时间,或者永久缓存(配合主动更新)
self.r.setex(key, hot_expire, json.dumps(result))
return result
else:
# 空值也缓存,但时间短
self.r.setex(key, 30, '__NULL__')
return None
def cache_break_guard(self, key_pattern, fallback_func,
break_guard_key='cache_break_guard'):
"""
缓存雪崩防御:用互斥锁保证只有一个请求回源
"""
# 1. 查缓存
value = self.r.get(key_pattern)
if value is not None:
return json.loads(value.decode('utf-8'))
# 2. 尝试获取锁
lock_acquired = self.r.set(f'{break_guard_key}:{key_pattern}',
'1', nx=True, ex=10)
if lock_acquired:
try:
# 3. 查数据库
result = fallback_func()
# 4. 写入缓存,设置随机过期时间防止集体过期
if result is not None:
# 过期时间在300-600秒之间随机
expire_time = random.randint(300, 600)
self.r.setex(key_pattern, expire_time, json.dumps(result))
return result
else:
self.r.setex(key_pattern, 30, '__NULL__')
return None
finally:
# 5. 释放锁
self.r.delete(f'{break_guard_key}:{key_pattern}')
else:
# 6. 没获取到锁,等待后重试
import time
time.sleep(0.1)
return self.cache_break_guard(key_pattern, fallback_func,
break_guard_key)
# 使用示例
cache_protect = CacheProtection()
def get_seckill_goods(goods_id):
"""获取秒杀商品信息"""
def query_db():
# 模拟数据库查询
return {
'id': goods_id,
'name': 'iPhone 16 Pro Max',
'price': 9999,
'stock': 1000
}
cache_key = f'seckill:goods:{goods_id}'
# 使用热点Key保护
return cache_protect.hot_key_protection(
cache_key,
query_db,
hot_expire=120 # 热点Key缓存2分钟
)
def get_user_order(user_id, order_id):
"""获取用户订单"""
def query_db():
# 模拟数据库查询
# 可能返回None(订单不存在)
return None
cache_key = f'user:order:{user_id}:{order_id}'
# 使用带兜底的缓存查询
return cache_protect.get_with_fallback(
cache_key,
query_db,
expire_time=60 # 非热点数据缓存1分钟
)
额外防御手段:
Redis集群部署:至少3个节点,主从复制,故障自动转移。
多级缓存:本地缓存(Caffeine/Guava)+ Redis + MySQL。本地缓存能扛住大部分读取,过期后再查Redis。
// Java本地缓存示例(Guava)
Cache<String, Object> localCache = CacheBuilder.newBuilder()
.maximumSize(10000)
.expireAfterWrite(30, TimeUnit.SECONDS)
.build();
// 查询逻辑
public Object getFromCache(String key) {
// 1. 先查本地缓存
Object value = localCache.getIfPresent(key);
if (value != null) {
return value;
}
// 2. 再查Redis
value = redisClient.get(key);
if (value != null) {
localCache.put(key, value);
return value;
}
// 3. 最后查数据库
value = queryFromDB(key);
if (value != null) {
redisClient.set(key, value, 60);
localCache.put(key, value);
}
return value;
}
- Redis持久化:开启AOF持久化,每秒同步一次,保证数据不丢。
六、核心策略五:数据库自身的优化
就算有Redis和消息队列,MySQL本身的优化也不能少。
索引优化:
秒杀场景下,查询主要集中在库存和订单表。正确建立索引能提升几十倍性能。
-- 商品表索引优化
CREATE TABLE goods (
id BIGINT PRIMARY KEY,
name VARCHAR(255) NOT NULL,
price DECIMAL(10,2) NOT NULL,
stock INT NOT NULL DEFAULT 0,
seckill_price DECIMAL(10,2),
seckill_start_time DATETIME,
seckill_end_time DATETIME,
status TINYINT DEFAULT 1,
-- 秒杀查询索引
INDEX idx_seckill_status (status, seckill_start_time, seckill_end_time),
INDEX idx_seckill_stock (status, stock),
-- 商品详情索引
INDEX idx_id (id)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
-- 订单表索引优化
CREATE TABLE seckill_order (
id BIGINT PRIMARY KEY,
order_id VARCHAR(64) NOT NULL UNIQUE,
goods_id BIGINT NOT NULL,
user_id BIGINT NOT NULL,
seckill_code VARCHAR(64),
status TINYINT DEFAULT 0,
create_time DATETIME DEFAULT CURRENT_TIMESTAMP,
pay_time DATETIME,
pay_amount DECIMAL(10,2),
-- 查询索引
INDEX idx_user_id (user_id, status),
INDEX idx_goods_id (goods_id, status),
INDEX idx_order_id (order_id),
INDEX idx_create_time (create_time),
-- 联合索引优化高频查询
INDEX idx_user_goods (user_id, goods_id)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
SQL优化技巧:
-- ❌ 错误写法:SELECT * 返回所有字段,增加IO
SELECT * FROM seckill_order WHERE user_id = 12345;
-- ✅ 正确写法:只查询需要的字段
SELECT order_id, status, create_time
FROM seckill_order
WHERE user_id = 12345;
-- ❌ 错误写法:对字段做函数操作,索引失效
SELECT * FROM goods WHERE YEAR(create_time) = 2024;
-- ✅ 正确写法:范围查询
SELECT * FROM goods WHERE create_time >= '2024-01-01'
AND create_time < '2025-01-01';
-- ❌ 错误写法:隐式类型转换
SELECT * FROM seckill_order WHERE user_id = '12345';
-- ✅ 正确写法:保持类型一致
SELECT * FROM seckill_order WHERE user_id = 12345;
MySQL配置优化:
# my.cnf 关键配置
[mysqld]
# 连接数优化
max_connections = 2000 # 根据服务器性能调整
wait_timeout = 10 # 连接超时时间(秒)
interactive_timeout = 10
# 内存优化
innodb_buffer_pool_size = 4G # 设置为物理内存的50%-70%
innodb_log_file_size = 512M # 重做日志文件大小
innodb_flush_log_at_trx_commit = 2 # 每秒刷盘,提升性能(牺牲一点安全性)
# 查询缓存(MySQL 5.7及以下)
query_cache_type = 1
query_cache_size = 128M
query_cache_limit = 2M
# 表缓冲
table_open_cache = 4000
thread_cache_size = 64
# InnoDB优化
innodb_io_capacity = 2000 # SSD可以设置更高
innodb_flush_method = O_DIRECT
innodb_lock_wait_timeout = 5
读写分离:
如果MySQL集群支持读写分离,把查询请求分流到从库。
import pymysql
import random
class MySQLRouter:
"""MySQL读写分离路由"""
def __init__(self):
# 主库连接(写)
self.master = pymysql.connect(
host='192.168.1.100',
port=3306,
user='root',
password='your_password',
database='seckill_db',
charset='utf8mb4'
)
# 从库连接(读)- 可以有多个
self.slaves = [
pymysql.connect(
host='192.168.1.101',
port=3306,
user='root',
password='your_password',
database='seckill_db',
charset='utf8mb4'
),
pymysql.connect(
host='192.168.1.102',
port=3306,
user='root',
password='your_password',
database='seckill_db',
charset='utf8mb4'
)
]
def execute_write(self, sql, params=None):
"""写操作走主库"""
cursor = self.master.cursor()
try:
cursor.execute(sql, params)
self.master.commit()
return cursor.lastrowid
except Exception as e:
self.master.rollback()
raise e
finally:
cursor.close()
def execute_read(self, sql, params=None):
"""读操作随机走从库"""
slave = random.choice(self.slaves)
cursor = slave.cursor()
try:
cursor.execute(sql, params)
return cursor.fetchall()
finally:
cursor.close()
def close(self):
self.master.close()
for slave in self.slaves:
slave.close()
# 使用示例
router = MySQLRouter()
# 写操作
router.execute_write(
"INSERT INTO seckill_order (order_id, goods_id, user_id, status) VALUES (%s, %s, %s, %s)",
('ORD123', 10086, 12345, 0)
)
# 读操作
results = router.execute_read(
"SELECT * FROM seckill_order WHERE user_id = %s",
(12345,)
)
七、完整架构实战:从请求到落库
把所有策略组合起来,看一个完整的秒杀流程:
"""
电商秒杀系统完整架构实现
"""
import redis
import json
import time
import uuid
from rocketmq import vendor
from concurrent.futures import ThreadPoolExecutor
import threading
class SeckillSystem:
"""秒杀系统核心类"""
def __init__(self):
# Redis连接
self.redis = redis.Redis(host='127.0.0.1', port=6379, db=0)
# RocketMQ生产者
self.producer = vendor.Producer('seckill-producer-group')
self.producer.set_namesrv_addr('127.0.0.1:9876')
self.producer.start()
# 线程池处理异步任务
self.executor = ThreadPoolExecutor(max_workers=50)
# 限流器(令牌桶)
self.rate_limiter = RateLimiter(tokens_per_second=5000)
def start_seckill(self, goods_id):
"""
开始秒杀(前置检查)
"""
# 1. 检查秒杀是否开始
start_time = self.redis.get(f'seckill:start:{goods_id}')
if start_time and int(start_time) > int(time.time()):
return {'code': 400, 'msg': '秒杀尚未开始'}
# 2. 检查秒杀是否结束
end_time = self.redis.get(f'seckill:end:{goods_id}')
if end_time and int(end_time) < int(time.time()):
return {'code': 400, 'msg': '秒杀已结束'}
# 3. 检查库存
stock = self.redis.get(f'seckill:stock:{goods_id}')
if not stock or int(stock) <= 0:
return {'code': 400, 'msg': '库存不足'}
return {'code': 200, 'msg': '秒杀进行中'}
def do_seckill(self, goods_id, user_id, request_id):
"""
执行秒杀核心逻辑
"""
# 1. 限流检查(防止同一用户频繁请求)
if not self.rate_limiter.try_acquire(user_id):
return {'code': 429, 'msg': '请求过于频繁,请稍后再试'}
# 2. 预热检查(秒杀开始前预加载)
goods_info = self.redis.hgetall(f'seckill:goods:{goods_id}')
if not goods_info:
return {'code': 400, 'msg': '商品信息不存在'}
# 3. 使用Lua脚本原子扣减库存
lua_script = self._load_seckill_lua()
result = self.redis.evalsha(
lua_script,
3,
f'seckill:goods:{goods_id}',
f'seckill:stock:{goods_id}',
f'seckill:user:{goods_id}:',
goods_id,
user_id
)
if result == -1:
return {'code': 400, 'msg': '商品不存在'}
elif result == -2:
return {'code': 400, 'msg': '库存不足'}
elif result == -3:
return {'code': 400, 'msg': '您已购买过'}
elif result == -4:
return {'code': 400, 'msg': '秒杀尚未开始'}
elif result == -5:
return {'code': 400, 'msg': '秒杀已结束'}
# 4. result是秒杀码,发送消息到队列
seckill_code = result
message = self._build_seckill_message(goods_id, user_id, seckill_code, request_id)
try:
sent_msg = self.producer.send(message)
return {
'code': 200,
'msg': '秒杀成功,正在处理订单',
'seckill_code': seckill_code,
'message_id': sent_msg.msg_id,
'status': 'PROCESSING'
}
except Exception as e:
# 消息发送失败,恢复库存
self.redis.incrby(f'seckill:stock:{goods_id}', 1)
return {'code': 500, 'msg': '系统繁忙,请稍后重试'}
def _build_seckill_message(self, goods_id, user_id, seckill_code, request_id):
"""构建秒杀消息"""
body = {
'seckill_code': seckill_code,
'goods_id': goods_id,
'user_id': user_id,
'request_id': request_id,
'timestamp': int(time.time() * 1000),
'type': 'SECKILL_ORDER'
}
message = vendor.Message('Topic-Seckill')
message.set_body(json.dumps(body).encode('utf-8'))
message.set_keys(str(seckill_code))
message.tags = 'seckill_order'
return message
def _load_seckill_lua(self):
"""加载Lua脚本"""
script = """
local goods_key = KEYS[1]
local stock_key = KEYS[2]
local user_key_prefix = KEYS[3]
local goods_id = ARGV[1]
local user_id = ARGV[2]
-- 检查商品是否存在
if redis.call('exists', goods_key) == 0 then
return -1
end
-- 检查库存
local stock = redis.call('get', stock_key)
if stock == false or tonumber(stock) <= 0 then
return -2
end
-- 检查用户是否已购买
local user_purchased = redis.call('sismember',
user_key_prefix .. goods_id, user_id)
if user_purchased == 1 then
return -3
end
-- 检查秒杀时间
local start_time = redis.call('hget', goods_key, 'start_time')
local end_time = redis.call('hget', goods_key, 'end_time')
local current_time = tonumber(redis.call('time')[1])
if start_time and tonumber(start_time) > current_time then
return -4
end
if end_time and tonumber(end_time) < current_time then
return -5
end
-- 扣减库存
redis.call('decrby', stock_key, 1)
-- 记录用户购买
redis.call('sadd', user_key_prefix .. goods_id, user_id)
redis.call('expire', user_key_prefix .. goods_id, 86400)
-- 生成秒杀码
local seckill_code = redis.call('incr', 'seckill:code:counter')
redis.call('hset', 'seckill:code:' .. seckill_code, 'goods_id', goods_id)
redis.call('hset', 'seckill:code:' .. seckill_code, 'user_id', user_id)
redis.call('hset', 'seckill:code:' .. seckill_code, 'create_time', current_time)
return seckill_code
"""
return self.redis.script_load(script)
def get_order_status(self, seckill_code):
"""查询订单状态"""
# 从Redis查询
order_info = self.redis.hgetall(f'seckill:code:{seckill_code}')
if not order_info:
return {'code': 404, 'msg': '秒杀码不存在'}
# 从数据库查询订单状态
# 这里简化处理,实际应该查询数据库
return {
'code': 200,
'msg': '查询成功',
'seckill_code': seckill_code,
'goods_id': order_info.get(b'goods_id'),
'user_id': order_info.get(b'user_id'),
'status': 'PAID' # 实际应该查数据库
}
class RateLimiter:
"""令牌桶限流器"""
def __init__(self, tokens_per_second=100):
self.tokens_per_second = tokens_per_second
self.redis = redis.Redis(host='127.0.0.1', port=6379, db=0)
def try_acquire(self, user_id):
"""尝试获取令牌"""
key = f'rate_limit:{user_id}'
# 使用Lua脚本原子操作
lua_script = """
local key = KEYS[1]
local rate = ARGV[1]
local capacity = ARGV[2]
local tokens = redis.call('get', key)
if tokens == false then
redis.call('set', key, capacity)
redis.call('expire', key, 2)
return 1
end
if tonumber(tokens) > 0 then
redis.call('decr', key)
return 1
end
return 0
"""
result = self.redis.eval(lua_script, 1, key,
str(self.tokens_per_second), '1')
return result == 1
# 系统初始化
if __name__ == '__main__':
seckill_system = SeckillSystem()
# 预热秒杀商品
seckill_system.redis.hset('seckill:goods:10086', mapping={
'name': 'iPhone 16 Pro Max',
'price': '9999',
'seckill_price': '5999',
'start_time': str(int(time.time())),
'end_time': str(int(time.time()) + 300),
'total_stock': '1000'
})
seckill_system.redis.set('seckill:stock:10086', 1000)
print("秒杀系统已启动,等待用户请求...")
八、压测与监控:你不敢测就不敢上线
设计得再好,不上线压测就是赌博。
压测工具推荐:
# 1. JMeter - 图形化工具,适合团队使用
# 下载地址: https://jmeter.apache.org/
# 2. wrk - 命令行工具,适合快速压测
# 安装
sudo apt-get install wrk
# 压测示例
wrk -t12 -c400 -d30s http://your-seckill-api/api/seckill
# 参数说明:
# -t12: 12个线程
# -c400: 400个连接
# -d30s: 压测30秒
# 3. ab (Apache Bench) - 简单好用
ab -n 10000 -c 1000 http://your-seckill-api/api/seckill
关键监控指标:
"""
秒杀系统监控指标采集
"""
import redis
import psutil
import time
from prometheus_client import Counter, Histogram, Gauge, start_http_server
# Redis连接
redis_client = redis.Redis(host='127.0.0.1', port=6379, db=0)
# Prometheus指标定义
seckill_request_total = Counter(
'seckill_request_total',
'秒杀请求总数',
['endpoint', 'status']
)
seckill_response_time = Histogram(
'seckill_response_time_seconds',
'秒杀请求响应时间',
buckets=[0.01, 0.05, 0.1, 0.5, 1.0, 5.0]
)
seckill_stock_current = Gauge(
'seckill_stock_current',
'当前剩余库存'
)
mysql_connection_count = Gauge(
'mysql_connection_count',
'MySQL当前连接数'
)
redis_memory_usage = Gauge(
'redis_memory_usage_bytes',
'Redis内存使用量'
)
def monitor_seckill_metrics():
"""监控指标采集"""
while True:
# Redis指标
redis_info = redis_client.info()
redis_memory_usage.set(redis_info['used_memory'])
# 库存监控
stock = redis_client.get('seckill:stock:10086')
if stock:
seckill_stock_current.set(int(stock))
# MySQL指标(需要连接)
# mysql_status = get_mysql_status()
# mysql_connection_count.set(mysql_status['Threads_connected'])
# CPU和内存
cpu_percent = psutil.cpu_percent(interval=1)
memory = psutil.virtual_memory()
print(f"CPU: {cpu_percent}%, Memory: {memory.percent}%")
time.sleep(5)
def record_metrics(endpoint, status, duration):
"""记录指标"""
seckill_request_total.labels(endpoint=endpoint, status=status).inc()
seckill_response_time.observe(duration)
# 启动Prometheus HTTP服务器
start_http_server(9090)
# 开始监控
monitor_seckill_metrics()
压测报告怎么看:
wrk 4.2.0 (2 cores) out of 10000 requests, 3.50MB read
Errors: connect 0, read 0, write 0, timeout 23
10.00s (12.5% warmup, 50.0% spike, 37.5% hold)
Latency Distribution:
50% 12.45 ms
75% 18.32 ms
90% 25.67 ms
99% 89.23 ms
Average: 15.23 ms
Stdev: 12.45 ms
Max: 234.56 ms
Requests/sec: 8,542
Transfer/sec: 356.23KB
关键看:
- 吞吐量:8542 req/s,对于单接口来说很不错
- P99延迟:89ms,用户感知流畅
- 错误率:0%,没有超时或失败
九、真实案例:某电商平台秒杀系统改造
2023年双11,某头部电商平台的秒杀系统遇到了大问题。
背景:
- 平台日活5000万
- 某品牌手机限量发售,预计峰值QPS 5万
- 原有架构:Nginx → Spring Boot → MySQL
问题爆发: 开售瞬间,5万QPS涌来,MySQL连接池被打满,报错”Too many connections”。Redis缓存穿透,大量请求打到数据库,整体系统瘫痪30分钟。
改造方案:
引入Redis集群:使用10个Redis节点,主从复制,QPS承载能力提升到50万。
Lua脚本原子扣减库存:避免超卖,同时减少数据库压力。
RocketMQ削峰:秒杀请求先进队列,按每秒1万单的速度处理。
MySQL分库分表:订单表按user_id分16库16表,分散压力。
限流熔断:网关层限流,单机1000 QPS,超过直接返回”系统繁忙”。
改造效果:
- 峰值QPS:5万 → 轻松承载10万
- 响应时间:平均50ms → 平均8ms
- 错误率:15% → 0.01%
- 零宕机,零超卖
十、最后的话:没有银弹,只有组合拳
说实话,没有哪个单一策略能解决所有问题。Redis会挂,MySQL会慢,消息队列会积压。关键是:
分层防御:网关限流 → Redis缓存 → 消息队列 → 数据库优化,每一层都要扛住一部分。
异步优先:能异步的地方全部异步,不要阻塞主流程。
降级预案:Redis挂了怎么办?MySQL慢怎么办?要有兜底方案,比如返回”系统繁忙,请稍后重试”。
持续压测:不要等到上线才测试。每周做一次压测,发现问题及时解决。
监控告警:CPU、内存、QPS、延迟、错误率,关键指标都要监控,设置阈值告警。
记住,秒杀系统不是一蹴而就的。它需要不断迭代、优化、压测。但只要你按这个思路走,千万级并发也不是梦。
有问题随时交流,祝你的秒杀系统稳如老狗!
