MongoDB中如何处理多用户同时下单同一款商品的并发问题
并发下单场景的串行处理与事务优化方案
问题分析
你当前的代码通过MongoDB事务实现了库存扣减与订单创建的原子性,但并发下单时,第二个请求会因库存检查不通过(第一个请求已扣减库存)直接抛出"Out of stock"错误并中止事务。这是因为MongoDB的bulkWrite虽然是原子操作,但并发请求会同时触发库存检查,后到的请求会因库存不足直接失败。你需要在应用层实现请求串行化,让同一商品的下单请求按顺序处理,避免并发竞争。
核心解决方案
1. 引入锁机制实现请求排队
- 单实例部署:使用内存锁(如
Map存储锁状态),确保同一商品的下单请求串行执行。 - 多实例集群:使用分布式锁(如Redis的SETNX命令),跨实例同步锁状态,避免锁失效。
2. 优化事务执行顺序
将Promise.all并行执行的操作调整为串行执行,先完成库存扣减,再执行清空购物车、创建订单的操作,逻辑更清晰,也便于排查问题。
代码修改示例
方案一:单实例本地锁实现
// 全局维护商品锁,key为商品ID,value为Promise(用于请求排队) const productLocks = new Map<string, Promise<void>>(); const placeOrder = async (req: NextApiRequest, res: NextApiResponse<Data>) => { let session: ClientSession | null = null; let productLockRelease: (() => void) | null = null; let targetProductId: string | null = null; try { if (req.headers.isauth === '0') { return handleError(req, res, { code: 401, message: 'unAuthorized' }); } const { address } = req.body as Omit<Address, 'defaultAdd'> & { defaultAdd?: boolean; }; const user = await User.findOne({ _id: req.headers._id }); if (!user) { return handleError(req, res, { code: 404, message: 'user not found' }); } // 取购物车中第一个商品ID(若购物车多商品,需扩展为多锁逻辑) targetProductId = user.cart[0].id; // 等待已有锁释放 if (productLocks.has(targetProductId)) { await productLocks.get(targetProductId); } // 创建新锁 let resolveLock: () => void; const lockPromise = new Promise<void>(resolve => { resolveLock = resolve; }); productLocks.set(targetProductId, lockPromise); productLockRelease = resolveLock; // 启动事务 session = await startSession(); session.startTransaction(); // 构建库存扣减操作 const bulkOps = user.cart.map((item: Cart) => ({ updateOne: { filter: { _id: item.id, inStock: { $gte: item.quantity }, }, update: { $inc: { inStock: -item.quantity, sold: +item.quantity } }, }, })); // 串行执行操作:先扣减库存 const productResponse = await Product.bulkWrite(bulkOps, { session }); if (productResponse.result.nModified < bulkOps.length) { throw new Error('Out of stock'); } // 再清空购物车 await user.update({ cart: [] }, { session }); // 最后创建订单 await Order.create( [ { userId: user._id.toString(), shippingFee: 15000, // fixed price address, products: user.cart, }, ], { session }, ); // 提交事务 await session.commitTransaction(); return res.status(200).send({ message: 'ok' }); } catch (err) { if (session) { await session.abortTransaction(); } return handleError(req, res, err instanceof Error ? { message: err.message } : {}); } finally { // 释放锁 if (productLockRelease && targetProductId) { productLockRelease(); productLocks.delete(targetProductId); } // 关闭会话 if (session) { session.endSession(); } } };
方案二:Redis分布式锁实现(多实例集群场景)
import redis from 'redis'; // 初始化Redis客户端 const redisClient = redis.createClient({ /* 你的Redis连接配置 */ }); await redisClient.connect(); // 获取分布式锁 const acquireLock = async (productId: string, timeout = 10000): Promise<boolean> => { const lockKey = `lock:product:${productId}`; // SETNX + EXPIRE 原子操作,避免死锁 const result = await redisClient.set(lockKey, '1', { NX: true, EX: timeout, }); return result === 'OK'; }; // 释放分布式锁 const releaseLock = async (productId: string) => { const lockKey = `lock:product:${productId}`; await redisClient.del(lockKey); }; const placeOrder = async (req: NextApiRequest, res: NextApiResponse<Data>) => { let session: ClientSession | null = null; let targetProductId: string | null = null; try { if (req.headers.isauth === '0') { return handleError(req, res, { code: 401, message: 'unAuthorized' }); } const { address } = req.body as Omit<Address, 'defaultAdd'> & { defaultAdd?: boolean; }; const user = await User.findOne({ _id: req.headers._id }); if (!user) { return handleError(req, res, { code: 404, message: 'user not found' }); } targetProductId = user.cart[0].id; // 尝试获取锁,最多重试3次 let lockAcquired = false; for (let i = 0; i < 3; i++) { lockAcquired = await acquireLock(targetProductId); if (lockAcquired) break; await new Promise(resolve => setTimeout(resolve, 500)); // 间隔500ms重试 } if (!lockAcquired) { return handleError(req, res, { code: 429, message: '请求过于频繁,请稍后再试' }); } // 启动事务 session = await startSession(); session.startTransaction(); const bulkOps = user.cart.map((item: Cart) => ({ updateOne: { filter: { _id: item.id, inStock: { $gte: item.quantity }, }, update: { $inc: { inStock: -item.quantity, sold: +item.quantity } }, }, })); const productResponse = await Product.bulkWrite(bulkOps, { session }); if (productResponse.result.nModified < bulkOps.length) { throw new Error('Out of stock'); } await user.update({ cart: [] }, { session }); await Order.create( [ { userId: user._id.toString(), shippingFee: 15000, address, products: user.cart, }, ], { session }, ); await session.commitTransaction(); return res.status(200).send({ message: 'ok' }); } catch (err) { if (session) { await session.abortTransaction(); } return handleError(req, res, err instanceof Error ? { message: err.message } : {}); } finally { // 释放锁 if (targetProductId) { await releaseLock(targetProductId); } // 关闭会话与Redis连接 if (session) { session.endSession(); } await redisClient.quit(); } };
关键说明
- 锁粒度控制:按商品ID加锁,避免全局锁导致所有请求串行,影响系统吞吐量。
- 死锁预防:本地锁必须在
finally块中释放;分布式锁需设置超时时间,即使服务崩溃也能自动释放锁。 - 重试机制:分布式锁场景下添加有限次数重试,可提升用户体验,避免直接返回失败。
内容的提问来源于stack exchange,提问作者Simba
相关产品推荐
相关产品推荐

