如何将队列消费的消息发送至Express应用?
Hey there! Let's walk through how to consume messages from your AMQP queue and hook that up to your Express app. You've already got the producer side working with your /hello endpoint, so let's build out the consumer piece step by step.
1. 核心思路
First, we need to set up an AMQP consumer that listens continuously to your target queue (producerQueue, the one your producer sends messages to). When a message arrives, we can pass it to your Express app's logic—whether that's storing it in memory, triggering a business process, or even pushing it to frontend clients (we'll cover basic use cases first).
2. 分步代码实现
First off, stop creating a new AMQP connection every time you send a message (your current sendMessageToQueue does this, which is inefficient). We'll refactor to reuse connections and channels, then add the consumer logic.
第一步:复用AMQP连接与通道
Let's extract the AMQP initialization into a reusable setup so we don't spin up new connections constantly:
const amqp = require('amqplib'); const express = require('express'); const app = express(); // 全局变量保存AMQP连接和通道 let amqpConn = null; let amqpChannel = null; // 初始化AMQP连接和通道 async function initAmqp() { try { amqpConn = await amqp.connect('amqp://localhost'); amqpChannel = await amqpConn.createChannel(); // 声明与生产者一致的队列(保持durable: false配置) await amqpChannel.assertQueue('producerQueue', { durable: false }); console.log('✅ AMQP连接成功,通道已就绪'); } catch (err) { console.error('❌ AMQP初始化失败:', err); // 添加重连逻辑,处理临时连接中断 setTimeout(initAmqp, 5000); } } // 启动Express前先初始化AMQP initAmqp();
第二步:搭建队列消费者
Now let's build the consumer that listens for incoming messages and processes them. For this example, we'll store messages in an array, but you can swap this out with your own business logic:
// 示例:存储收到的消息(可替换为你的业务逻辑) const receivedMessages = []; // 启动队列消费者 async function startConsumer() { // 如果通道还未就绪,等待1秒后重试 if (!amqpChannel) { await new Promise(resolve => setTimeout(resolve, 1000)); return startConsumer(); } // 监听队列并处理消息 amqpChannel.consume('producerQueue', (msg) => { if (msg) { const messageText = msg.content.toString(); console.log(`📥 收到消息: ${messageText}`); // 在这里与Express集成: // - 将消息存入数据库 // - 触发内部API逻辑 // - 通过WebSocket推送给前端 // 目前我们先简单存入内存 receivedMessages.push(messageText); // 手动确认消息处理完成,队列会移除该消息 // 务必在消息处理成功后再执行此操作! amqpChannel.ack(msg); } }, { noAck: false }); // noAck: false 表示需要手动确认消息 } // 启动消费者 startConsumer();
第三步:在Express中使用消费到的消息
Let's update your existing producer endpoint to use the reused channel, then add a new endpoint to retrieve the consumed messages:
// 重构sendMessageToQueue,使用复用的通道 async function sendMessageToQueue(msg, q) { if (!amqpChannel) { throw new Error('AMQP通道尚未就绪'); } await amqpChannel.assertQueue(q, { durable: false }); amqpChannel.sendToQueue(q, Buffer.from(msg)); } // 你已有的生产者接口(新增错误处理) app.get('/hello', async function(req, res) { try { await sendMessageToQueue('messageSent', 'producerQueue'); res.send('✅ 消息已成功发送至队列'); } catch (err) { res.status(500).send(`❌ 消息发送失败: ${err.message}`); } }); // 新增接口:获取所有已收到的消息 app.get('/received-messages', function(req, res) { res.json({ 消息总数: receivedMessages.length, 消息列表: receivedMessages }); }); // 启动Express服务 const PORT = 3000; app.listen(PORT, () => { console.log(`🚀 Express服务运行于 http://localhost:${PORT}`); });
3. 重要注意事项
- 连接复用:复用AMQP连接和通道对性能至关重要——每次发送消息都新建连接会极大浪费资源。
- 消息确认:使用
amqpChannel.ack(msg)确保只有在消息处理成功后才从队列移除,避免应用崩溃时丢失消息。 - 重连逻辑:AMQP连接可能意外断开,添加重连循环(比如
initAmqp中的setTimeout)能提升服务的韧性。 - 实时推送:如果需要向前端实时推送消息,可以搭配
ws等WebSocket库——当消费者收到消息时,广播给所有连接的客户端。
内容的提问来源于stack exchange,提问作者Tom Sawyer

