本文目录导读:

- 方法一:Python (最推荐,代码简洁、功能完善)
- 方法二:Shell 脚本 (使用
amqp-consume或者rabbitmqadmin) - 方法三:Node.js (JavaScript)
- 方法四:Ruby
- 常见注意事项
- 推荐方案
要用脚本消费 RabbitMQ 消息,主要取决于你使用的编程语言,以下是几种常用语言的脚本示例,包括 Python (推荐)、Shell (配合工具) 和 Node.js。
核心思路
- 连接:连接到 RabbitMQ 服务器(IP、端口、vhost、凭据)。
- 声明/绑定:确保队列存在,并将队列绑定到交换机。
- 消费:监听消息并处理(打印、写入文件、调用API等)。
- 确认:处理成功后发送 ACK,避免消息丢失。
Python (最推荐,代码简洁、功能完善)
需要安装库:pip install pika
简单消费脚本 (自动确认)
import pika
import json
# 连接参数
credentials = pika.PlainCredentials('guest', 'guest')
parameters = pika.ConnectionParameters('localhost', 5672, '/', credentials)
# 建立连接
connection = pika.BlockingConnection(parameters)
channel = connection.channel()
# 声明队列 (如果不存在则创建)
queue_name = 'my_queue'
channel.queue_declare(queue=queue_name, durable=True) # durable 持久化
# 定义回调函数:处理收到的消息
def callback(ch, method, properties, body):
print(f" [x] Received: {body.decode()}")
# 如果你希望手动确认,在这里调用 ch.basic_ack(delivery_tag=method.delivery_tag)
# 这里自动确认(auto_ack=True),所以不需要手动 ACK
# 消费消息
# auto_ack=True 代表收到即确认(丢失风险低,但可能重复消费)
# auto_ack=False 则必须手动 ACK
channel.basic_consume(queue=queue_name,
on_message_callback=callback,
auto_ack=True)
print(' [*] Waiting for messages. To exit press CTRL+C')
channel.start_consuming()
带手动确认 & JSON解析
import pika
import json
def callback(ch, method, properties, body):
try:
data = json.loads(body.decode())
print(f"Received: {data}")
# 你的业务逻辑
# 处理成功后手动确认
ch.basic_ack(delivery_tag=method.delivery_tag)
except Exception as e:
print(f"处理失败: {e}")
# 可以选择拒绝并重新入队:ch.basic_nack(delivery_tag=method.delivery_tag, requeue=True)
# ...(同上连接部分)...
channel.basic_qos(prefetch_count=1) # 一次只处理一条,公平分发
channel.basic_consume(queue='my_queue', on_message_callback=callback, auto_ack=False)
channel.start_consuming()
Shell 脚本 (使用 amqp-consume 或者 rabbitmqadmin)
需要安装工具:rabbitmqadmin (来自 RabbitMQ 官方) 或 amqp-consume (来自 amqp-utils 包)。
使用 rabbitmqadmin(只支持获取单条或清空,不适合持续监听)
# 获取单条消息并删除 rabbitmqadmin get queue=my_queue --format=json # 或者一次性获取多条 rabbitmqadmin get queue=my_queue count=5 --format=raw
缺点:rabbitmqadmin get 是“拉模式”,一次一次轮询,不是真正的“推模式”阻塞监听,适合小规模检查,不适合持续消费。
使用 amqp-consume (更好)
# 安装 (macOS: brew install amqp-utils, Linux: pip install amqp-utils) # 持续消费并输出消息体 amqp-consume -u "amqp://guest:guest@localhost:5672/%2f" -q my_queue -c 10 ./handler.sh
这里 -c 10 表示最多处理10条退出,handler.sh 是处理脚本,如果没有 handler,可以:
amqp-consume -u "amqp://..." -q my_queue -c 1 cat # 只输出消息体到 stdout
Node.js (JavaScript)
需要安装 amqplib:npm install amqplib
const amqp = require('amqplib');
async function consume() {
const conn = await amqp.connect('amqp://guest:guest@localhost:5672');
const channel = await conn.createChannel();
const queue = 'my_queue';
await channel.assertQueue(queue, { durable: true });
console.log(' [*] Waiting for messages...');
channel.consume(queue, (msg) => {
if (msg !== null) {
const content = msg.content.toString();
console.log(" [x] Received:", content);
// 手动确认
channel.ack(msg);
}
}, { noAck: false });
}
consume().catch(console.error);
Ruby
require 'bunny'
connection = Bunny.new(host: 'localhost', user: 'guest', pass: 'guest')
connection.start
channel = connection.create_channel
queue = channel.queue('my_queue', durable: true)
puts ' [*] Waiting for messages...'
queue.subscribe(manual_ack: true, block: true) do |delivery_info, properties, body|
puts "Received: #{body}"
channel.ack(delivery_info.delivery_tag)
end
常见注意事项
| 问题 | 解决方案 |
|---|---|
| 连接错误 (Connection refused) | 确认 RabbitMQ 服务已启动:sudo systemctl start rabbitmq-server |
| 认证失败 (ACCESS_REFUSED) | 检查用户名/密码,默认是 guest/guest,但只允许 localhost |
| 队列不存在 | 使用 queue_declare 声明(幂等的,不存在才创建) |
| 消息不消费?卡住? | 检查交换机-路由键-队列绑定是否正确 |
| 消费速度慢? | 使用 channel.basic_qos(prefetch_count=1) 公平分发 |
| 脚本退出后消息丢失? | 使用 auto_ack=False + 手动 ack,并开启队列持久化 + 消息 delivery_mode=2 |
推荐方案
- 临时调试 / 一次运行:用 Python 或 Node.js 脚本最直接。
- 集成到自动化流水线(如CI/CD):Python 库
pika最稳定。 - 简单的监控/抓取:
rabbitmqadmin get或amqp-consume快速查看。
如果有特定的业务逻辑(比如把消息写入文件、转发到HTTP接口),在回调函数里加上相应代码即可。