You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

多服务器环境下限时订单锁定的定时器持久化方案咨询

分布式环境下限量库存锁定方案优化

核心思路转变

把单服务器内存中的定时器逻辑,迁移到共享持久化存储+分布式状态管理,解决跨服务器调度和服务器崩溃丢数据的问题。主要有两种可行方向:基于数据库的定时扫描,或者用分布式延迟任务队列。

方案一:数据库状态标记+定时扫描(无额外依赖)

这是最轻量化的方案,适合中小规模电商,不需要引入第三方中间件。

实现步骤

  • 数据库表设计:新增库存锁定记录表,追踪锁定状态、过期时间;拆分商品库存为可用、锁定、已售出三个字段,避免直接修改可用库存导致并发问题。
  • 购买流程:开启事务,检查可用库存→扣减可用库存并增加锁定库存→创建锁定记录(含过期时间)。
  • 过期恢复:用分布式定时任务(如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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.07 01:29:53