rabbitmq

rabbitmq

rabbitmq客户端

快速开始

import rabbitmq from "rabbitmq"
//创建连接
let connection = rabbitmq.open("amqp://admin:123456@localhost:5672/")
//  创建通道
let ch = connection.channel()
//声明队列
let queue = ch.queueDeclare({name: "queue1"})
//发布消息
ch.publish({key: "queue1", msg: {contextType: "text/plain", body: "hello world"}})
//订阅消息
ch.consume({queue: "queue1"}, (msg) => {
    //收到并处理消息
    console.log(msg)
})
//线程暂停
sleep()

Rabbitmq

名称类型参数返回值说明
open方法url:stringConnection创建一个连接对象

open

通过 rabbitmq.open(url) 方法创建一个连接。

import rabbitmq from "rabbitmq"
//创建连接
let connection = rabbitmq.open("amqp://admin:123456@localhost:5672/")

Connection

Connection对象表示一个rabbitmq的连接。

名称类型参数返回值说明
channel方法Channel打开一个频道
close方法关闭连接
isClosed方法bool判断当前连接是否关闭

channel

connection.channel()方法打开一个频道

示例

import rabbitmq from "rabbitmq"
//创建连接
let connection = rabbitmq.open("amqp://admin:123456@localhost:5672/")
//  打开一个频道
let ch = connection.channel()

close

connection.close()方法关闭连接

示例

import rabbitmq from "rabbitmq"
//创建连接
let connection = rabbitmq.open("amqp://admin:123456@localhost:5672/")
//  关闭连接
connection.close()

isClosed

connection.isClosed()方法关闭连接

示例

import rabbitmq from "rabbitmq"
//创建连接
let connection = rabbitmq.open("amqp://admin:123456@localhost:5672/")
// 判断连接是否关闭,下面输出: false
console.log(connection.isClosed())
//  关闭连接
connection.close()
// 下面输出: true
console.log(connection.isClosed())

Channel

Channel对象表示一个rabbitmq的频道。

名称类型参数返回值说明
ack方法tag:uint64, multiple:boolerror确认一条或多条消息(ACK)
cancel方法consumer:string, noWait:boolerror取消指定的消费者(consumer)
close方法error关闭通道
confirm方法noWait:boolerror启用发布确认模式(publisher confirms)
consume方法queue:string, consumer:string, autoAck:bool, exclusive:bool, noLocal:bool, noWait:bool, args:Table<-chan Delivery, error从队列消费消息,返回消息通道(异步消费)
exchangeBind方法destination:string, key:string, source:string, noWait:bool, args:Tableerror将交换机绑定到另一个交换机(exchange bind)
exchangeDeclare方法name:string, kind:string, durable:bool, autoDelete:bool, internal:bool, noWait:bool, args:Tableerror声明(创建)交换机
exchangeDeclarePassive方法name:string, kind:string, durable:bool, autoDelete:bool, internal:bool, noWait:bool, args:Tableerror被动声明交换机(仅检查是否存在)
exchangeDelete方法name:string, ifUnused:bool, noWait:boolerror删除交换机
exchangeUnbind方法destination:string, key:string, source:string, noWait:bool, args:Tableerror解除交换机绑定
flow方法active:boolerror暂停或恢复消息流(流量控制)
get方法queue:string, autoAck:boolmsg:Delivery, ok:bool, err:error从队列同步拉取单条消息(poll)
getNextPublishSeqNo方法uint64获取下一个发布序列号(用于确认追踪)
isClosed方法bool检查通道是否已关闭
nack方法tag:uint64, multiple:bool, requeue:boolerror否定确认(NACK),可选择重入队列或丢弃
notifyCancel方法c:chan stringchan string注册消费者取消通知通道
notifyClose方法c:chan *Errorchan *Error注册通道关闭通知通道
notifyConfirm方法ack:chan uint64, nack:chan uint64(chan uint64, chan uint64)注册发布确认(ack/nack)通知通道
notifyFlow方法c:chan boolchan bool注册流量控制(flow)通知通道
notifyPublish方法confirm:chan Confirmationchan Confirmation注册发布确认(Confirmation)通知通道
notifyReturn方法c:chan Returnchan Return注册返回(unroutable message)通知通道
publish方法exchange:string, key:string, mandatory:bool, immediate:bool, msg:Publishingerror发布消息(同步)
publishWithDeferredConfirm方法exchange:string, key:string, mandatory:bool, immediate:bool, msg:Publishing(*DeferredConfirmation, error)发布并返回 DeferredConfirmation,用于后续等待确认
qos方法prefetchCount:int, prefetchSize:int, global:boolerror设置预取(QOS)参数以控制消费速率
queueBind方法name:string, key:string, exchange:string, noWait:bool, args:Tableerror将队列绑定到交换机(queue bind)
queueDeclare方法name:string, durable:bool, autoDelete:bool, exclusive:bool, noWait:bool, args:Table(Queue, error)声明(创建)队列并返回队列信息
queueDeclarePassive方法name:string, durable:bool, autoDelete:bool, exclusive:bool, noWait:bool, args:Table(Queue, error)被动声明队列(仅检查是否存在)
queueDelete方法name:string, ifUnused:bool, ifEmpty:bool, noWait:bool(int, error)删除队列,返回被删除消息数
queueInspect方法name:string(Queue, error)检查并返回队列状态(消息数、消费者数等)
queuePurge方法name:string, noWait:bool(int, error)清空队列,返回被移除的消息数
queueUnbind方法name:string, key:string, exchange:string, args:Tableerror解除队列与交换机的绑定
recover方法requeue:boolerror重新排队未确认的消息(requeue=true 则重入队列)
reject方法tag:uint64, requeue:boolerror拒绝消息(等同于 nack/reject)
tx方法error启动事务模式
txCommit方法error提交事务
txRollback方法error回滚事务

