如何用脚本消费RabbitMQ消息?

wen 实用脚本 5

本文目录导读:

如何用脚本消费RabbitMQ消息?

  1. 方法一:Python (最推荐,代码简洁、功能完善)
  2. 方法二:Shell 脚本 (使用 amqp-consume 或者 rabbitmqadmin)
  3. 方法三:Node.js (JavaScript)
  4. 方法四:Ruby
  5. 常见注意事项
  6. 推荐方案

要用脚本消费 RabbitMQ 消息,主要取决于你使用的编程语言,以下是几种常用语言的脚本示例,包括 Python (推荐)Shell (配合工具)Node.js

核心思路

  1. 连接:连接到 RabbitMQ 服务器(IP、端口、vhost、凭据)。
  2. 声明/绑定:确保队列存在,并将队列绑定到交换机。
  3. 消费:监听消息并处理(打印、写入文件、调用API等)。
  4. 确认:处理成功后发送 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)

需要安装 amqplibnpm 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

推荐方案

  • 临时调试 / 一次运行:用 PythonNode.js 脚本最直接。
  • 集成到自动化流水线(如CI/CD):Python 库 pika 最稳定。
  • 简单的监控/抓取rabbitmqadmin getamqp-consume 快速查看。

如果有特定的业务逻辑(比如把消息写入文件、转发到HTTP接口),在回调函数里加上相应代码即可。

抱歉,评论功能暂时关闭!