如何实现兼具Producer与Consumer功能的Python RabbitMQ微服务?
能否用Python创建同时作为RabbitMQ生产者和消费者的微服务?
当然可以,这种模式在很多业务场景中都很实用——比如服务需要接收消息处理后,再将结果转发到另一个队列的场景。核心思路是通过多线程或异步IO分离生产者和消费者的逻辑,避免两者互相阻塞。
多线程实现示例
用Python的threading模块拆分消费者和生产者逻辑,适合简单场景:
import pika import threading import time def consumer(): # 建立消费者连接 connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() channel.queue_declare(queue='input_queue') def callback(ch, method, properties, body): print(f"收到消息: {body.decode()}") # 处理消息后,调用生产者逻辑转发结果 producer(f"{body.decode()} - 已处理") # 手动确认消息已消费 ch.basic_ack(delivery_tag=method.delivery_tag) channel.basic_consume(queue='input_queue', on_message_callback=callback) print("消费者已启动,等待消息...") channel.start_consuming() def producer(message): # 建立生产者连接并发送消息 connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() channel.queue_declare(queue='output_queue') channel.basic_publish(exchange='', routing_key='output_queue', body=message) print(f"发送处理后的消息: {message}") connection.close() if __name__ == "__main__": # 启动消费者后台线程 consumer_thread = threading.Thread(target=consumer) consumer_thread.daemon = True consumer_thread.start() # 主线程模拟用户输入触发生产者逻辑 while True: user_input = input("输入要发送的消息(输入q退出):") if user_input.lower() == 'q': break producer(user_input)
异步IO实现示例
用aio-pika库基于asyncio实现,更适合高并发场景,避免线程切换开销:
import asyncio import aio_pika async def consumer(): connection = await aio_pika.connect_robust("amqp://localhost/") async with connection: channel = await connection.channel() queue = await channel.declare_queue("input_queue") async with queue.iterator() as queue_iter: async for message in queue_iter: async with message.process(): received_msg = message.body.decode() print(f"收到消息: {received_msg}") # 处理后转发到输出队列 await producer(f"{received_msg} - 已处理") async def producer(message): connection = await aio_pika.connect_robust("amqp://localhost/") async with connection: channel = await connection.channel() await channel.declare_queue("output_queue") await channel.default_exchange.publish( aio_pika.Message(body=message.encode()), routing_key="output_queue" ) print(f"发送处理后的消息: {message}") async def main(): # 启动消费者异步任务 consumer_task = asyncio.create_task(consumer()) # 模拟用户输入触发生产者 while True: user_input = await asyncio.to_thread(input, "输入要发送的消息(输入q退出):") if user_input.lower() == 'q': consumer_task.cancel() break await producer(user_input) if __name__ == "__main__": asyncio.run(main())
关键注意事项
- 连接复用:尽量复用RabbitMQ连接,避免每次发送/接收都创建新连接,减少资源消耗。
- 错误处理:需添加连接断开重连、消息重试、异常捕获等逻辑,保证服务稳定性。
- 消息确认:消费者必须正确调用
basic_ack(多线程)或message.process()(异步),避免消息重复消费。 - 资源隔离:多线程模式下注意共享资源的线程安全,异步模式则无需担心此问题。
内容的提问来源于stack exchange,提问作者flying_loaf_3
相关产品推荐
相关产品推荐