ack

ch.ack(tag, multiple)方法确认一条或多条消息。

参数说明

参数类型必需说明
taguint64消息的投递标签(delivery tag)
multiplebool如果为true,则确认所有比指定标签更早的未确认消息;如果为false,只确认指定的消息

返回值

  • 成功时返回true
  • 如果参数错误或执行失败,会抛出异常

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 消费消息
let consumer = ch.consume({queue: "queue1"}, (msg) => {
    // 处理消息...

    // 确认单条消息
    ch.ack(msg.deliveryTag, false)

    // 或者确认此消息及之前的所有消息
    // ch.ack(msg.deliveryTag, true)
})

cancel

ch.cancel(consumer, noWait)方法取消指定的消费者。

参数说明

参数类型必需说明
consumerstring消费者标识符(consumer tag)
noWaitbool如果为true,不等待服务器响应;如果为false,会等待服务器取消操作的响应结果

返回值

  • 成功时返回true
  • 如果参数错误或执行失败,会抛出异常

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 开始消费
let consumer = ch.consume({queue: "queue1"}, (msg) => {
    console.log("收到消息:", msg)
})

// 一段时间后取消消费者
setTimeout(() => {
    // 不等待服务器响应
    ch.cancel(consumer, true)

    // 或者等待服务器响应
    // ch.cancel(consumer, false)

    console.log("消费者已取消")
}, 10000)

close

ch.close()方法关闭通道。

参数说明

  • 无参数

返回值

  • 成功时返回void
  • 如果执行失败,会抛出异常

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 使用通道进行各种操作...

// 关闭通道
ch.close()

// 此时如果再使用该通道会抛出异常

confirm

ch.confirm(noWait)方法启用发布确认模式。

参数说明

参数类型必需默认值说明
noWaitboolfalse如果为true,不等待服务器响应;如果为false,会等待服务器确认模式设置的响应结果

返回值

  • 成功时返回true
  • 如果执行失败,会抛出异常

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 启用发布确认模式(等待服务器响应)
ch.confirm(false)

// 启用发布确认模式(不等待服务器响应)
// ch.confirm(true)

// 现在可以使用发布确认相关功能
ch.publish({
    exchange: "",
    key: "queue1",
    msg: {body: "Hello World"}
})

consume

ch.consume(options, callback)方法从指定队列消费消息,启动一个异步消费者。

参数说明

第一个参数是一个包含消费选项的对象:

参数类型必需默认值说明
queuestring-要消费的队列名称
consumerstring自动生成消费者标识符,用于后续取消消费
autoAckboolfalse是否自动确认消息。如果为true,消息会在发送给消费者后立即被确认;如果为false,需要手动调用msg.ack()等方法确认消息
exclusiveboolfalse是否排他消费。如果为true,只有当前消费者可以访问队列
noLocalboolfalse如果为true,消费者不会接收自己发布的消息
noWaitboolfalse是否不等待服务器响应
argsTable空对象额外的参数表

第二个参数是一个回调函数:

参数类型必需说明
callbackFunction收到消息时的回调函数,接收一个消息对象作为参数

返回值

  • 成功时返回true
  • 如果参数错误或执行失败,会抛出异常

消息对象说明

回调函数接收的消息对象包含以下属性和方法:

属性/方法类型说明
bodystring消息内容(字符串形式)
deliveryTagint64消息投递标签,用于确认消息
redeliveredbool是否重新投递的消息
exchangestring消息来源的交换机名称
routingKeystring消息的路由键
ack(multiple)Function确认消息。参数multiple(可选,默认false)表示是否确认所有比此标签更早的未确认消息
nack(multiple, requeue)Function否定确认消息。参数multiple(可选,默认false)表示是否拒绝所有比此标签更早的未确认消息;requeue(可选,默认true)表示是否将消息重新入队
reject(requeue)Function拒绝消息。参数requeue(可选,默认true)表示是否将消息重新入队

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 消费消息,不自动确认
ch.consume({
    queue: "queue1",
    consumer: "myConsumer",  // 可选
    autoAck: false,          // 手动确认
    exclusive: false,
    noLocal: false,
    noWait: false,
    args: {}                 // 可选参数
}, (msg) => {
    console.log("收到消息:", msg.body)
    console.log("来自交换机:", msg.exchange)
    console.log("路由键:", msg.routingKey)

    // 处理消息...

    // 手动确认消息
    msg.ack(false)  // 只确认此条消息

    // 或者拒绝消息并重新入队
    // msg.reject(true)

    // 或者否定确认并重新入队
    // msg.nack(false, true)
})

exchangeBind

ch.exchangeBind(options)方法将一个交换机绑定到另一个交换机。

参数说明

