第27章:Mongodb连接池与高并发写入——秒杀库存如何抗住
1. 项目背景
业务场景:本地生活电商的"双十一"大促进入倒计时。运营团队放出 100 件飞天茅台,活动页面刚上线不到 1 秒,库存归零——但后端日志显示超过 800 次的写入请求"成功"了。开发查看代码发现致命问题:库存扣减逻辑是"先 findOne 查库存 → if > 0 → updateOne 扣减"的三步走。在高并发下,100 个人同时读到库存 > 0,然后 100 个人都执行了 updateOne——最终库存变成了 -85。更糟糕的是连接池配置——Spring Boot 默认 maxPoolSize=100,秒杀期间 5 万 QPS,连接池满导致 3000 多个请求排队等 2 分钟后超时。
痛点:高并发写入场景下,连接池和写入策略是相生相克的。连接池太小——请求排队积压,上游超时重试加剧雪崩;连接池太大——MongoDB 的连接线程数爆增,上下文切换吃掉 CPU。写入策略错误——查后写(read-then-write)导致超卖;不用 bulkWrite 导致每秒只能写入几百条而非数万条;没有幂等键导致相同请求被重试时创建了多条重复记录。
2. 项目设计
小胖(盯着压测报告目瞪口呆):大师!秒杀接口模拟 5000 并发,连接池满、超时报错、库存还变成了负数!这怎么搞?
大师:秒杀的高并发写入有三个要点,缺一不可——原子扣减、幂等去重、连接池压测调优。你现在的代码是三步走的查后写(read-then-write),自然超卖。
小胖:那怎么写才对?
大师:一步到位——用原子操作符$inc配合 filter 条件。关键代码一行就够了:
db.products.updateOne({_id:productId,stock:{$gt:0}},{$inc:{stock:-1},$set:{updatedAt:newDate()}})如果modifiedCount > 0——扣减成功,抢到了;modifiedCount === 0——库存不足,没抢到。这条 update 在 MongoDB 内部是原子执行的——stock > 0的检查和stock - 1的更新在同一个原子操作中完成,不存在时间窗口。
技术映射:原子操作符($inc+ filter 条件)替代 read-then-write 是 MongoDB 并发写入的第一原则。这是 MongoDB 的硬拳头——充分利用单文档原子性,不给竞态条件留缝隙。
小胖:那bulkWrite呢?秒杀完之后还得插领取记录吧?
大师:对。秒杀结束、用户扣库成功后,还需要插入条领取记录。但如果把扣库存和插记录拆成两次独立的操作,中间有网络往返延迟。正确写法是用bulkWrite把多个操作打成一个包发到服务器——一次网络往返完成。
技术映射:bulkWrite是 MongoDB 的高吞吐写入利器。ordered: false(乱序模式)下,批量中的每一步独立执行、相互不影响——单条失败不阻塞整批。
小白(追问):那幂等呢?如果网络抖动,用户点了两次"抢购",会生成两条领取记录吗?
大师:这就是唯一索引 + 幂等键的作用。在领取记录表上建{userId: 1, productId: 1}的唯一索引。不管客户端发多少次请求,第二次insertOne会因为唯一冲突被优雅拒绝——不会产生重复记录。
技术映射:幂等 = 重试安全的写入。唯一索引是实现幂等的最简单且最可靠的手段。
小胖:那连接池怎么调?默认 100 够不够?
大师:秒杀的连接池调整要结合实际的并发数和响应延迟。粗略公式:maxPoolSize = (目标 QPS × P95 延迟 ms) / 1000 + 20% 余量。比如目标 QPS 5000,P95 延迟 5ms → 需要 25 个连接,加 20% → 约 30 个连接。但你配 500 个连接是没有用的——MongoDB 每个连接是一个线程,500 个连接 × 5ms 延迟意味着大量线程在上下文切换上浪费 CPU。
技术映射:连接池不是越大越好——它受限于 MongoDB 服务端的并发处理能力。过多的连接导致服务端线程暴涨、上下文切换开销吃掉真正的处理时间。
大师(总结):秒杀三把斧——原子操作防超卖、bulkWrite 提高吞吐、幂等唯一键防重放。连接池大小要根据实测调,别拍脑袋。建议压测时先从 50 个连接开始,逐步加,找到吞吐不再增长的"饱和点"。
3. 项目实战
3.1 环境准备
沿用 Docker MongoDB(复制集环境更好,能模拟真实并发场景)。
dockercompose-fmongodb-lab/docker-compose.ymlps# 或启动 3 节点复制集3.2 分步实现
步骤一:模拟超卖——错误写法 vs 正确写法
目标:对比 read-then-write 和原子操作在高并发下的行为差异。
use local_life// 准备秒杀商品db.seckill_products.drop()db.seckill_products.insertOne({_id:"SKU_MOUTAI",name:"飞天茅台 53度",totalStock:100,stock:100,price:NumberDecimal("1499.00"),status:"active",updatedAt:newDate()})// 领取记录(建唯一索引防止重复)db.seckill_records.drop()db.seckill_records.createIndex({userId:1,productId:1},{unique:true,name:"uk_user_product"})print("秒杀商品就绪:库存 100 瓶")// === 错误写法:read-then-write(超卖演示) ===functionbadSeckill(userId){constproduct=db.seckill_products.findOne({_id:"SKU_MOUTAI"})if(!product||product.stock<=0){return{success:false,reason:"库存不足"}}// ⚠️ 时间窗口!从 read 到 write 之间,其他请求可能已经扣了库存constresult=db.seckill_products.updateOne({_id:"SKU_MOUTAI"},{$inc:{stock:-1},$set:{updatedAt:newDate()}})return{success:result.modifiedCount>0,userId}}// === 正确写法:原子操作 ===functiongoodSeckill(userId){// 步骤一:原子扣库存(一步到位)constdeductResult=db.seckill_products.updateOne({_id:"SKU_MOUTAI",stock:{$gt:0}},// filter 中检查库存{$inc:{stock:-1},$set:{updatedAt:newDate()}})if(deductResult.modifiedCount===0){return{success:false,reason:"库存不足"}}// 步骤二:幂等插入领取记录try{db.seckill_records.insertOne({userId:userId,productId:"SKU_MOUTAI",claimedAt:newDate()})}catch(e){if(e.code===11000){// DuplicateKeyreturn{success:true,reason:"幂等——重复请求视为成功"}}throwe}return{success:true,userId}}步骤二:连接池配置与压测
目标:配置 Driver 连接池参数并观察不同配置的吞吐差异。
// === Spring Boot 连接池配置(application.yml 关键部分) ===// spring:// data:// mongodb:// uri: mongodb://user:pwd@host:27017/db?authSource=admin// connection-pool:// max-size: 50 # 目标连接数// min-size: 10 # 保持 10 个热连接,避免冷启动// max-wait-time: 2000ms # 等待可用连接的超时——2秒快速失败// max-connection-idle-time: 600s// max-connection-life-time: 1800s// 手动设置连接池参数(mongosh 中无法直接演示,此处用说明)// Driver 端的 MongoClientSettings:// MongoClientSettings settings = MongoClientSettings.builder()// .applyConnectionString(new ConnectionString(uri))// .applyToConnectionPoolSettings(b -> b// .maxSize(50)// .minSize(10)// .maxWaitTime(2000, TimeUnit.MILLISECONDS)// .maxConnectionIdleTime(10, TimeUnit.MINUTES))// .build();// === 连接池监控:查看当前连接使用情况 ===use adminconstconnStats=db.serverStatus().connectionsprint("=== MongoDB 连接统计 ===")print("当前连接:",connStats.current)print("可用连接:",connStats.available)// maxIncomingConnections = 65536 默认print("活跃连接:",connStats.active)print("总创建数(自启动):",connStats.totalCreated)print("线程客户端:",connStats.threadedConnections||"N/A")// 如果 current / available > 80%,连接即将耗尽步骤三:bulkWrite 批量写入优化
目标:用 bulkWrite 批量插入记录,对比逐条插入的性能差异。
// === 逐条插入(慢) ===functioninsertOneByOne(count){conststart=Date.now()for(leti=0;i<count;i++){db.bulk_test.insertOne({userId:"U"+i,productId:"SKU_BULK",claimedAt:newDate()})}constelapsed=(Date.now()-start)/1000return{count,elapsed,qps:(count/elapsed).toFixed(0)}}// === bulkWrite 批量(快) ===functioninsertBulkWrite(count,batchSize=500){conststart=Date.now()lettotalInserted=0for(letround=0;round<count;round+=batchSize){constbatch=[]constactualSize=Math.min(batchSize,count-round)for(leti=0;i<actualSize;i++){batch.push({insertOne:{document:{userId:"BULK_U"+(round+i),productId:"SKU_BULK",claimedAt:newDate()}}})}constresult=db.bulk_test.bulkWrite(batch,{ordered:false})totalInserted+=result.insertedCount}constelapsed=(Date.now()-start)/1000return{count:totalInserted,elapsed,qps:(totalInserted/elapsed).toFixed(0)}}// 对比测试db.bulk_test.drop()constsingleResult=insertOneByOne(1000)print("逐条插入 1000条:",singleResult.qps,"条/秒")db.bulk_test.drop()constbulkResult=insertBulkWrite(5000,500)print("bulkWrite 5000条:",bulkResult.qps,"条/秒 (批次大小 500)")// 期望:bulkWrite QPS 是逐条的 5-20 倍步骤四:消费者队列削峰 + 前台快速返回
目标:演示秒杀架构中"前台快速返回 + 后台写入"的削峰模式。
// === Java 削峰代码示例 ===// 前台:秒杀接口只做原子扣库存,快速返回结果@PostMapping("/seckill")publicSeckillResultseckill(@RequestParamStringuserId,@RequestParamStringproductId){// 原子扣库存Queryquery=newQuery(Criteria.where("_id").is(productId).and("stock").gt(0));Updateupdate=newUpdate().inc("stock",-1);longmodified=mongoTemplate.updateFirst(query,update,SeckillProduct.class).getModifiedCount();if(modified==0){returnSeckillResult.fail("已抢光");}// 丢一个事件到队列(异步处理:插领取记录、发消息、刷新缓存等)eventQueue.offer(newSeckillEvent(userId,productId,Instant.now()));returnSeckillResult.success("抢到啦!");}// 后台:消费者从队列读取事件,处理剩余的写入(领券记录、统计等)@EventListenerpublicvoidhandleSeckillEvent(SeckillEventevent){try{mongoTemplate.insert(newSeckillRecord(event.userId,event.productId,event.time),"seckill_records");}catch(DuplicateKeyExceptione){// 幂等——忽略重复}}步骤五:压测与监控——如何找到连接池的最佳值
目标:通过逐步增大连接池,找到应用吞吐的饱和点。
// === 压测监控脚本(在 mongosh 中观察)===// 压测期间在不同窗口运行此脚本functionmonitorDuringLoadTest(intervalMs=2000,rounds=10){for(leti=0;i<rounds;i++){constconn=db.serverStatus().connectionsconstops=db.serverStatus().opcounters// 计算 QPS(两次采样的差值/时间差)consttotalOps=ops.insert+ops.query+ops.update+ops.deleteprint(`[${newDate().toISOString().slice(11,19)}]`+`连接:${conn.current}/${conn.available}活跃:${conn.active}`+`操作累计:${totalOps}`)sleep(intervalMs)}}monitorDuringLoadTest(2000,15)3.3 完整代码清单
| 文件 | 用途 |
|---|---|
mongodb-lab/scripts/ch27-seckill-setup.js | 秒杀场景初始化 |
mongodb-lab/scripts/ch27-oversell-demo.js | 超卖 vs 原子扣减对比 |
mongodb-lab/scripts/ch27-bulkwrite-perf.js | bulkWrite 性能对比 |
mongodb-lab/scripts/ch27-connection-monitor.js | 连接池监控脚本 |
3.4 测试验证
use local_life// 1. 原子扣减验证——模拟并发扣库存db.seckill_products.updateOne({_id:"SKU_MOUTAI"},{$set:{stock:3}})// 连续发起 10 次扣减请求(只有前 3 次成功)letsuccessCount=0for(leti=0;i<10;i++){constr=db.seckill_products.updateOne({_id:"SKU_MOUTAI",stock:{$gt:0}},{$inc:{stock:-1}})if(r.modifiedCount>0)successCount++}print("尝试 10 次, 成功:",successCount,successCount===3?"PASS":"FAIL")// 2. 幂等验证——相同 userId+productId 不能重复插入db.seckill_records.deleteMany({})db.seckill_records.insertOne({userId:"IDEMPOTENT_1",productId:"SKU_MOUTAI",claimedAt:newDate()})try{db.seckill_records.insertOne({userId:"IDEMPOTENT_1",productId:"SKU_MOUTAI",claimedAt:newDate()})print("幂等验证: FAIL (允许重复插入)")}catch(e){print("幂等验证:",e.code===11000?"PASS (唯一冲突拦截)":"FAIL")}// 3. 连接池监控constconn=db.serverStatus().connectionsprint("连接:",conn.current,"| 可用:",conn.available,conn.current<conn.available*0.8?"PASS":"⚠ 接近上限")print("\n=== 高并发写入验证完成 ===")4. 项目总结
4.1 秒杀架构决策矩阵
| 场景 | 写入策略 | 连接池 | 幂等 | 削峰 |
|---|---|---|---|---|
| 普通商品下单 | 原子 $inc + filter | 默认 50 | 订单号唯一索引 | 无需 |
| 限时秒杀 | 原子 $inc + 快速返回 | 调大到 100 | userId+productId 唯一索引 | 前台扣库 + 后台 consumer |
| 批量消息发送 | bulkWrite (ordered:false) | 调大到 200 | messageId 唯一索引 | Kafka/Redis 队列 |
| 日志采集 | bulkWrite 单批次 500-1000 | 调大到 100 | 日志 ID 唯一索引或无需 | 批量攒批写入 |
4.2 适用场景
本章优化适用:
- 秒杀/抢购——原子扣库存 + 幂等记录 + 前台快速返回。
- 高并发的优惠券发放——原子库存扣减 + 兜底唯一冲突。
- 批量数据导入——bulkWrite 替代逐条 insert。
- 日志/事件流灌入——bulkWrite + 异步 consumer 削峰。
- 热点写入场景的专项优化——调优连接池、增大 opLog 窗口、开启写关注 majority。
4.3 注意事项
| 注意事项 | 说明 |
|---|---|
| 原子操作仅限于单文档 | 跨文档的原子性必须用事务(第 11 章) |
ordered: false的副作用 | 批量中的错误不会阻止其他操作执行——检查每条结果 |
| 唯一索引和正常索引不可混用 | 幂等键是唯一索引,不要在这个字段上再建普通索引 |
连接池的maxWaitTime设为 2s | 失败快速返回,由上游重试比排队等 2 分钟更健康 |
| bulkWrite 单批次不可太大 | 单批次 > 1000 条可能导致单次执行时间过长,影响复制延迟 |
4.4 常见踩坑经验
故障案例一:连接池 maxPoolSize=10 打满后服务雪崩
某秒杀服务部署了 20 个 Pod,每个 Pod 的 maxPoolSize=10。秒杀开始时 20×10=200 个连接全部被秒杀请求占满。其他正常业务请求(如查订单、查物流)也需要连接但连接池是独立的——问题在于秒杀服务的连接池满后,该 Pod 的健康检查也走 MongoDB 查询——健康检查拿不到连接超时,K8s 把 Pod 标记为 Unhealthy 重启,重启后瞬间又被打满——反复 CrashLoopBackOff。解决:① 秒杀服务单独部署、独立的 MongoDB 连接池;② 健康检查不走 MongoDB(用简单的 TCP 端口检查);③ maxWaitTime 设为 1 秒快速失败。
故障案例二:read-then-write 的库存超卖在分片集群中更严重
某团队用 read-then-write 做分片集群中的库存扣减。分片集群中 read 可能被路由到 Secondary(从库),write 必须去 Primary——从库读到的库存和主库的最新库存之间存在复制延迟,超卖比单机更严重。解决:必须用原子$inc + filter、writeConcern majority、readPreference primary 三位一体。
故障案例三:bulkWrite 单条异常导致整批回滚
某团队 bulkWrite 设置了ordered: true(默认),批量 1000 条中的第 500 条因为唯一键冲突失败,后 500 条全部没执行——但前 499 条已成功(已提交不回滚)。应用层以为"批量失败了",重新发送了这 1000 条——导致前 499 条被重复处理的副作用。解决:ordered: false+ 每个 write 结果检查 + 幂等键兜底。批量异常不代表全部失败。
4.5 思考题
- 在分片集群中做秒杀,如果分片键不是
productId,updateOne({_id: productId, stock: {$gt: 0}}, {$inc: {stock: -1}})能否精确路由到单个分片?如果不能,会有什么后果? - 如果连接池的
maxWaitTime设为 0(永远等待),在高并发下会发生什么?
(答案将在第 28 章末尾揭晓)
上一章思考题答案:
删除文档不会自动减少 Chunk 数——Chunk 划分基于分片键的值范围,而非文档数量。即使一个 Chunk 内的文档被清空,这个 Chunk 仍然存在(成为空 Chunk),Chunk 数不变。MongoDB 会定期合并小的空 Chunk,但这是后台操作而非实时。
reshardCollection如果在执行到一半时 mongos 挂了——重启后 reshard 操作不会自动恢复,需要手动重新发起。但旧集合的数据完好无损——reshard 在内部创建了一个新的目标集合,原子切换发生在最后。如果切换前失败,旧集合不受任何影响;如果切换后失败,新集合已经生效,数据完整。可通过sh.status()查看残留的 resharding 操作。
延伸阅读与资源
MongoDB 实战进阶与内核修炼
python入门:Rquests从菜鸟脚本到企业级SDK的网络实战圣经
Milvus向量数据库实战修炼:从 0 到 1精通向量检索与生产落地
后端工程师的 AI 转型第一课:Ollama 与私有化大模型实战
10倍开发者的 Dify 魔法书:从零构建全栈 AI 应用
后端工程师转型AI第一课-Ollama 与私有化大模型实战
大型语言模型(LLM) vLLM 高性能推理落地实战
Agent开发之LlamaIndex 实战修炼与源码进阶
大语言模型Transformers 实战修炼与源码剖析