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

基于RabbitMQ的Node.js多发布订阅者架构及Docker部署求助

问题描述

已实现基于Node.js的Publisher和Subscriber服务,通过RabbitMQ完成加法、乘法数学任务的分发与处理,REST接口调用正常。但需实现以下需求:

  • 创建2个状态不同的Publisher实例、3个Sub-worker/Subscriber实例
  • 工作者完成任务后向发布者反馈结果
  • 发布者在本地临时文件记录所有历史任务及其结果或待处理状态
  • 用Docker创建上述多实例

现有代码

Publisher.js

const express = require("express");
const amqp = require("amqplib");
const app = express();
const bodyParser = require("body-parser");
const PORT = process.env.PORT || 3000;
let channel, connection;
app.use(express.json());

app.get("/math-task/sum", (req, res) => {
  let inputOfA = parseInt(req.body.a);
  let inputOfB = parseInt(req.body.b);
  let sum = Number(inputOfA + inputOfB);
  sendData(sum); // pass the data to the function we defined
  console.log("A message is sent to queue");
  res.send("Message Sent For Addition:" + Number(sum)); //response to the API request
});

app.get("/math-task/mul", (req, res) => {
  let inputOfA = parseInt(req.body.a);
  let inputOfB = parseInt(req.body.b);
  let product = Number(inputOfA * inputOfB);
  sendData(product); // pass the data to the function we defined
  console.log("A message is sent to queue");
  res.send("Message Sent For Multiplication:" + Number(product)); //response to the API request
});

app.use(bodyParser.urlencoded({extended:false}));
app.use(bodyParser.json()); 

app.listen(PORT, () => console.log("Server running at port " + PORT));

async function connectQueue() {
  try {
    connection = await amqp.connect("amqp://localhost:5672");
    channel = await connection.createChannel();

    await channel.assertQueue("test-queue");
  } catch (error) {
    console.log(error);
  }
}

async function sendData(data) {
  // send data to queue
  await channel.sendToQueue("test-queue", Buffer.from(JSON.stringify(data)));

  // close the channel and connection
  await channel.close();
  await connection.close();
}

connectQueue();

Subscriber.js

const express = require("express");
const app = express();
const PORT = process.env.PORT || 3001;
app.use(express.json());
app.listen(PORT, () => console.log("Server running at port " + PORT));

const amqp = require("amqplib");
var channel, connection;
connectQueue()  // call the connect function
 
async function connectQueue() {
    try {
        connection = await amqp.connect("amqp://localhost:5672");
        channel    = await connection.createChannel()
        
        await channel.assertQueue("test-queue")
        
        channel.consume("test-queue", data => {
            console.log(`${Buffer.from(data.content)}`);
            channel.ack(data);
        })
    } catch (error) {
        console.log(error);
    }
}

解决方案

1. 多实例创建(Publisher/Subscriber)

Publisher多实例(状态区分)

通过环境变量区分实例状态,设置PUBLISHER_ID标识不同实例,同时指定不同端口避免冲突。启动命令:

# 实例1
PUBLISHER_ID=pub1 PORT=3000 node Publisher.js
# 实例2
PUBLISHER_ID=pub2 PORT=3001 node Publisher.js

代码中可通过process.env.PUBLISHER_ID获取标识,用于任务记录区分。

Subscriber多实例

启动时指定不同端口即可,若需标识可添加SUBSCRIBER_ID环境变量:

# 实例1
PORT=3002 node Subscriber.js
# 实例2
PORT=3003 node Subscriber.js
# 实例3
PORT=3004 node Subscriber.js

2. 工作者结果反馈机制

添加回调队列实现结果回传,修改核心代码如下:

修改后的Publisher.js(核心片段)

const { v4: uuidv4 } = require('uuid'); // 生成唯一任务ID

// 初始化时同时监听回调队列
async function connectQueue() {
  try {
    connection = await amqp.connect(process.env.AMQP_URL || "amqp://localhost:5672");
    channel = await connection.createChannel();

    await channel.assertQueue("math-tasks"); // 统一任务队列
    const callbackQueue = await channel.assertQueue('', { exclusive: true }); // 独占临时回调队列

    // 监听回调队列接收结果
    channel.consume(callbackQueue.queue, async (msg) => {
      const resultData = JSON.parse(msg.content.toString());
      console.log(`收到任务${resultData.taskId}结果:${resultData.result}`);
      // 触发任务记录更新(对应需求3)
      await updateTaskRecord(resultData.taskId, {
        status: 'completed',
        result: resultData.result
      });
      channel.ack(msg);
    });

    global.callbackQueueName = callbackQueue.queue;
  } catch (error) {
    console.log(error);
  }
}

// 修改为发送带回调信息的任务对象
async function sendTask(taskType, a, b) {
  const taskId = uuidv4();
  const task = {
    taskId,
    type: taskType,
    a,
    b,
    callbackQueue: global.callbackQueueName
  };

  await channel.sendToQueue("math-tasks", Buffer.from(JSON.stringify(task)));
  // 初始化任务记录(对应需求3)
  await createTaskRecord(taskId, {
    type: taskType,
    a,
    b,
    status: 'pending',
    publisherId: process.env.PUBLISHER_ID || 'default'
  });
  return taskId;
}