参数是一个对象,包含以下属性:

参数类型必需说明
destinationstring目标交换机名称(被绑定的交换机)
sourcestring源交换机名称(绑定的交换机)
keystring绑定键(routing key)
noWaitbool是否不等待服务器响应,默认false
argsTable额外的参数表

返回值

  • 成功时返回true
  • 如果参数错误或执行失败,会抛出异常

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 声明两个交换机
ch.exchangeDeclare({
    name: "logs",
    kind: "fanout"
})

ch.exchangeDeclare({
    name: "backup",
    kind: "fanout"
})

// 将 backup 交换机绑定到 logs 交换机
ch.exchangeBind({
    destination: "backup",
    source: "logs",
    key: "",          // fanout 交换机忽略 routing key
    noWait: false,
    args: {}          // 可选参数
})

exchangeDeclare

ch.exchangeDeclare(options)方法声明(创建)一个交换机。

参数说明

参数是一个对象,包含以下属性:

参数类型必需默认值说明
namestring-交换机名称
kindstring-交换机类型:"direct""fanout""topic""headers"
durableboolfalse是否持久化。如果为true,交换机会在服务器重启后依然存在
autoDeleteboolfalse是否自动删除。如果为true,当所有绑定队列都解绑后,交换机自动删除
internalboolfalse是否内部交换机。如果为true,此交换机不能直接接收客户端发布的消息
noWaitboolfalse是否不等待服务器响应
argsTable空对象额外的参数表

返回值

  • 成功时返回true
  • 如果参数错误或执行失败,会抛出异常

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 声明一个持久化的 direct 交换机
ch.exchangeDeclare({
    name: "direct_logs",
    kind: "direct",
    durable: true,
    autoDelete: false,
    internal: false,
    noWait: false,
    args: {
        "alternate-exchange": "ae"  // 设置备用交换机
    }
})

// 声明一个自动删除的 fanout 交换机
ch.exchangeDeclare({
    name: "fanout_logs",
    kind: "fanout",
    durable: false,
    autoDelete: true
})

exchangeDeclarePassive

ch.exchangeDeclarePassive(options)方法被动声明一个交换机,仅检查交换机是否存在。

参数说明

参数是一个对象,包含以下属性:

参数类型必需默认值说明
namestring-交换机名称
kindstring-交换机类型:"direct""fanout""topic""headers"
durableboolfalse是否持久化
autoDeleteboolfalse是否自动删除
internalboolfalse是否内部交换机
noWaitboolfalse是否不等待服务器响应
argsTable空对象额外的参数表

返回值

  • 成功时返回true,表示交换机存在
  • 如果交换机不存在或参数错误,会抛出异常

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

try {
    // 检查交换机是否存在
    ch.exchangeDeclarePassive({
        name: "existing_exchange",
        kind: "direct"
    })
    console.log("交换机存在")
} catch (err) {
    console.log("交换机不存在或发生错误:", err.message)
    // 如果不存在,可以创建它
    ch.exchangeDeclare({
        name: "existing_exchange",
        kind: "direct"
    })
}

exchangeDelete

ch.exchangeDelete(options)方法删除一个交换机。

参数说明

参数是一个对象,包含以下属性:

参数类型必需默认值说明
namestring-要删除的交换机名称
ifUnusedboolfalse是否只在没有队列绑定时才删除。如果为true且仍有队列绑定到此交换机,删除会失败
noWaitboolfalse是否不等待服务器响应

返回值

  • 成功时返回true
  • 如果参数错误或执行失败,会抛出异常

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 删除交换机(无条件)
ch.exchangeDelete({
    name: "unused_exchange"
})

// 只在没有绑定时删除
ch.exchangeDelete({
    name: "maybe_used_exchange",
    ifUnused: true
})

exchangeUnbind

ch.exchangeUnbind(options)方法解除两个交换机之间的绑定关系。

参数说明

参数是一个对象,包含以下属性:

参数类型必需说明
destinationstring目标交换机名称(被绑定的交换机)
sourcestring源交换机名称(绑定的交换机)
keystring绑定键(routing key)
noWaitbool是否不等待服务器响应,默认false
argsTable额外的参数表(必须与绑定时使用的参数一致)

返回值

  • 成功时返回true
  • 如果参数错误或执行失败,会抛出异常

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 解除交换机绑定
ch.exchangeUnbind({
    destination: "backup",
    source: "logs",
    key: "",
    noWait: false,
    args: {}  // 必须与绑定时使用的参数一致
})

flow

ch.flow(active)方法暂停或恢复消息流(流量控制)。

参数说明

参数类型必需说明
activebool如果为false,暂停消息流,消费者将不再接收新消息;如果为true,恢复消息流

返回值

  • 成功时返回true
  • 如果参数错误或执行失败,会抛出异常

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 暂停消息流
ch.flow(false)
console.log("消息流已暂停")

// 一段时间后恢复消息流
setTimeout(() => {
    ch.flow(true)
    console.log("消息流已恢复")
}, 10000)

get

ch.get(queue, autoAck)方法从队列中同步获取(拉取)单条消息。

参数说明

参数类型必需默认值说明
queuestring-队列名称
autoAckboolfalse是否自动确认消息。如果为true,消息会在获取后立即被确认

