You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何实现兼具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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.12 00:21:09