// 更新API接口
app.get("/math-task/sum", async (req, res) => {
  let a = parseInt(req.body.a);
  let b = parseInt(req.body.b);
  const taskId = await sendTask('sum', a, b);
  res.send(`加法任务已提交,任务ID:${taskId}`);
});

app.get("/math-task/mul", async (req, res) => {
  let a = parseInt(req.body.a);
  let b = parseInt(req.body.b);
  const taskId = await sendTask('mul', a, b);
  res.send(`乘法任务已提交,任务ID:${taskId}`);
});

修改后的Subscriber.js(核心片段)

async function connectQueue() {
  try {
    connection = await amqp.connect(process.env.AMQP_URL || "amqp://localhost:5672");
    channel = await connection.createChannel();
    
    await channel.assertQueue("math-tasks");
    
    channel.consume("math-tasks", async (msg) => {
      const task = JSON.parse(msg.content.toString());
      let result;
      // 执行任务
      if (task.type === 'sum') {
        result = task.a + task.b;
      } else if (task.type === 'mul') {
        result = task.a * task.b;
      }
      console.log(`完成任务${task.taskId}:${result}`);
      // 发送结果到回调队列
      await channel.sendToQueue(task.callbackQueue, Buffer.from(JSON.stringify({
        taskId: task.taskId,
        result: result
      })));
      channel.ack(msg);
    });
  } catch (error) {
    console.log(error);
  }
}

3. 发布者任务历史记录

使用Node.jsfs模块操作临时文件(如tasks.json),实现任务的创建与更新:

添加到Publisher.js的代码

const fs = require('fs').promises;
const TASK_FILE = './tasks.json';

// 初始化任务文件
async function initTaskFile() {
  try {
    await fs.access(TASK_FILE);
  } catch {
    await fs.writeFile(TASK_FILE, JSON.stringify([]));
  }
}

// 创建任务记录
async function createTaskRecord(taskId, taskData) {
  await initTaskFile();
  const tasks = JSON.parse(await fs.readFile(TASK_FILE, 'utf8'));
  tasks.push({
    taskId,
    ...taskData,
    createdAt: new Date().toISOString()
  });
  await fs.writeFile(TASK_FILE, JSON.stringify(tasks, null, 2));
}

// 更新任务记录
async function updateTaskRecord(taskId, updateData) {
  await initTaskFile();
  const tasks = JSON.parse(await fs.readFile(TASK_FILE, 'utf8'));
  const updatedTasks = tasks.map(task => 
    task.taskId === taskId ? { ...task, ...updateData, completedAt: new Date().toISOString() } : task
  );
  await fs.writeFile(TASK_FILE, JSON.stringify(updatedTasks, null, 2));
}

// 启动时初始化任务文件
initTaskFile();

4. Docker多实例部署

步骤1:编写Dockerfile

创建Dockerfile(Publisher和Subscriber共用):

FROM node:18-alpine
WORKDIR /app
COPY package*.json ./
RUN npm install --production
COPY . .
CMD ["node", "Publisher.js"] # 默认启动Publisher,可通过command覆盖

步骤2:编写docker-compose.yml

创建docker-compose.yml,定义所有服务:

version: '3.8'

services:
  rabbitmq:
    image: rabbitmq:3-management-alpine
    ports:
      - "5672:5672" # AMQP端口
      - "15672:15672" # 管理界面端口
    volumes:
      - rabbitmq_data:/var/lib/rabbitmq

  publisher1:
    build: .
    environment:
      - PORT=3000
      - PUBLISHER_ID=pub1
      - AMQP_URL=amqp://rabbitmq:5672
    ports:
      - "3000:3000"
    volumes:
      - ./tasks_pub1:/app # 挂载目录保存任务记录
    depends_on:
      - rabbitmq

  publisher2:
    build: .
    environment:
      - PORT=3001
      - PUBLISHER_ID=pub2
      - AMQP_URL=amqp://rabbitmq:5672
    ports:
      - "3001:3001"
    volumes:
      - ./tasks_pub2:/app
    depends_on:
      - rabbitmq

  subscriber1:
    build: .
    environment:
      - PORT=3002
      - AMQP_URL=amqp://rabbitmq:5672
    command: ["node", "Subscriber.js"]
    depends_on:
      - rabbitmq

  subscriber2:
    build: .
    environment:
      - PORT=3003
      - AMQP_URL=amqp://rabbitmq:5672
    command: ["node", "Subscriber.js"]
    depends_on:
      - rabbitmq

  subscriber3:
    build: .
    environment:
      - PORT=3004
      - AMQP_URL=amqp://rabbitmq:5672
    command: ["node", "Subscriber.js"]
    depends_on:
      - rabbitmq

volumes:
  rabbitmq_data:

步骤3:启动服务

执行命令启动所有实例:

docker-compose up -d --build

可通过http://localhost:3000/3001访问两个Publisher的API,RabbitMQ管理界面在http://localhost:15672(默认账号guest/guest)。


内容的提问来源于stack exchange,提问作者Ronak07

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 08:15:32