返回值

  • 如果队列中有消息,返回一个消息对象
  • 如果队列中没有消息,返回null
  • 如果参数错误或执行失败,会抛出异常

消息对象说明

返回的消息对象与consume方法回调中的消息对象类似,但额外包含一个messageCount属性:

属性/方法类型说明
bodystring消息内容(字符串形式)
deliveryTagint64消息投递标签,用于确认消息
redeliveredbool是否重新投递的消息
exchangestring消息来源的交换机名称
routingKeystring消息的路由键
messageCountint64队列中剩余的消息数量
ack(multiple)Function确认消息。参数multiple表示是否确认所有比此标签更早的未确认消息
nack(multiple, requeue)Function否定确认消息。参数multiple表示是否拒绝所有比此标签更早的未确认消息;requeue表示是否将消息重新入队
reject(requeue)Function拒绝消息。参数requeue表示是否将消息重新入队

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 声明队列
ch.queueDeclare({
    name: "queue1",
    durable: false
})

// 发布一条消息
ch.publish({
    exchange: "",
    key: "queue1",
    msg: {body: "Hello World"}
})

// 获取消息(不自动确认)
let msg = ch.get("queue1", false)
if (msg) {
    console.log("获取到消息:", msg.body)
    console.log("队列中剩余消息数:", msg.messageCount)

    // 确认消息
    msg.ack(false)

    // 或者拒绝消息并重新入队
    // msg.reject(true)
} else {
    console.log("队列中没有消息")
}

// 再次获取(应该为空)
let msg2 = ch.get("queue1", false)
if (!msg2) {
    console.log("队列已空")
}

getNextPublishSeqNo

ch.getNextPublishSeqNo()方法获取下一个发布序列号,用于在发布确认模式中追踪消息。

参数说明

  • 无参数

返回值

  • 返回一个int64类型的发布序列号
  • 如果没有启用发布确认模式,返回值可能不可预测

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 启用发布确认模式
ch.confirm(false)

// 获取下一个发布序列号
let seqNo = ch.getNextPublishSeqNo()
console.log("下一个发布序列号:", seqNo)

// 发布消息
ch.publish({
    exchange: "",
    key: "queue1",
    msg: {body: "Hello World"}
})

// 序列号会增加
let nextSeqNo = ch.getNextPublishSeqNo()
console.log("下一个发布序列号:", nextSeqNo)

isClosed

ch.isClosed()方法检查通道是否已关闭。

参数说明

  • 无参数

返回值

  • 返回一个bool值:
    • true:通道已关闭
    • false:通道仍然打开

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 检查通道状态
console.log("通道是否关闭:", ch.isClosed()) // false

// 使用通道...

// 关闭通道
ch.close()

// 再次检查
console.log("通道是否关闭:", ch.isClosed()) // true

nack

ch.nack(tag, multiple, requeue)方法否定确认一条或多条消息。

参数说明

参数类型必需说明
taguint64消息的投递标签(delivery tag)
multiplebool如果为true,则拒绝所有比指定标签更早的未确认消息;如果为false,只拒绝指定的消息
requeuebool如果为true,消息会被重新放入队列(可能会被其他消费者消费);如果为false,消息会被丢弃

返回值

  • 成功时返回true
  • 如果参数错误或执行失败,会抛出异常

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 消费消息(手动确认)
ch.consume({queue: "queue1", autoAck: false}, (msg) => {
    console.log("收到消息:", msg.body)

    try {
        // 处理消息...
        // 假设处理失败,需要拒绝消息

        // 拒绝单条消息并重新入队
        ch.nack(msg.deliveryTag, false, true)

        // 或者拒绝多条消息并丢弃
        // ch.nack(msg.deliveryTag, true, false)

    } catch (err) {
        console.log("处理消息失败,消息已重新入队")
    }
})

notifyCancel

ch.notifyCancel(callback)方法注册消费者取消通知,当消费者被取消时会调用回调函数。

参数说明

参数类型必需说明
callbackFunction消费者取消时的回调函数,接收消费者标识符(consumer tag)作为参数

返回值

  • 成功时返回true
  • 如果参数错误或执行失败,会抛出异常

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 注册消费者取消通知
ch.notifyCancel((consumerTag) => {
    console.log("消费者被取消:", consumerTag)
    // 可以进行清理操作...
})

// 开始消费
let consumerTag = ch.consume({queue: "queue1"}, (msg) => {
    console.log("收到消息:", msg.body)
    msg.ack(false)
})

// 一段时间后取消消费者(会触发通知)
setTimeout(() => {
    ch.cancel(consumerTag, false)
}, 10000)

notifyClose

ch.notifyClose(callback)方法注册通道关闭通知,当通道关闭时会调用回调函数。

参数说明

参数类型必需说明
callbackFunction通道关闭时的回调函数,接收一个错误对象或null作为参数

返回值

  • 成功时返回true
  • 如果参数错误或执行失败,会抛出异常

错误对象说明

如果通道因错误而关闭,回调函数会收到一个包含以下属性的错误对象:

属性类型说明
codeint64错误代码
reasonstring错误原因描述
serverbool是否为服务器错误
recoverbool是否可恢复

