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:string | Connection | 创建一个连接对象 |
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:bool | error | 确认一条或多条消息(ACK) |
| cancel | 方法 | consumer:string, noWait:bool | error | 取消指定的消费者(consumer) |
| close | 方法 | — | error | 关闭通道 |
| confirm | 方法 | noWait:bool | error | 启用发布确认模式(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:Table | error | 将交换机绑定到另一个交换机(exchange bind) |
| exchangeDeclare | 方法 | name:string, kind:string, durable:bool, autoDelete:bool, internal:bool, noWait:bool, args:Table | error | 声明(创建)交换机 |
| exchangeDeclarePassive | 方法 | name:string, kind:string, durable:bool, autoDelete:bool, internal:bool, noWait:bool, args:Table | error | 被动声明交换机(仅检查是否存在) |
| exchangeDelete | 方法 | name:string, ifUnused:bool, noWait:bool | error | 删除交换机 |
| exchangeUnbind | 方法 | destination:string, key:string, source:string, noWait:bool, args:Table | error | 解除交换机绑定 |
| flow | 方法 | active:bool | error | 暂停或恢复消息流(流量控制) |
| get | 方法 | queue:string, autoAck:bool | msg:Delivery, ok:bool, err:error | 从队列同步拉取单条消息(poll) |
| getNextPublishSeqNo | 方法 | — | uint64 | 获取下一个发布序列号(用于确认追踪) |
| isClosed | 方法 | — | bool | 检查通道是否已关闭 |
| nack | 方法 | tag:uint64, multiple:bool, requeue:bool | error | 否定确认(NACK),可选择重入队列或丢弃 |
| notifyCancel | 方法 | c:chan string | chan string | 注册消费者取消通知通道 |
| notifyClose | 方法 | c:chan *Error | chan *Error | 注册通道关闭通知通道 |
| notifyConfirm | 方法 | ack:chan uint64, nack:chan uint64 | (chan uint64, chan uint64) | 注册发布确认(ack/nack)通知通道 |
| notifyFlow | 方法 | c:chan bool | chan bool | 注册流量控制(flow)通知通道 |
| notifyPublish | 方法 | confirm:chan Confirmation | chan Confirmation | 注册发布确认(Confirmation)通知通道 |
| notifyReturn | 方法 | c:chan Return | chan Return | 注册返回(unroutable message)通知通道 |
| publish | 方法 | exchange:string, key:string, mandatory:bool, immediate:bool, msg:Publishing | error | 发布消息(同步) |
| publishWithDeferredConfirm | 方法 | exchange:string, key:string, mandatory:bool, immediate:bool, msg:Publishing | (*DeferredConfirmation, error) | 发布并返回 DeferredConfirmation,用于后续等待确认 |
| qos | 方法 | prefetchCount:int, prefetchSize:int, global:bool | error | 设置预取(QOS)参数以控制消费速率 |
| queueBind | 方法 | name:string, key:string, exchange:string, noWait:bool, args:Table | error | 将队列绑定到交换机(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:Table | error | 解除队列与交换机的绑定 |
| recover | 方法 | requeue:bool | error | 重新排队未确认的消息(requeue=true 则重入队列) |
| reject | 方法 | tag:uint64, requeue:bool | error | 拒绝消息(等同于 nack/reject) |
| tx | 方法 | — | error | 启动事务模式 |
| txCommit | 方法 | — | error | 提交事务 |
| txRollback | 方法 | — | error | 回滚事务 |
ack
ch.ack(tag, multiple)方法确认一条或多条消息。
参数说明
| 参数 | 类型 | 必需 | 说明 |
|---|---|---|---|
| tag | uint64 | 是 | 消息的投递标签(delivery tag) |
| multiple | bool | 是 | 如果为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)方法取消指定的消费者。
参数说明
| 参数 | 类型 | 必需 | 说明 |
|---|---|---|---|
| consumer | string | 是 | 消费者标识符(consumer tag) |
| noWait | bool | 是 | 如果为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)方法启用发布确认模式。
参数说明
| 参数 | 类型 | 必需 | 默认值 | 说明 |
|---|---|---|---|---|
| noWait | bool | 否 | false | 如果为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)方法从指定队列消费消息,启动一个异步消费者。
参数说明
第一个参数是一个包含消费选项的对象:
| 参数 | 类型 | 必需 | 默认值 | 说明 |
|---|---|---|---|---|
queue | string | 是 | - | 要消费的队列名称 |
consumer | string | 否 | 自动生成 | 消费者标识符,用于后续取消消费 |
autoAck | bool | 否 | false | 是否自动确认消息。如果为true,消息会在发送给消费者后立即被确认;如果为false,需要手动调用msg.ack()等方法确认消息 |
exclusive | bool | 否 | false | 是否排他消费。如果为true,只有当前消费者可以访问队列 |
noLocal | bool | 否 | false | 如果为true,消费者不会接收自己发布的消息 |
noWait | bool | 否 | false | 是否不等待服务器响应 |
args | Table | 否 | 空对象 | 额外的参数表 |
第二个参数是一个回调函数:
| 参数 | 类型 | 必需 | 说明 |
|---|---|---|---|
callback | Function | 是 | 收到消息时的回调函数,接收一个消息对象作为参数 |
返回值
- 成功时返回
true - 如果参数错误或执行失败,会抛出异常
消息对象说明
回调函数接收的消息对象包含以下属性和方法:
| 属性/方法 | 类型 | 说明 |
|---|---|---|
body | string | 消息内容(字符串形式) |
deliveryTag | int64 | 消息投递标签,用于确认消息 |
redelivered | bool | 是否重新投递的消息 |
exchange | string | 消息来源的交换机名称 |
routingKey | string | 消息的路由键 |
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)方法将一个交换机绑定到另一个交换机。
参数说明
参数是一个对象,包含以下属性:
| 参数 | 类型 | 必需 | 说明 |
|---|---|---|---|
destination | string | 是 | 目标交换机名称(被绑定的交换机) |
source | string | 是 | 源交换机名称(绑定的交换机) |
key | string | 是 | 绑定键(routing key) |
noWait | bool | 否 | 是否不等待服务器响应,默认false |
args | Table | 否 | 额外的参数表 |
返回值
- 成功时返回
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)方法声明(创建)一个交换机。
参数说明
参数是一个对象,包含以下属性:
| 参数 | 类型 | 必需 | 默认值 | 说明 |
|---|---|---|---|---|
name | string | 是 | - | 交换机名称 |
kind | string | 是 | - | 交换机类型:"direct"、"fanout"、"topic"、"headers" |
durable | bool | 否 | false | 是否持久化。如果为true,交换机会在服务器重启后依然存在 |
autoDelete | bool | 否 | false | 是否自动删除。如果为true,当所有绑定队列都解绑后,交换机自动删除 |
internal | bool | 否 | false | 是否内部交换机。如果为true,此交换机不能直接接收客户端发布的消息 |
noWait | bool | 否 | false | 是否不等待服务器响应 |
args | Table | 否 | 空对象 | 额外的参数表 |
返回值
- 成功时返回
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)方法被动声明一个交换机,仅检查交换机是否存在。
参数说明
参数是一个对象,包含以下属性:
| 参数 | 类型 | 必需 | 默认值 | 说明 |
|---|---|---|---|---|
name | string | 是 | - | 交换机名称 |
kind | string | 是 | - | 交换机类型:"direct"、"fanout"、"topic"、"headers" |
durable | bool | 否 | false | 是否持久化 |
autoDelete | bool | 否 | false | 是否自动删除 |
internal | bool | 否 | false | 是否内部交换机 |
noWait | bool | 否 | false | 是否不等待服务器响应 |
args | Table | 否 | 空对象 | 额外的参数表 |
返回值
- 成功时返回
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)方法删除一个交换机。
参数说明
参数是一个对象,包含以下属性:
| 参数 | 类型 | 必需 | 默认值 | 说明 |
|---|---|---|---|---|
name | string | 是 | - | 要删除的交换机名称 |
ifUnused | bool | 否 | false | 是否只在没有队列绑定时才删除。如果为true且仍有队列绑定到此交换机,删除会失败 |
noWait | bool | 否 | false | 是否不等待服务器响应 |
返回值
- 成功时返回
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)方法解除两个交换机之间的绑定关系。
参数说明
参数是一个对象,包含以下属性:
| 参数 | 类型 | 必需 | 说明 |
|---|---|---|---|
destination | string | 是 | 目标交换机名称(被绑定的交换机) |
source | string | 是 | 源交换机名称(绑定的交换机) |
key | string | 是 | 绑定键(routing key) |
noWait | bool | 否 | 是否不等待服务器响应,默认false |
args | Table | 否 | 额外的参数表(必须与绑定时使用的参数一致) |
返回值
- 成功时返回
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)方法暂停或恢复消息流(流量控制)。
参数说明
| 参数 | 类型 | 必需 | 说明 |
|---|---|---|---|
active | bool | 是 | 如果为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)方法从队列中同步获取(拉取)单条消息。
参数说明
| 参数 | 类型 | 必需 | 默认值 | 说明 |
|---|---|---|---|---|
queue | string | 是 | - | 队列名称 |
autoAck | bool | 否 | false | 是否自动确认消息。如果为true,消息会在获取后立即被确认 |
返回值
- 如果队列中有消息,返回一个消息对象
- 如果队列中没有消息,返回
null - 如果参数错误或执行失败,会抛出异常
消息对象说明
返回的消息对象与consume方法回调中的消息对象类似,但额外包含一个messageCount属性:
| 属性/方法 | 类型 | 说明 |
|---|---|---|
body | string | 消息内容(字符串形式) |
deliveryTag | int64 | 消息投递标签,用于确认消息 |
redelivered | bool | 是否重新投递的消息 |
exchange | string | 消息来源的交换机名称 |
routingKey | string | 消息的路由键 |
messageCount | int64 | 队列中剩余的消息数量 |
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)方法否定确认一条或多条消息。
参数说明
| 参数 | 类型 | 必需 | 说明 |
|---|---|---|---|
tag | uint64 | 是 | 消息的投递标签(delivery tag) |
multiple | bool | 是 | 如果为true,则拒绝所有比指定标签更早的未确认消息;如果为false,只拒绝指定的消息 |
requeue | bool | 是 | 如果为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)方法注册消费者取消通知,当消费者被取消时会调用回调函数。
参数说明
| 参数 | 类型 | 必需 | 说明 |
|---|---|---|---|
callback | Function | 是 | 消费者取消时的回调函数,接收消费者标识符(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)方法注册通道关闭通知,当通道关闭时会调用回调函数。
参数说明
| 参数 | 类型 | 必需 | 说明 |
|---|---|---|---|
callback | Function | 是 | 通道关闭时的回调函数,接收一个错误对象或null作为参数 |
返回值
- 成功时返回
true - 如果参数错误或执行失败,会抛出异常
错误对象说明
如果通道因错误而关闭,回调函数会收到一个包含以下属性的错误对象:
| 属性 | 类型 | 说明 |
|---|---|---|
code | int64 | 错误代码 |
reason | string | 错误原因描述 |
server | bool | 是否为服务器错误 |
recover | bool | 是否可恢复 |
如果通道正常关闭(无错误),回调函数会收到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)的通知。
参数说明
| 参数 | 类型 | 必需 | 说明 |
|---|---|---|---|
ackCallback | Function | 是 | 消息发布成功时的回调函数,接收消息的投递标签(delivery tag)作为参数 |
nackCallback | Function | 是 | 消息发布失败时的回调函数,接收消息的投递标签(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暂停或恢复消息传递时会调用回调函数。
参数说明
| 参数 | 类型 | 必需 | 说明 |
|---|---|---|---|
callback | Function | 是 | 流量控制状态变化的回调函数,接收一个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)方法注册发布确认通知,返回一个包含投递标签和确认状态的对象。
参数说明
| 参数 | 类型 | 必需 | 说明 |
|---|---|---|---|
callback | Function | 是 | 发布确认的回调函数,接收一个确认对象作为参数 |
返回值
- 成功时返回
true - 如果参数错误或执行失败,会抛出异常
确认对象说明
回调函数接收的确认对象包含以下属性:
| 属性 | 类型 | 说明 |
|---|---|---|
deliveryTag | int64 | 消息的投递标签 |
ack | bool | 确认状态: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)方法注册消息返回通知,当消息无法路由到任何队列时(且设置了mandatory或immediate标志)会调用回调函数。
参数说明
| 参数 | 类型 | 必需 | 说明 |
|---|---|---|---|
callback | Function | 是 | 消息返回时的回调函数,接收一个返回对象作为参数 |
返回值
- 成功时返回
true - 如果参数错误或执行失败,会抛出异常
返回对象说明
回调函数接收的返回对象包含以下属性:
| 属性 | 类型 | 说明 |
|---|---|---|
replyCode | int64 | 返回代码(如404表示未找到) |
replyText | string | 返回文本说明 |
exchange | string | 消息发布的交换机名称 |
routingKey | string | 消息的路由键 |
body | string | 消息内容(字符串形式) |
使用示例
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)方法发布一条消息。
参数说明
参数是一个对象,包含以下属性:
| 参数 | 类型 | 必需 | 默认值 | 说明 |
|---|---|---|---|---|
exchange | string | 是 | - | 交换机名称。空字符串表示默认交换机 |
key | string | 是 | - | 路由键(routing key) |
mandatory | bool | 否 | false | 如果为true,当消息无法路由到任何队列时,消息会被返回给发布者(需要配合notifyReturn使用) |
immediate | bool | 否 | false | 如果为true,当消息无法立即被消费者接收时,消息会被返回给发布者 |
msg | Object | 是 | - | 消息对象,包含消息内容和属性 |
消息对象说明
msg参数是一个对象,支持以下属性:
| 属性 | 类型 | 必需 | 默认值 | 说明 |
|---|---|---|---|---|
contentType | string | 否 | "text/plain" | 消息内容类型 |
body | string | 是 | - | 消息内容 |
返回值
- 成功时返回
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)方法发布一条消息,并返回一个延迟确认对象,用于后续等待确认。
参数说明
参数是一个对象,包含以下属性:
| 参数 | 类型 | 必需 | 默认值 | 说明 |
|---|---|---|---|---|
exchange | string | 是 | - | 交换机名称。空字符串表示默认交换机 |
key | string | 是 | - | 路由键(routing key) |
mandatory | bool | 否 | false | 如果为true,当消息无法路由到任何队列时,消息会被返回给发布者 |
immediate | bool | 否 | false | 如果为true,当消息无法立即被消费者接收时,消息会被返回给发布者 |
msg | Object | 是 | - | 消息对象,包含消息内容和属性 |
返回值
- 成功时返回一个延迟确认对象(DeferredConfirmation)
- 如果参数错误或执行失败,会抛出异常
延迟确认对象说明
返回的对象包含以下属性和方法:
| 属性/方法 | 类型 | 说明 |
|---|---|---|
deliveryTag | int64 | 消息的投递标签 |
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),控制消费者从队列中预取消息的数量。
参数说明
| 参数 | 类型 | 必需 | 默认值 | 说明 |
|---|---|---|---|---|
prefetchCount | int | 是 | - | 预取消息数量。0表示没有限制 |
prefetchSize | int | 是 | - | 预取消息大小(字节)。0表示没有限制 |
global | bool | 是 | - | 如果为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)方法将队列绑定到交换机。
参数说明
参数是一个对象,包含以下属性:
| 参数 | 类型 | 必需 | 默认值 | 说明 |
|---|---|---|---|---|
name | string | 是 | - | 队列名称 |
key | string | 是 | - | 绑定键(routing key) |
exchange | string | 是 | - | 交换机名称 |
noWait | bool | 否 | false | 是否不等待服务器响应 |
args | Table | 否 | 空对象 | 额外的参数表 |
返回值
- 成功时返回
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)方法声明(创建)一个队列。
参数说明
参数是一个对象,包含以下属性:
| 参数 | 类型 | 必需 | 默认值 | 说明 |
|---|---|---|---|---|
name | string | 否 | 自动生成 | 队列名称。如果为空或未提供,服务器会生成一个唯一的队列名 |
durable | bool | 否 | false | 是否持久化。如果为true,队列会在服务器重启后依然存在(消息本身也需要持久化) |
autoDelete | bool | 否 | false | 是否自动删除。如果为true,当所有消费者断开连接后,队列会自动删除 |
exclusive | bool | 否 | false | 是否排他队列。如果为true,队列仅对当前连接可见,连接关闭后队列自动删除 |
noWait | bool | 否 | false | 是否不等待服务器响应 |
args | Table | 否 | 空对象 | 额外的参数表 |
返回值
- 成功时返回一个队列信息对象
- 如果参数错误或执行失败,会抛出异常
队列信息对象说明
返回的对象包含以下属性:
| 属性 | 类型 | 说明 |
|---|---|---|
name | string | 队列名称(如果是自动生成的队列,这是生成的名称) |
consumers | int64 | 当前消费者数量 |
messages | int64 | 队列中的消息数量 |
使用示例
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)方法被动声明队列,仅检查队列是否存在而不创建新队列。
参数说明
参数是一个对象,包含以下属性:
| 参数 | 类型 | 必需 | 默认值 | 说明 |
|---|---|---|---|---|
name | string | 是 | - | 队列名称 |
durable | bool | 否 | false | 是否持久化 |
autoDelete | bool | 否 | false | 是否自动删除 |
exclusive | bool | 否 | false | 是否排他队列 |
noWait | bool | 否 | false | 是否不等待服务器响应 |
args | Table | 否 | 空对象 | 额外的参数表 |
返回值
- 成功时返回一个队列信息对象
- 如果队列不存在或参数错误,会抛出异常
使用示例
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)方法删除一个队列。
参数说明
参数是一个对象,包含以下属性:
| 参数 | 类型 | 必需 | 默认值 | 说明 |
|---|---|---|---|---|
name | string | 是 | - | 要删除的队列名称 |
ifUnused | bool | 否 | false | 如果为true,只在没有消费者时删除队列 |
ifEmpty | bool | 否 | false | 如果为true,只在队列为空时删除队列 |
noWait | bool | 否 | false | 是否不等待服务器响应 |
返回值
- 成功时返回被删除队列中的消息数量(
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)方法检查队列状态并返回队列的当前信息。
参数说明
| 参数 | 类型 | 必需 | 说明 |
|---|---|---|---|
name | string | 是 | 要检查的队列名称 |
返回值
- 成功时返回一个队列信息对象
- 如果队列不存在或参数错误,会抛出异常
队列信息对象说明
返回的对象包含以下属性:
| 属性 | 类型 | 说明 |
|---|---|---|
name | string | 队列名称 |
consumers | int64 | 当前消费者数量 |
messages | int64 | 队列中的消息数量 |
使用示例
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)方法清空队列中的所有消息。
参数说明
| 参数 | 类型 | 必需 | 默认值 | 说明 |
|---|---|---|---|---|
name | string | 是 | - | 要清空的队列名称 |
noWait | bool | 否 | false | 是否不等待服务器响应 |
返回值
- 成功时返回被清除的消息数量(
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)方法解除队列与交换机的绑定关系。
参数说明
参数是一个对象,包含以下属性:
| 参数 | 类型 | 必需 | 说明 |
|---|---|---|---|
name | string | 是 | 队列名称 |
key | string | 是 | 绑定键(routing key) |
exchange | string | 是 | 交换机名称 |
args | Table | 否 | 额外的参数表(必须与绑定时使用的参数一致) |
返回值
- 成功时返回
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)方法请求重新投递所有未确认的消息。
参数说明
| 参数 | 类型 | 必需 | 默认值 | 说明 |
|---|---|---|---|---|
requeue | bool | 否 | true | 如果为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)方法拒绝一条消息。
参数说明
| 参数 | 类型 | 必需 | 说明 |
|---|---|---|---|
tag | uint64 | 是 | 消息的投递标签(delivery tag) |
requeue | bool | 是 | 如果为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
