多服务器环境下限时订单锁定的定时器持久化方案咨询
分布式环境下限量库存锁定方案优化
核心思路转变
把单服务器内存中的定时器逻辑,迁移到共享持久化存储+分布式状态管理,解决跨服务器调度和服务器崩溃丢数据的问题。主要有两种可行方向:基于数据库的定时扫描,或者用分布式延迟任务队列。
方案一:数据库状态标记+定时扫描(无额外依赖)
这是最轻量化的方案,适合中小规模电商,不需要引入第三方中间件。
实现步骤
- 数据库表设计:新增库存锁定记录表,追踪锁定状态、过期时间;拆分商品库存为可用、锁定、已售出三个字段,避免直接修改可用库存导致并发问题。
- 购买流程:开启事务,检查可用库存→扣减可用库存并增加锁定库存→创建锁定记录(含过期时间)。
- 过期恢复:用分布式定时任务(如Node.js的
node-schedule、Java的Quartz)周期性扫描过期且未处理的锁定记录,执行库存恢复并更新状态。 - 支付/取消流程:更新锁定记录状态为已支付/已取消,同时调整库存(锁定转售出,或锁定转可用)。
关键代码示例
数据库表结构(MySQL)
CREATE TABLE inventory_locks ( id INT AUTO_INCREMENT PRIMARY KEY, buyer_address VARCHAR(255) NOT NULL, product_id INT NOT NULL, locked_quantity INT NOT NULL, expire_at DATETIME NOT NULL, status ENUM('LOCKED', 'PAID', 'EXPIRED', 'CANCELED') NOT NULL DEFAULT 'LOCKED', created_at DATETIME DEFAULT CURRENT_TIMESTAMP, UNIQUE KEY idx_buyer_product (buyer_address, product_id) ); ALTER TABLE products ADD COLUMN locked_stock INT DEFAULT 0, ADD COLUMN sold_stock INT DEFAULT 0;
购买操作(Node.js)
async function handleBuy(buyerAddress, productId, quantity) { const tx = await db.beginTransaction(); try { // 加行锁避免并发下单 const [product] = await tx.query('SELECT available_stock FROM products WHERE id = ? FOR UPDATE', [productId]); if (!product || product.available_stock < quantity) { throw new Error('库存不足'); } // 更新库存状态 await tx.query( 'UPDATE products SET available_stock = available_stock - ?, locked_stock = locked_stock + ? WHERE id = ?', [quantity, quantity, productId] ); // 创建锁定记录 const expireAt = new Date(Date.now() + 75000); await tx.query( 'INSERT INTO inventory_locks (buyer_address, product_id, locked_quantity, expire_at) VALUES (?, ?, ?, ?)', [buyerAddress, productId, quantity, expireAt] ); await tx.commit(); } catch (err) { await tx.rollback(); throw err; } }
过期扫描任务
const schedule = require('node-schedule'); // 每10秒执行一次过期扫描 schedule.scheduleJob('*/10 * * * * *', async () => { const now = new Date(); const [expiredLocks] = await db.query( 'SELECT * FROM inventory_locks WHERE status = ? AND expire_at < ?', ['LOCKED', now] ); for (const lock of expiredLocks) { // 恢复库存 await db.query( 'UPDATE products SET available_stock = available_stock + ?, locked_stock = locked_stock - ? WHERE id = ?', [lock.locked_quantity, lock.locked_quantity, lock.product_id] ); // 更新锁定记录状态 await db.query('UPDATE inventory_locks SET status = ? WHERE id = ?', ['EXPIRED', lock.id]); } });
方案二:分布式延迟任务队列(适合高并发场景)
当并发量较高时,定时扫描可能存在性能瓶颈,此时可以用基于Redis的分布式延迟任务队列,精准控制任务执行时间,且支持跨服务器取消任务。
常用工具
- BullMQ(Node.js):轻量级分布式任务队列,支持延迟任务、任务取消,基于Redis持久化。
- Redisson(Java):提供分布式延迟队列实现,自带分布式锁。
- Celery(Python):通过Redis/RabbitMQ实现延迟任务。
实现示例(BullMQ)
初始化队列与Worker
const { Queue, Worker } = require('bullmq'); // 连接Redis const lockQueue = new Queue('inventory-lock', { connection: { host: 'localhost', port: 6379 } }); // 处理延迟任务的Worker const worker = new Worker('inventory-lock', async (job) => { const { buyerAddress, productId, quantity } = job.data; // 幂等检查:确认锁定状态仍为有效 const [lock] = await db.query( 'SELECT status FROM inventory_locks WHERE buyer_address = ? AND product_id = ?', [buyerAddress, productId] ); if (lock?.status !== 'LOCKED') return; // 执行恢复逻辑 await db.query( 'UPDATE products SET available_stock = available_stock + ?, locked_stock = locked_stock - ? WHERE id = ?', [quantity, quantity, productId] ); await db.query('UPDATE inventory_locks SET status = ? WHERE buyer_address = ? AND product_id = ?', ['EXPIRED', buyerAddress, productId]); });
购买时创建延迟任务
async function handleBuy(buyerAddress, productId, quantity) { const tx = await db.beginTransaction(); try { // 库存检查与更新(同方案一) const [product] = await tx.query('SELECT available_stock FROM products WHERE id = ? FOR UPDATE', [productId]); if (!product || product.available_stock < quantity) { throw new Error('库存不足'); } await tx.query( 'UPDATE products SET available_stock = available_stock - ?, locked_stock = locked_stock + ? WHERE id = ?', [quantity, quantity, productId] ); // 创建75秒后执行的延迟任务,用买家+商品ID作为任务ID const job = await lockQueue.add( 'restore-lock', { buyerAddress, productId, quantity }, { delay: 75000, jobId: `${buyerAddress}-${productId}` } ); // 存储任务ID到锁定记录 const expireAt = new Date(Date.now() + 75000); await tx.query( 'INSERT INTO inventory_locks (buyer_address, product_id, locked_quantity, expire_at, job_id) VALUES (?, ?, ?, ?, ?)', [buyerAddress, productId, quantity, expireAt, job.id] ); await tx.commit(); } catch (err) { await tx.rollback(); throw err; } }
支付成功时取消任务
async function handlePaymentSuccess(buyerAddress, productId) { const tx = await db.beginTransaction(); try { const [lock] = await db.query( 'SELECT job_id, locked_quantity FROM inventory_locks WHERE buyer_address = ? AND product_id = ? AND status = ? FOR UPDATE', [buyerAddress, productId, 'LOCKED'] ); if (!lock) throw new Error('无效锁定记录'); // 取消延迟任务 await lockQueue.remove(lock.job_id); // 确认库存状态更新 await tx.query( 'UPDATE products SET locked_stock = locked_stock - ?, sold_stock = sold_stock + ? WHERE id = ?', [lock.locked_quantity, lock.locked_quantity, productId] ); await tx.query('UPDATE inventory_locks SET status = ? WHERE id = ?', ['PAID', lock.id]); await tx.commit(); } catch (err) { await tx.rollback(); throw err; } }
关键注意事项
- 事务与行锁:所有库存操作必须包裹在数据库事务中,且查询库存时加行锁(如
SELECT ... FOR UPDATE),避免超卖。 - 幂等性检查:延迟任务执行前必须再次验证锁定状态,防止支付成功后任务仍被执行。
- 任务持久化:分布式任务队列依赖Redis等持久化存储,确保服务器崩溃后任务不会丢失。
内容的提问来源于stack exchange,提问作者berger
相关产品推荐
相关产品推荐