如果通道正常关闭(无错误),回调函数会收到null

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 注册通道关闭通知
ch.notifyClose((err) => {
    if (err) {
        console.log("通道因错误关闭:")
        console.log("错误代码:", err.code)
        console.log("错误原因:", err.reason)
        console.log("是否服务器错误:", err.server)
        console.log("是否可恢复:", err.recover)
    } else {
        console.log("通道正常关闭")
    }
})

// 执行一些操作...

// 正常关闭通道(会触发通知,err为null)
ch.close()

// 或者通道可能因错误关闭,例如服务器断开连接

notifyConfirm

ch.notifyConfirm(ackCallback, nackCallback)方法注册发布确认通知,用于接收消息发布成功(ack)或失败(nack)的通知。

参数说明

参数类型必需说明
ackCallbackFunction消息发布成功时的回调函数,接收消息的投递标签(delivery tag)作为参数
nackCallbackFunction消息发布失败时的回调函数,接收消息的投递标签(delivery tag)作为参数

返回值

  • 成功时返回true
  • 如果参数错误或执行失败,会抛出异常

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 启用发布确认模式
ch.confirm(false)

// 注册确认通知
ch.notifyConfirm(
    (tag) => {
        console.log("消息发布成功,标签:", tag)
    },
    (tag) => {
        console.log("消息发布失败,标签:", tag)
    }
)

// 发布消息(启用确认模式后)
ch.publish({
    exchange: "",
    key: "queue1",
    msg: {body: "Hello World"}
})

ch.publish({
    exchange: "",
    key: "queue1",
    msg: {body: "Another Message"}
})

notifyFlow

ch.notifyFlow(callback)方法注册流量控制通知,当RabbitMQ暂停或恢复消息传递时会调用回调函数。

参数说明

参数类型必需说明
callbackFunction流量控制状态变化的回调函数,接收一个bool值作为参数:false表示消息流被暂停,true表示消息流被恢复

返回值

  • 成功时返回true
  • 如果参数错误或执行失败,会抛出异常

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 注册流量控制通知
ch.notifyFlow((active) => {
    if (active) {
        console.log("消息流已恢复,可以继续发送消息")
    } else {
        console.log("消息流已暂停,暂时不能发送消息")
    }
})

// 消费消息
ch.consume({queue: "queue1"}, (msg) => {
    console.log("收到消息:", msg.body)
    msg.ack(false)
})

// 注意:流量控制通常由RabbitMQ服务器触发,
// 当服务器负载过高时可能会暂停消息流

notifyPublish

ch.notifyPublish(callback)方法注册发布确认通知,返回一个包含投递标签和确认状态的对象。

参数说明

参数类型必需说明
callbackFunction发布确认的回调函数,接收一个确认对象作为参数

返回值

  • 成功时返回true
  • 如果参数错误或执行失败,会抛出异常

确认对象说明

回调函数接收的确认对象包含以下属性:

属性类型说明
deliveryTagint64消息的投递标签
ackbool确认状态:true表示消息已确认,false表示消息未被确认(nack)

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 启用发布确认模式
ch.confirm(false)

// 注册发布确认通知
ch.notifyPublish((confirmation) => {
    if (confirmation.ack) {
        console.log("消息已确认,标签:", confirmation.deliveryTag)
    } else {
        console.log("消息未确认(nack),标签:", confirmation.deliveryTag)
        // 可以进行重试等操作
    }
})

// 发布多条消息
for (let i = 0; i < 5; i++) {
    ch.publish({
        exchange: "",
        key: "queue1",
        msg: {body: `Message ${i}`}
    })
}

// 等待所有确认
setTimeout(() => {
    console.log("所有消息确认完成")
}, 1000)

notifyReturn

ch.notifyReturn(callback)方法注册消息返回通知,当消息无法路由到任何队列时(且设置了mandatoryimmediate标志)会调用回调函数。

参数说明

参数类型必需说明
callbackFunction消息返回时的回调函数,接收一个返回对象作为参数

返回值

  • 成功时返回true
  • 如果参数错误或执行失败,会抛出异常

返回对象说明

回调函数接收的返回对象包含以下属性:

属性类型说明
replyCodeint64返回代码(如404表示未找到)
replyTextstring返回文本说明
exchangestring消息发布的交换机名称
routingKeystring消息的路由键
bodystring消息内容(字符串形式)

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 注册返回通知
ch.notifyReturn((returned) => {
    console.log("消息无法路由,已返回:")
    console.log("返回代码:", returned.replyCode)
    console.log("返回原因:", returned.replyText)
    console.log("交换机:", returned.exchange)
    console.log("路由键:", returned.routingKey)
    console.log("消息内容:", returned.body)
})

// 发布消息,设置mandatory标志为true
// 如果消息无法路由到队列,会触发返回通知
ch.publish({
    exchange: "",
    key: "non.existent.queue",
    mandatory: true,  // 重要:必须设置为true才能触发返回
    immediate: false,
    msg: {
        contentType: "text/plain",
        body: "This message cannot be routed"
    }
})

publish

ch.publish(options)方法发布一条消息。

参数说明

参数是一个对象,包含以下属性:

参数类型必需默认值说明
exchangestring-交换机名称。空字符串表示默认交换机
keystring-路由键(routing key)
mandatoryboolfalse如果为true,当消息无法路由到任何队列时,消息会被返回给发布者(需要配合notifyReturn使用)
immediateboolfalse如果为true,当消息无法立即被消费者接收时,消息会被返回给发布者
msgObject-消息对象,包含消息内容和属性

