基于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
相关产品推荐
相关产品推荐

