如何在Node.js中集成job queue实现定时价格监控下单功能?
Node.js 实现定时价格监控任务队列方案
核心思路
你的需求完全可以基于Bull实现,核心是利用重复任务+状态控制处理每小时查价逻辑,通过任务元数据隔离不同用户请求,结合Redis实现任务状态持久化,确保并发场景下的正确性。
具体实现步骤
1. 初始化Bull队列
Bull依赖Redis存储任务数据,先完成基础配置:
const Queue = require('bull'); const redis = require('redis'); // 初始化Redis客户端(用于存储任务状态) const redisClient = redis.createClient({ host: 'localhost', port: 6379 }); redisClient.connect(); // 初始化价格监控队列 const priceMonitorQueue = new Queue('price-monitor', { redis: { host: 'localhost', port: 6379 } });
2. 为用户创建监控任务
每个用户的监控任务需包含唯一标识、商品信息、目标价格及下单/退款参数,同时设置重复规则和过期时间:
async function createPriceMonitorTask(userId, productId, targetPrice, orderParams, refundParams) { const expireTime = Date.now() + 5 * 24 * 60 * 60 * 1000; // 5天后过期 const uniqueJobId = `monitor-${userId}-${productId}`; // 避免重复添加同一用户同一商品的任务 const existingJob = await priceMonitorQueue.getJob(uniqueJobId); if (existingJob) return existingJob; return await priceMonitorQueue.add( { userId, productId, targetPrice, orderParams, refundParams, expireTime }, { jobId: uniqueJobId, repeat: { every: 60 * 60 * 1000 }, // 每小时执行一次 attempts: 3, // 任务失败时重试3次 backoff: { type: 'exponential', delay: 5000 } // 指数退避重试 } ); }
3. 任务处理器核心逻辑
实现价格检查、下单、退款逻辑,通过状态控制避免重复执行:
// 设置并发处理数(根据服务器资源调整) priceMonitorQueue.process(10, async (job) => { const { userId, productId, targetPrice, orderParams, refundParams, expireTime } = job.data; const jobStatusKey = `job-status:${job.id}`; // 检查任务是否已完成/退款,避免重复处理 const currentStatus = await redisClient.hGet(jobStatusKey, 'status'); if (currentStatus === 'completed' || currentStatus === 'refunded') { await job.removeRepeatable(); // 取消后续重复任务 return; } // 检查任务是否已过期 if (Date.now() > expireTime) { // 执行退款逻辑 await processRefund(refundParams); await redisClient.hSet(jobStatusKey, 'status', 'refunded'); await job.removeRepeatable(); return; } try { // 获取当前商品价格(替换为你的价格API调用) const currentPrice = await fetchProductPrice(productId); if (currentPrice <= targetPrice) { // 执行下单逻辑 await processOrder(orderParams); await redisClient.hSet(jobStatusKey, 'status', 'completed'); await job.removeRepeatable(); // 取消后续检查 } else { // 未达标,标记为待下次检查 await redisClient.hSet(jobStatusKey, 'status', 'pending'); } } catch (err) { // 价格查询失败,抛出错误让Bull重试 throw new Error(`Price fetch failed: ${err.message}`); } }); // 模拟价格查询函数 async function fetchProductPrice(productId) { // 替换为真实的商品价格API调用 return await fetch(`https://your-price-api.com/products/${productId}`) .then(res => res.json()) .then(data => data.price); } // 模拟下单函数 async function processOrder(orderParams) { // 替换为真实的下单逻辑 console.log(`Executing order for user ${orderParams.userId}:`, orderParams); } // 模拟退款函数 async function processRefund(refundParams) { // 替换为真实的退款逻辑 console.log(`Executing refund for user ${refundParams.userId}:`, refundParams); }
4. 并发与用户隔离
- 通过
process方法的第一个参数设置并发数(如上例的10),支持同时处理多个用户的任务。 - 每个任务用
jobId(用户ID+商品ID)唯一标识,避免同一用户重复添加任务。 - 任务元数据(用户ID、下单参数等)完全隔离,处理器仅处理当前任务对应的用户逻辑。
关键注意事项
- 幂等性:下单和退款逻辑必须实现幂等,避免因任务重试或重复执行导致重复操作。
- Redis持久化:配置Redis持久化(RDB/AOF),避免服务器重启后任务丢失。
- 任务清理:通过
removeOnComplete和removeOnFail参数自动清理已完成/失败的任务,减少Redis存储占用。 - 异常处理:对价格API、下单/退款接口的异常进行捕获,利用Bull的重试机制保证任务可靠性。
内容的提问来源于stack exchange,提问作者Tahir Mirza
相关产品推荐
相关产品推荐