消息对象说明

msg参数是一个对象,支持以下属性:

属性类型必需默认值说明
contentTypestring"text/plain"消息内容类型
bodystring-消息内容

返回值

  • 成功时返回true
  • 如果参数错误或执行失败,会抛出异常

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 声明队列
ch.queueDeclare({
    name: "queue1",
    durable: false
})

// 发布简单消息
ch.publish({
    exchange: "",
    key: "queue1",
    mandatory: false,
    immediate: false,
    msg: {
        contentType: "text/plain",
        body: "Hello World"
    }
})

// 发布到指定交换机
ch.exchangeDeclare({
    name: "logs",
    kind: "fanout"
})

ch.publish({
    exchange: "logs",
    key: "",  // fanout交换机忽略路由键
    msg: {
        contentType: "text/plain",
        body: "Log message"
    }
})

publishWithDeferredConfirm

ch.publishWithDeferredConfirm(options)方法发布一条消息,并返回一个延迟确认对象,用于后续等待确认。

参数说明

参数是一个对象,包含以下属性:

参数类型必需默认值说明
exchangestring-交换机名称。空字符串表示默认交换机
keystring-路由键(routing key)
mandatoryboolfalse如果为true,当消息无法路由到任何队列时,消息会被返回给发布者
immediateboolfalse如果为true,当消息无法立即被消费者接收时,消息会被返回给发布者
msgObject-消息对象,包含消息内容和属性

返回值

  • 成功时返回一个延迟确认对象(DeferredConfirmation)
  • 如果参数错误或执行失败,会抛出异常

延迟确认对象说明

返回的对象包含以下属性和方法:

属性/方法类型说明
deliveryTagint64消息的投递标签
confirmed()Function等待确认的方法,返回一个bool值:true表示消息已确认,false表示消息未被确认(nack)

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 启用发布确认模式
ch.confirm(false)

// 发布消息并获取延迟确认对象
let deferredConf = ch.publishWithDeferredConfirm({
    exchange: "",
    key: "queue1",
    mandatory: false,
    immediate: false,
    msg: {
        contentType: "text/plain",
        body: "Hello World"
    }
})

console.log("消息已发布,投递标签:", deferredConf.deliveryTag)

// 在需要时等待确认
setTimeout(() => {
    let isConfirmed = deferredConf.confirmed()
    if (isConfirmed) {
        console.log("消息已确认")
    } else {
        console.log("消息未被确认(nack)")
    }
}, 1000)

qos

ch.qos(prefetchCount, prefetchSize, global)方法设置服务质量(Quality of Service),控制消费者从队列中预取消息的数量。

参数说明

参数类型必需默认值说明
prefetchCountint-预取消息数量。0表示没有限制
prefetchSizeint-预取消息大小(字节)。0表示没有限制
globalbool-如果为true,设置应用于整个通道;如果为false,设置仅应用于每个消费者

返回值

  • 成功时返回true
  • 如果参数错误或执行失败,会抛出异常

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 设置预取:每个消费者最多同时处理1条消息
ch.qos(1, 0, false)

// 设置预取:整个通道最多同时处理10条消息
ch.qos(10, 0, true)

// 消费消息(将受QoS设置影响)
ch.consume({queue: "queue1"}, (msg) => {
    console.log("处理消息:", msg.body)

    // 模拟耗时处理
    setTimeout(() => {
        msg.ack(false)
        console.log("消息处理完成")
    }, 2000)
})

queueBind

ch.queueBind(options)方法将队列绑定到交换机。

参数说明

参数是一个对象,包含以下属性:

参数类型必需默认值说明
namestring-队列名称
keystring-绑定键(routing key)
exchangestring-交换机名称
noWaitboolfalse是否不等待服务器响应
argsTable空对象额外的参数表

返回值

  • 成功时返回true
  • 如果参数错误或执行失败,会抛出异常

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 声明队列和交换机
ch.queueDeclare({
    name: "error_logs",
    durable: false
})

ch.exchangeDeclare({
    name: "direct_logs",
    kind: "direct"
})

// 将队列绑定到交换机,只接收错误级别的日志
ch.queueBind({
    name: "error_logs",
    key: "error",      // 路由键为"error"
    exchange: "direct_logs",
    noWait: false,
    args: {}           // 可选参数
})

// 再绑定一个队列接收所有日志
ch.queueDeclare({
    name: "all_logs",
    durable: false
})

ch.queueBind({
    name: "all_logs",
    key: "#",          // "#"匹配所有路由键
    exchange: "direct_logs"
})

queueDeclare

ch.queueDeclare(options)方法声明(创建)一个队列。

参数说明

参数是一个对象,包含以下属性:

参数类型必需默认值说明
namestring自动生成队列名称。如果为空或未提供,服务器会生成一个唯一的队列名
durableboolfalse是否持久化。如果为true,队列会在服务器重启后依然存在(消息本身也需要持久化)
autoDeleteboolfalse是否自动删除。如果为true,当所有消费者断开连接后,队列会自动删除
exclusiveboolfalse是否排他队列。如果为true,队列仅对当前连接可见,连接关闭后队列自动删除
noWaitboolfalse是否不等待服务器响应
argsTable空对象额外的参数表

