Python RabbitMQ消费者多线程处理脚本无响应问题求助
问题分析与修复方案
我来帮你排查这个RabbitMQ消费者的问题,你的代码里有几个关键逻辑错误,导致程序执行后没有任何反应:
1. 消费线程未启动消息监听
你的runq方法仅仅把channel赋值给self.channel,但没有调用read_queue()启动RabbitMQ的消息消费循环。这意味着消费线程根本没开始从队列拉取消息,自然不会有后续的消息处理动作。
2. 错误地在run方法中重新实例化con对象
在run方法里,你又创建了一个新的cn = con(),并用这个新对象获取连接和channel,但这个新对象和当前的self实例完全无关,后续self.channel并没有关联到正确的连接通道,消息消费逻辑等于完全没绑定到有效连接上。
3. 线程启动后立刻调用join()导致主线程阻塞
你启动线程后马上调用t1.join()和t2.join(),这会让主线程等待这两个线程执行完毕才继续。但t1执行完runq就结束了,t2会卡在q.get()(因为没有消息被放入队列),加上你把线程设为守护线程,最终程序就会陷入无响应的状态。
4. run方法中未定义q变量
在run方法创建线程时,你用到了args=(channel, q),但run方法内部并没有定义q变量——虽然你在main里创建了q,但没有传递给run方法,这里本该抛出NameError,可能是你粘贴代码时的疏漏。
修复后的完整代码
下面是修正后的代码,我标注了所有关键改动:
import pika, Consumer_Config, queue, threading import Message_Class class con: def __init__(self): self.config = Consumer_Config._config() self.path = self.config.path self.active = 0 self.channel = None self.tag = None self.tb = None def channsel(self): pika_conn_params = pika.ConnectionParameters( host=self.config.url, port=self.config.port, credentials=pika.credentials.PlainCredentials(self.config.user_id, self.config.password)) connection = pika.BlockingConnection(pika_conn_params) return connection # 改动1:让消费线程启动消息监听 def runq(self, channel, q): self.channel = channel # 调用read_queue开始监听RabbitMQ队列 self.read_queue(q) # 改动2:新增q参数,传递给消息回调 def read_queue(self, q): queue = self.channel.queue_declare( queue="queue", durable=True, exclusive=False, auto_delete=False) self.channel.queue_bind( "queue", 'exchange', routing_key=str('text')) self.channel.basic_qos( prefetch_count=500) # 把q传递给on_msg回调,方便放入消息队列 self.tag = self.channel.basic_consume("queue", lambda ch, method, props, body: self.on_msg(ch, method, props, body, q)) self.channel.start_consuming() # 改动3:新增q参数,收到消息后放入队列 def on_msg(self, _unused_channel, basic_deliver, properties, body, q): rk = basic_deliver.routing_key self.tb = body q.put(self.tb) self.acknowledge_message(basic_deliver.delivery_tag) def acknowledge_message(self, delivery_tag): """Acknowledge the message delivery from RabbitMQ by sending a Basic.Ack RPC method for the delivery tag. """ self.channel.basic_ack(delivery_tag) # 改动4:循环处理队列消息,避免线程执行一次就结束 def printq(self, q): while True: tb = q.get() try: item = Message_Class.Messages().Process_msg(tb, 'text') print(item) finally: # 改动5:标记任务完成,让q.join()能正确等待所有消息处理完毕 q.task_done() # 改动6:修复run方法的逻辑,正确传递参数和管理线程 def run(self, q): # 用当前实例self获取连接,不要重新创建新的con实例 connct = self.channsel() channel = connct.channel() num_worker_threads = 2 # 启动单独的消费线程 consume_thread = threading.Thread(target=self.runq, args=(channel, q)) consume_thread.daemon = True consume_thread.start() # 启动多个工作线程处理队列消息 for i in range(num_worker_threads): worker_thread = threading.Thread(target=self.printq, args=(q,)) worker_thread.daemon = True worker_thread.start() # 等待队列所有消息处理完成 q.join() if __name__ == "__main__": q = queue.Queue() r = con() # 改动7:把q传递给run方法 r.run(q)
关键改动说明
- 启动消息消费循环:在
runq里调用read_queue(),让消费线程真正开始监听RabbitMQ队列。 - 传递队列到消息回调:通过lambda表达式把队列
q传递给on_msg,确保收到消息后能正确放入本地队列。 - 工作线程持续处理:给
printq加入while True循环,让工作线程持续从队列取消息处理,而不是只执行一次。 - 正确使用当前实例:用
self获取连接,避免创建多余的con实例,保证channel绑定到当前对象的有效连接。 - 传递队列到run方法:在
main里把q传给r.run(q),解决变量未定义的问题。 - 标记任务完成:在
printq的finally块调用q.task_done(),确保q.join()能正确等待所有消息处理完毕。
这样修改后,你的消费者应该能正常拉取RabbitMQ消息,并在独立线程中处理数据了。
内容的提问来源于stack exchange,提问作者kabik
相关产品推荐
相关产品推荐

