Bull队列:如何在单个队列中触发多个函数执行
实现基于Bull队列的HTTP触发任务方案
核心思路
要实现不同HTTP接口触发对应队列任务,关键是在queue.add时传递任务类型标识和业务参数,再在queue.process中根据类型分发到对应函数执行。
完整实现步骤
1. 安装依赖
先安装Bull队列和Express HTTP框架:
npm install bull express
2. 完整代码示例
import Queue from "bull"; import express from "express"; const app = express(); app.use(express.json()); // 解析JSON请求体 // 创建MyQueue队列实例 const myQueue = new Queue("MyQueue"); // 定义需要执行的业务函数 const func1 = async (params) => { console.log("执行func1,参数:", params); // 这里编写func1的业务逻辑 return `func1处理完成,参数:${JSON.stringify(params)}`; }; const func2 = async (params) => { console.log("执行func2,参数:", params); // 这里编写func2的业务逻辑 return `func2处理完成,参数:${JSON.stringify(params)}`; }; const func3 = async (params) => { console.log("执行func3,参数:", params); // 这里编写func3的业务逻辑 return `func3处理完成,参数:${JSON.stringify(params)}`; }; // 队列任务处理器:根据任务类型分发到对应函数 myQueue.process(async (job) => { const { type, params } = job.data; switch (type) { case "func1": return await func1(params); case "func2": return await func2(params); case "func3": return await func3(params); default: throw new Error(`未知任务类型:${type}`); } }); // 定义HTTP接口,触发对应队列任务 app.post("/func1", async (req, res) => { try { const job = await myQueue.add({ type: "func1", params: req.body // 接收请求中的业务参数 }); res.status(200).json({ message: "func1任务已加入队列", jobId: job.id }); } catch (err) { res.status(500).json({ error: err.message }); } }); app.post("/func2", async (req, res) => { try { const job = await myQueue.add({ type: "func2", params: req.body }); res.status(200).json({ message: "func2任务已加入队列", jobId: job.id }); } catch (err) { res.status(500).json({ error: err.message }); } }); app.post("/func3", async (req, res) => { try { const job = await myQueue.add({ type: "func3", params: req.body }); res.status(200).json({ message: "func3任务已加入队列", jobId: job.id }); } catch (err) { res.status(500).json({ error: err.message }); } }); // 启动HTTP服务 const PORT = 3000; app.listen(PORT, () => { console.log(`服务运行在 http://localhost:${PORT}`); });
关键细节说明
- 任务标识传递:通过
queue.add的data字段同时传递type(任务类型)和params(业务参数),让处理器能精准匹配执行函数。 - 异步处理优化:使用
async/await替代传统done回调,代码更简洁易维护,同时支持任务结果返回。 - 参数接收:HTTP接口通过
req.body获取请求携带的参数,直接传入队列任务。
测试方式
用curl或Postman发送POST请求即可测试:
# 测试func1 curl -X POST http://localhost:3000/func1 -H "Content-Type: application/json" -d '{"name": "John", "age": 30}' # 测试func2 curl -X POST http://localhost:3000/func2 -H "Content-Type: application/json" -d '{"orderId": "12345", "amount": 99.9}'
内容的提问来源于stack exchange,提问作者user1862965
相关产品推荐
相关产品推荐