返回值

  • 成功时返回一个队列信息对象
  • 如果参数错误或执行失败,会抛出异常

队列信息对象说明

返回的对象包含以下属性:

属性类型说明
namestring队列名称(如果是自动生成的队列,这是生成的名称)
consumersint64当前消费者数量
messagesint64队列中的消息数量

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 声明一个持久化队列
let queueInfo = ch.queueDeclare({
    name: "durable_queue",
    durable: true,
    autoDelete: false,
    exclusive: false,
    args: {
        "x-max-length": 10000  // 限制队列最大长度
    }
})

console.log("队列声明成功:", queueInfo.name)
console.log("当前消费者数:", queueInfo.consumers)
console.log("当前消息数:", queueInfo.messages)

// 声明一个自动生成的临时队列
let tempQueue = ch.queueDeclare({
    name: "",  // 空字符串表示自动生成
    durable: false,
    autoDelete: true,
    exclusive: true
})

console.log("自动生成的队列名:", tempQueue.name)

queueDeclarePassive

ch.queueDeclarePassive(options)方法被动声明队列,仅检查队列是否存在而不创建新队列。

参数说明

参数是一个对象,包含以下属性:

参数类型必需默认值说明
namestring-队列名称
durableboolfalse是否持久化
autoDeleteboolfalse是否自动删除
exclusiveboolfalse是否排他队列
noWaitboolfalse是否不等待服务器响应
argsTable空对象额外的参数表

返回值

  • 成功时返回一个队列信息对象
  • 如果队列不存在或参数错误,会抛出异常

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

try {
    // 检查队列是否存在
    let queueInfo = ch.queueDeclarePassive({
        name: "existing_queue"
    })

    console.log("队列存在:")
    console.log("名称:", queueInfo.name)
    console.log("消费者数:", queueInfo.consumers)
    console.log("消息数:", queueInfo.messages)
} catch (err) {
    console.log("队列不存在或发生错误:", err.message)

    // 如果不存在,可以创建它
    ch.queueDeclare({
        name: "existing_queue",
        durable: false
    })
}

queueDelete

ch.queueDelete(options)方法删除一个队列。

参数说明

参数是一个对象,包含以下属性:

参数类型必需默认值说明
namestring-要删除的队列名称
ifUnusedboolfalse如果为true,只在没有消费者时删除队列
ifEmptyboolfalse如果为true,只在队列为空时删除队列
noWaitboolfalse是否不等待服务器响应

