Skip to content

Node.js MQ 进阶 — 工作模式

消息确认机制 (ack)

js
channel.consume(queue, (msg) => {
  // 手动确认
  channel.ack(msg)       // 确认消费
  // channel.nack(msg)    // 否定,消息重新入队
  // channel.reject(msg)  // 拒绝,丢弃消息
}, { noAck: false })     // false = 手动确认模式

四种工作模式

模式特点场景
简单模式一生产者,一消费者单任务
Work Queues多消费者竞争消费同一队列任务分发
Pub/Subfanout 广播,每个消费者收到全部通知推送
Routingdirect 模式 + 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
})