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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 15:40:22