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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 06:49:59