发布于 2025-12-10 2 阅读
0

使用 bee-queue 和 redis 的简单 Node.js 任务队列

使用 bee-queue 和 redis 的简单 Node.js 任务队列

Unsplash上的封面照片由Bimo Luki拍摄

正如您在上一篇文章中看到的那样,任务队列非常棒🌟,在本教程中,我们将在我们自己的应用程序中使用任务队列,让我们开始编写一些代码。

我们将按照上一篇文章所述建造我们的餐厅。

本教程主要是一个演示,而不是一个可以运行的应用程序,所以如果你想了解如何在你的应用中插入任务队列,请继续关注我。
下一篇文章我们将构建一个真正的应用程序。(我知道这很令人兴奋,你肯定迫不及待地想看😉)。

👨‍💻 该项目的完整 GitHub 仓库链接位于文章末尾⬇

让我们开始吧。

先决条件

  • 在您的机器上安装Node.js,然后运行以下命令来验证安装是否正确。```bash

$ 节点--版本

v12.16.1

+ Redis running on your pc or the cloud. Install [Redis](https://redis.io/) or create an instance on [RedisLabs](https://redislabs.com/) for free.

And we're good to go :grin:

## Initialization
Run:
```bash


$ npm init


Enter fullscreen mode Exit fullscreen mode

之后运行以下命令安装必要的软件包



$ npm install express bee-queue dotenv


Enter fullscreen mode Exit fullscreen mode

如果你想知道每个包的作用,这里有一些信息:

  • express帮助我们创建服务器并轻松处理传入的请求。
  • bee-queue是我们的任务队列管理器,将帮助创建和运行作业
  • dotenv.env帮助我们从本地文件加载环境变量

之后创建一个文件restaurant.js并编辑它package.json,使其看起来像这样



{
  ...
  "main": "restaurant.js",
  "scripts": {
    "start": "node restaurant.js"
  }
  ...
}


Enter fullscreen mode Exit fullscreen mode

是时候编写一些真正的代码了

在您选择的编辑器中打开restaurant.js并添加以下代码行



require('dotenv').config();
const express = require('express');
const http = require('http');

// Inits
const app = express();
app.use(express.json());
app.use(express.urlencoded({ extended: false }));

// Routes
app.get('/', (req, res) => {
    res.send("😋 We are serving freshly cooked food 🍲");
});


// Create and start the server
const server = http.createServer(app);
const PORT = process.env.PORT || 5000;
server.listen(PORT, () => {
    console.log(`Restaurant open at:${PORT}`);
});


Enter fullscreen mode Exit fullscreen mode

它的作用基本上是在指定端口(这里是 5000)启动本地网络服务器,并GET在基本 URL 上监听传入请求/并用文本回复。

运行以下命令启动服务器并转到localhost:5000您的浏览器。



$ npm start
> restaurant@1.0.0 start /mnt/code/dev/queue
> node restaurant.js

Restaurant open at port:5000


Enter fullscreen mode Exit fullscreen mode

您将获得一张空白页,上面有一条简洁的小😋 We are serving freshly cooked food 🍲信息


现在是时候创建我们的任务队列了

首先创建一个名为的文件.env,并在其中粘贴您的数据库凭据,如下所示(您也可以在这里使用本地 redis 实例),并记住,永远不要提交.env到您的源代码控制。



DB_HOST=redis-random-cloud.redislabs.com
DB_PORT=14827
DB_PASS=pTAl.not-my-password.rUlJq


Enter fullscreen mode Exit fullscreen mode

您已完成基本配置。

让我们继续创建我们的waiter。首先创建一个文件waiter.js并添加以下代码块:




const Queue = require('bee-queue');

const options = {
    removeOnSuccess: true,
    redis: {
        host: process.env.DB_HOST,
        port: process.env.DB_PORT,
        password: process.env.DB_PASS,
    },
}

const cookQueue = new Queue('cook', options);
const serveQueue = new Queue('serve', options);


