Python后端如何实现RabbitMQ队列的消息读取与写入操作?
Python 操作 RabbitMQ 消息读写实现方案
前置准备
首先安装Python对应的RabbitMQ客户端依赖:pip install pika
请确保你已部署并启动RabbitMQ服务,默认服务端口为5672,默认账号密码为guest/guest,生产环境请自行修改为自定义账号密码。
1. 消息写入(生产者实现)
生产者负责往RabbitMQ队列中写入消息,示例代码如下:
import pika # 配置连接参数 credentials = pika.PlainCredentials('guest', 'guest') connection = pika.BlockingConnection( pika.ConnectionParameters( host='localhost', port=5672, virtual_host='/', credentials=credentials ) ) channel = connection.channel() # 声明队列,durable=True表示队列持久化,RabbitMQ重启后队列不会丢失 channel.queue_declare(queue='business_queue', durable=True) # 写入消息 msg = "业务测试消息" channel.basic_publish( exchange='', # 使用默认交换机 routing_key='business_queue', # 路由键和队列名匹配即可发送到对应队列 body=msg.encode('utf-8'), properties=pika.BasicProperties( delivery_mode=2 # 消息持久化,RabbitMQ重启后消息不会丢失 ) ) print(f"消息已写入队列:{msg}") # 关闭连接 connection.close()
2. 消息读取(消费者实现)
消费者负责从RabbitMQ队列中读取并处理消息,示例代码如下:
import pika def msg_process_callback(ch, method, properties, body): """消息处理回调函数,收到消息后自动触发""" recv_msg = body.decode('utf-8') print(f"读取到队列消息:{recv_msg}") # 手动确认消息已处理,RabbitMQ会删除该消息 ch.basic_ack(delivery_tag=method.delivery_tag) # 建立连接,参数和生产者保持一致 credentials = pika.PlainCredentials('guest', 'guest') connection = pika.BlockingConnection( pika.ConnectionParameters( host='localhost', port=5672, virtual_host='/', credentials=credentials ) ) channel = connection.channel() channel.queue_declare(queue='business_queue', durable=True) # 配置预取数,每次只拉取1条未处理消息,避免消费端负载过高 channel.basic_qos(prefetch_count=1) # 绑定回调函数到对应队列 channel.basic_consume(queue='business_queue', on_message_callback=msg_process_callback) print("消费者已启动,等待接收消息...") # 启动持续监听 channel.start_consuming()
注意事项
- 队列声明是幂等操作,重复执行不会覆盖已存在的队列,仅当队列参数和已有队列不一致时才会抛出异常
- 生产环境建议添加连接异常重试逻辑,避免RabbitMQ服务波动导致业务中断
- 如果开启了手动ack确认,一定要确保消息处理完成后发送ack信号,否则消息会一直留在队列中,消费者重启后会重新消费
- 复杂业务场景可以使用交换机(Exchange)实现广播、路由匹配、主题匹配等更灵活的消息分发规则
内容的提问来源于stack exchange,提问作者BeefNoodle
相关产品推荐
相关产品推荐

