如何在双线程中运行Pika消费者与生产者并正确停止关闭?
解决Pika多线程下停止消费者并关闭通道/连接的AssertionError问题
问题分析
你遇到的AssertionError: ('_AsyncTransportBase._initate_abort() expected non-_STATE_COMPLETED', 4)错误,主要源于以下几点:
- 消费者线程未启动消费循环,后续停止操作无有效目标
- 多线程环境下未按正确顺序关闭资源,且未保证线程安全
- 停止消费和关闭操作的时机不对,导致连接状态异常
修复方案及代码示例
关键修复点
- 消费者线程必须启动消费循环,同时预留停止信号检测逻辑
- 通过线程事件对象实现安全的线程间通信,避免跨线程直接操作通道
- 严格遵循停止消费 → 关闭通道 → 关闭连接的资源释放顺序
修改后的完整代码
import time import pika import threading url = 'amqp://user:user@localhost:5672/' # 修正拼写错误:localhosh → localhost def run_consumer(connection, tag, stop_event): def on_message(channel, method, properties, body): print(f"收到消息: {body.decode()}") channel.basic_ack(delivery_tag=method.delivery_tag) # 手动确认消息,避免堆积 channel = connection.channel() channel.queue_declare(queue='queue', auto_delete=True) channel.basic_consume( queue='queue', auto_ack=False, on_message_callback=on_message, consumer_tag=tag ) # 循环处理事件,同时检测停止信号 while not stop_event.is_set(): try: # 限时处理事件,避免永久阻塞 channel.process_data_events(time_limit=1) except (pika.exceptions.ConnectionClosedByBroker, pika.exceptions.AMQPChannelError): break # 取消消费者并关闭通道 channel.basic_cancel(consumer_tag=tag) channel.close() print("消费者通道已关闭") def run_publisher(connection, stop_event): channel = connection.channel() time.sleep(5) # 缩短等待时间,便于测试 for i in range(3): channel.basic_publish( exchange='', routing_key='queue', body=f'Test message {i+1}'.encode(), ) print(f"发送消息: Test message {i+1}") time.sleep(2) # 发送完成后触发停止信号 stop_event.set() channel.close() print("生产者通道已关闭") if __name__ == '__main__': parameters = pika.URLParameters(url) connection = pika.BlockingConnection(parameters) consumer_tag = 'consumer' # 创建事件对象用于线程间停止信号传递 stop_event = threading.Event() thread_consumer = threading.Thread( target=run_consumer, args=(connection, consumer_tag, stop_event), name='Consumer' ) thread_consumer.start() thread_publisher = threading.Thread( target=run_publisher, args=(connection, stop_event), name='Publisher' ) thread_publisher.start() # 等待所有线程执行完毕 thread_consumer.join() thread_publisher.join() # 最后关闭连接 if connection.is_open: connection.close() print("连接已关闭")
核心说明
- 线程安全控制:用
threading.Event作为停止信号,避免跨线程直接调用通道方法引发状态冲突 - 消费循环优化:使用
process_data_events(time_limit=1)替代start_consuming(),让消费者线程能定期检查停止信号,而非永久阻塞 - 资源释放顺序:生产者完成发送后触发停止信号,消费者收到信号后先取消消费者、关闭通道,最后主线程关闭连接
- 消息可靠性:添加
basic_ack手动确认消息,防止队列中堆积未确认的消息
内容的提问来源于stack exchange,提问作者AlexanderKondrat
相关产品推荐
相关产品推荐