const placeOrder = (order) => {
    return cookQueue.createJob(order).save();
};

serveQueue.process((job, done) => {
    console.log(`🧾 ${job.data.qty}x ${job.data.dish} ready to be served 😋`);
    // Notify the client via push notification, web socket or email etc.
    done();
})
    // Notify the client via push notification, web socket or email etc.
    done();
})


module.exports.placeOrder = placeOrder;


Enter fullscreen mode Exit fullscreen mode

🤯 哇!那是什么?好吧,让我解释一下。

我们首先将bee-queue包导入为Queue
然后将数据库配置传递给两个新Queue对象。其中一个队列将包含厨师准备的订单列表,另一个队列将包含服务员已准备好的订单列表。

然后,我们创建一个新的函数placeOrder,接受一个order作为参数。我们稍后会定义这个订单对象,但请记住它的结构如下:



order = {
    dish: "Pizza 🍕", 
    qty: 2,
    orderNo: "kbv9euic"
}


Enter fullscreen mode Exit fullscreen mode

该函数接收此订单,并通过调用 Queue 对象的方法placeOrder将其添加到队列中。它充当任务发布者.createJob(order).save()cookQueue

最后,每当订单准备好并准备服务时, Queue 对象process上的方法serveQueue都会执行处理函数。这充当了任务消费者的角色。(job, done) => {...}

我们调用done()来确认任务队列中的作业已完成,以便它可以从队列中发送下一个待处理的任务。我们只需调用done()来指示任务已成功,并done(err)使用第一个参数(其中err为错误消息)调用 来指示作业失败。您也可以调用 来指示作业done(null, msg)成功,并使用第二个参数msg(其中 为成功消息)来指示作业成功。

我们的服务员👨‍💼已经准备好了


现在是时候让厨师们去厨房了👨‍🍳
创建另一个文件kitchen.js并在其中粘贴以下代码行:



const Queue = require('bee-queue');

const options = {
    removeOnSuccess: true,
    redis: {
        host: process.env.DB_HOST,
        port: process.env.DB_PORT,
        password: process.env.DB_PASS,
    },
}

const cookQueue = new Queue('cook', options);
const serveQueue = new Queue('serve', options);

cookQueue.process(3, (job, done) => {
    setTimeout(() => console.log("Getting the ingredients ready 🥬 🧄 🧅 🍄"), 1000);
    setTimeout(() => console.log(`🍳 Preparing ${job.data.dish}`), 1500);
    setTimeout(() => {
        console.log(`🧾 Order ${job.data.orderNo}: ${job.data.dish} ready`);
        done();
    }, job.data.qty * 5000);
});

cookQueue.on('succeeded', (job, result) => {
    serveQueue.createJob(job.data).save();
});


Enter fullscreen mode Exit fullscreen mode

😌 嗯,看起来很熟悉。

是的,但唯一的变化是,我们的厨师从那里消费cookQueue,然后发布到那里,serveQueue以便服务员接受并处理订单。

这里需要注意的是,通过 发布的所有内容都可以在方法的处理函数createJob(order)供消费者使用,但如果仔细观察,还会发现一些不同之处。没错,我们在实际的处理函数之前传入了一个数字。它被称为并发数(队列中可以同时处理的任务数)。这里我们将其设置为 3,因为我们的厨房有 3 个厨师,他们可以一起工作。job.dataQueue.process()(job, done) => {...}cookQueue.process(3, (job, done) => {...})

cookQueue.on('succeeded', (job, result) => {...})每当任务成功时(即每当您调用done()process()方法时),我们就使用该方法调用处理程序函数。

相信我,我们快完成了🤞


最后一步:把所有东西连接起来

打开restaurant.js并添加以下最后几行代码



// ...
// Add these lines before the Inits.
require('./kitchen');
const { placeOrder } = require('./waiter');

// Inits
// ...
// Routes

// ...