返回值

  • 成功时返回被删除队列中的消息数量(int64
  • 如果参数错误或执行失败,会抛出异常

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 创建测试队列并发布一些消息
let queueInfo = ch.queueDeclare({
    name: "test_queue",
    durable: false
})

// 发布5条消息
for (let i = 0; i < 5; i++) {
    ch.publish({
        exchange: "",
        key: "test_queue",
        msg: {body: `Message ${i}`}
    })
}

// 无条件删除队列(返回被删除的消息数)
let deletedCount = ch.queueDelete({
    name: "test_queue"
})
console.log(`队列已删除,删除了${deletedCount}条消息`)

// 创建新队列
ch.queueDeclare({
    name: "another_queue",
    durable: false
})

// 只在队列为空时删除
try {
    let count = ch.queueDelete({
        name: "another_queue",
        ifEmpty: true
    })
    console.log(`队列已删除,删除了${count}条消息`)
} catch (err) {
    console.log("删除失败(可能队列不为空):", err.message)
}

queueInspect

ch.queueInspect(name)方法检查队列状态并返回队列的当前信息。

参数说明

参数类型必需说明
namestring要检查的队列名称

返回值

  • 成功时返回一个队列信息对象
  • 如果队列不存在或参数错误,会抛出异常

队列信息对象说明

返回的对象包含以下属性:

属性类型说明
namestring队列名称
consumersint64当前消费者数量
messagesint64队列中的消息数量

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 检查队列状态
try {
    let queueInfo = ch.queueInspect("my_queue")

    console.log("队列状态:")
    console.log("名称:", queueInfo.name)
    console.log("消费者数:", queueInfo.consumers)
    console.log("消息数:", queueInfo.messages)

    if (queueInfo.messages > 1000) {
        console.log("警告: 队列积压了大量消息!")
    }
} catch (err) {
    console.log("队列不存在:", err.message)
}

queuePurge

ch.queuePurge(name, noWait)方法清空队列中的所有消息。

参数说明

参数类型必需默认值说明
namestring-要清空的队列名称
noWaitboolfalse是否不等待服务器响应

返回值

  • 成功时返回被清除的消息数量(int64
  • 如果参数错误或执行失败,会抛出异常

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 先向队列发布一些消息
for (let i = 0; i < 10; i++) {
    ch.publish({
        exchange: "",
        key: "queue_to_purge",
        msg: {body: `Message ${i}`}
    })
}

// 清空队列
let purgedCount = ch.queuePurge("queue_to_purge", false)
console.log(`清空了${purgedCount}条消息`)

// 验证队列已空
let queueInfo = ch.queueInspect("queue_to_purge")
console.log("队列中剩余消息:", queueInfo.messages) // 应为0

queueUnbind

ch.queueUnbind(options)方法解除队列与交换机的绑定关系。

参数说明

参数是一个对象,包含以下属性:

参数类型必需说明
namestring队列名称
keystring绑定键(routing key)
exchangestring交换机名称
argsTable额外的参数表(必须与绑定时使用的参数一致)

返回值

  • 成功时返回true
  • 如果参数错误或执行失败,会抛出异常

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 声明队列和交换机
ch.queueDeclare({
    name: "logs_queue",
    durable: false
})

ch.exchangeDeclare({
    name: "logs_exchange",
    kind: "fanout"
})

// 绑定队列到交换机
ch.queueBind({
    name: "logs_queue",
    key: "",  // fanout交换机忽略路由键
    exchange: "logs_exchange"
})

// ... 一段时间后,解除绑定
ch.queueUnbind({
    name: "logs_queue",
    key: "",
    exchange: "logs_exchange",
    args: {}  // 必须与绑定时使用的参数一致
})

console.log("队列已从交换机解绑")

recover

ch.recover(requeue)方法请求重新投递所有未确认的消息。

参数说明

参数类型必需默认值说明
requeuebooltrue如果为true,未确认的消息会被重新入队(可能会被其他消费者消费);如果为false,消息会重新投递给相同的消费者

返回值

  • 成功时返回true
  • 如果执行失败,会抛出异常

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 消费消息(不自动确认)
ch.consume({queue: "queue1", autoAck: false}, (msg) => {
    console.log("收到消息:", msg.body)

    // 不确认消息,模拟消费者崩溃
    // 注意:这里故意不调用 msg.ack() 或 msg.nack()
})

// 一段时间后,由于消费者崩溃,我们可以重新投递未确认的消息
setTimeout(() => {
    try {
        // 重新投递未确认的消息(重新入队)
        ch.recover(true)
        console.log("已重新投递未确认的消息")
    } catch (err) {
        console.log("重新投递失败:", err.message)
    }
}, 5000)

reject

ch.reject(tag, requeue)方法拒绝一条消息。

参数说明

参数类型必需说明
taguint64消息的投递标签(delivery tag)
requeuebool如果为true,消息会被重新放入队列(可能会被其他消费者消费);如果为false,消息会被丢弃

返回值

  • 成功时返回true
  • 如果参数错误或执行失败,会抛出异常

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 消费消息(手动确认)
ch.consume({queue: "queue1", autoAck: false}, (msg) => {
    console.log("收到消息:", msg.body)

    // 检查消息是否有效
    if (msg.body.includes("error")) {
        console.log("发现错误消息,拒绝并丢弃")

        // 拒绝消息并丢弃(不重新入队)
        ch.reject(msg.deliveryTag, false)
    } else if (msg.body.includes("retry")) {
        console.log("需要重试的消息,拒绝并重新入队")

        // 拒绝消息并重新入队
        ch.reject(msg.deliveryTag, true)
    } else {
        // 正常处理消息
        console.log("正常处理消息")
        msg.ack(false)
    }
})

tx

ch.tx()方法启动事务模式。

参数说明

  • 无参数

返回值

  • 成功时返回true
  • 如果执行失败,会抛出异常

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 启动事务模式
ch.tx()

console.log("事务模式已启动")

// 现在所有发布和确认操作都在事务中
// 可以执行多个操作...

// 提交或回滚事务

txCommit

ch.txCommit()方法提交当前事务。

参数说明

  • 无参数

返回值

  • 成功时返回true
  • 如果执行失败,会抛出异常

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 启动事务模式
ch.tx()

try {
    // 执行一系列操作
    ch.publish({
        exchange: "",
        key: "queue1",
        msg: {body: "Message 1"}
    })

    ch.publish({
        exchange: "",
        key: "queue1",
        msg: {body: "Message 2"}
    })

    // 确认一些消息(如果在同一个事务中)
    // ch.ack(someTag, false)

    // 提交事务(所有操作要么全部成功,要么全部失败)
    ch.txCommit()

    console.log("事务提交成功")
} catch (err) {
    console.log("操作失败,事务将回滚:", err.message)
    ch.txRollback()
}

txRollback

ch.txRollback()方法回滚当前事务。

参数说明

  • 无参数

返回值

  • 成功时返回true
  • 如果执行失败,会抛出异常

使用示例

import rabbitmq from "rabbitmq"

// 创建连接和通道
let conn = rabbitmq.open("amqp://admin:123456@localhost:5672/")
let ch = conn.channel()

// 启动事务模式
ch.tx()

try {
    // 执行第一个操作
    ch.publish({
        exchange: "",
        key: "queue1",
        msg: {body: "Message 1"}
    })

    // 执行第二个操作(可能失败)
    // 假设这个操作失败了
    throw new Error("模拟操作失败")

    // 如果到达这里,提交事务
    ch.txCommit()
} catch (err) {
    console.log("操作失败:", err.message)

    // 回滚事务(之前的所有操作都会被撤销)
    ch.txRollback()

    console.log("事务已回滚")
}

// 验证队列是否真的收到了消息
// 由于事务回滚,队列应该没有收到任何消息
let queueInfo = ch.queueInspect("queue1")
console.log("队列中的消息数:", queueInfo.messages) // 应为0
更新时间 3/18/2026, 7:21:33 PM