Node.js MQ 进阶 — 工作模式
消息确认机制 (ack)
js
channel.consume(queue, (msg) => {
// 手动确认
channel.ack(msg) // 确认消费
// channel.nack(msg) // 否定,消息重新入队
// channel.reject(msg) // 拒绝,丢弃消息
}, { noAck: false }) // false = 手动确认模式四种工作模式
| 模式 | 特点 | 场景 |
|---|---|---|
| 简单模式 | 一生产者,一消费者 | 单任务 |
| Work Queues | 多消费者竞争消费同一队列 | 任务分发 |
| Pub/Sub | fanout 广播,每个消费者收到全部 | 通知推送 |
| Routing | direct 模式 + routingKey | 条件路由 |
公平分发 (prefetch)
js
channel.prefetch(1) // 每个消费者每次只取一条防止能力强的消费者闲着、能力弱的积压。
持久化
js
channel.assertQueue(queue, { durable: true }) // 队列持久化
channel.sendToQueue(queue, buf, { persistent: true }) // 消息持久化发布/订阅模式
js
// 生产者
const exchange = 'logs'
channel.assertExchange(exchange, 'fanout', { durable: false })
channel.publish(exchange, '', Buffer.from('广播消息'))
// 消费者(每个都收到)
channel.assertExchange(exchange, 'fanout')
const q = await channel.assertQueue('', { exclusive: true }) // 临时队列
channel.bindQueue(q.queue, exchange, '')
channel.consume(q.queue, (msg) => console.log(msg.content.toString()))RPC 模式
js
// 使用 correlationId + replyTo 实现请求-响应模式
channel.sendToQueue('rpc_queue', buf, {
correlationId,
replyTo: callbackQueue.queue
})