app.post('/order', (req, res) => {
    let order = {
        dish: req.body.dish,
        qty: req.body.qty,
        orderNo: Date.now().toString(36)
    }

    if (order.dish && order.qty) {
        placeOrder(order)
            .then(() => res.json({ done: true, message: "Your order will be ready in a while" }))
            .catch(() => res.json({ done: false, message: "Your order could not be placed" }));
    } else {
        res.status(422);
    }
})

// Create and start the server
// ...


Enter fullscreen mode Exit fullscreen mode

我们在这里所做的是导入我们的kitchenwaiter并添加了一个 POST 路由/order来接收来自客户的订单。还记得订单对象吗?



order = {
    dish: "Pizza 🍕", 
    qty: 2,
    orderNo: "kbv9euic"
}


Enter fullscreen mode Exit fullscreen mode

我们正在根据 POST 请求的 JSON 主体创建一个订单对象,并将其传递给服务员,并发送 JSON 响应来确认客户。如果请求不正确,我们还会发送一些错误消息。这样就完成了✌。


是啊,真的完成了。现在该测试一下了😁

  • $ npm start通过在终端上运行来启动服务器。
  • 发送一个 get 请求localhost:5000并查看是否收到如下响应:餐厅营业
  • 接下来发送一个 POST 请求localhost:5000/order并检查响应并查看您的控制台。测试 API

您可以依次发送多个请求,以检查它不会挂起任何请求。

让我们添加另一条POST路线,将其与没有任务队列的普通餐厅进行比较。

将这些行添加到restaurant.js



//  ...
app.post('/order-legacy', (req, res) => {
    let order = {
        dish: req.body.dish,
        qty: req.body.qty,
        orderNo: Date.now().toString(36)
    }
    if (order.dish && order.qty) {
        setTimeout(() => console.log("Getting the ingredients ready... 🥬 🧄 🧅 🍄"), 1000);
        setTimeout(() => console.log(`🍳 Preparing ${order.dish}`), 1500);
        setTimeout(() => {
            console.log(`🧾 Order ${order.orderNo}: ${order.dish} ready`);
            res.json({ done: true, message: `Your ${order.qty}x ${order.dish} is ready` })
        }, order.qty * 5000);
    } else {
        console.log("Incomplete order rejected");
        res.status(422).json({ done: false, message: "Your order could not be placed" });
    }
});


// Create and start the server
// ...


Enter fullscreen mode Exit fullscreen mode
  • 接下来发送一个 POST 请求localhost:5000/order-legacy并检查响应并查看您的控制台。替代文本

注意响应时间的差异🤯

使用任务队列
带有任务队列

无任务队列
无任务队列


这是 Github repo,包含完整的项目

GitHub 徽标 sarbikbetal / nodejs-任务队列

本 repo 包含文章“使用 bee-queue 和 redis 的简单 Node.js 任务队列”的示例代码

如果您有任何问题或建议,请在下面发表评论,并随时与我联系😄并查看下面的问答部分。

📸 Instagram 📨电子邮件 👨‍💼 LinkedIn Github

🤔 嗯...不过我有一些问题。

我知道,所以这里有一些常见的问题,欢迎在下面的评论部分提问。

  • 食物做好后,我们如何将其发送给顾客?

    为此,我们需要在服务器端和客户端应用程序中实现一些额外的逻辑。例如,我们可以通过 Websockets、推送通知、电子邮件等方式来实现。不用担心,我会在下一篇文章中详细介绍这些内容。

  • 难道没有像 RabbitMQ 这样的更好的东西吗?

    是的,当然有,但是对于不需要很多高级功能但仍想维护良好的后端基础设施的小规模项目来说,RabbitMQ 就有点大材小用了,而 bee-queue 可能只是简单易用。

鏂囩珷鏉ユ簮锛�https://dev.to/sarbikbetal/simple-node-js-task-queue-with-bee-queue-and-redis-105b