【问题标题】:Python Pika - Consumer into ThreadPython Pika - 消费者进入线程
【发布时间】:2013-05-03 12:52:53
【问题描述】:

我正在开发一个带有后台线程的 Python 应用程序,用于消费来自 RabbitMQ 队列的消息(主题场景)。

我在按钮的 on_click 事件上启动线程。 这是我的代码,请注意“#self.receive_command()”。

def on_click_start_call(self,widget):


    t_msg = threading.Thread(target=self.receive_command)
    t_msg.start()
    t_msg.join(0)
    #self.receive_command()


def receive_command(self):

    syslog.syslog("ENTERED")

    connection = pika.BlockingConnection(pika.ConnectionParameters(host='localhost'))
    syslog.syslog("1")

    channel = connection.channel()
    syslog.syslog("2")

    channel.exchange_declare(exchange='STORE_CMD', type='topic')
    syslog.syslog("3")

    result = channel.queue_declare(exclusive=True)
    syslog.syslog("4")

    queue_name = result.method.queue
    syslog.syslog("5")

    def callback_rabbit(ch,method,properties,body):
        syslog.syslog("RICEVUTO MSG: RKEY:"+method.routing_key+" MSG: "+body+"\n")

    syslog.syslog("6")

    channel.queue_bind(exchange='STORE_CMD', queue=queue_name , routing_key='test.routing.key')
    syslog.syslog("7")

    channel.basic_consume(callback_rabbit,queue=queue_name,no_ack=True)
    syslog.syslog("8")

    channel.start_consuming()

如果我运行此代码,我无法在 syslog 上看到消息 1、2、3、5、6、7、8 但我只能看到“ENTERED”。因此,代码被锁定在 pika.BlokingConnection 上。

如果我运行相同的代码(注释线程指令并取消对函数的直接调用),所有工作都按预期进行,并且消息被正确接收。

有什么解决方案可以让消费者进入线程?

提前致谢

大卫

【问题讨论】:

    标签: python multithreading rabbitmq pika


    【解决方案1】:

    我已经在我的机器上使用最新版本的 Pika 测试了代码。它工作正常。 Pika 存在线程问题,但只要您为每个线程创建一个连接,这应该不是问题。

    如果您遇到问题,很可能是因为旧版 Pika 中的错误,或者与您的线程无关的问题导致了问题。

    我建议您避免使用 0.9.13,因为存在多个错误,但 0.9.14 0.10.0 应该很快就会发布™。

    [编辑] Pika 0.9.14 已经发布。

    这是我使用的代码。

    def receive_command():
        print("ENTERED")
        connection = pika.BlockingConnection(pika.ConnectionParameters(host='localhost'))
        print("1")
        channel = connection.channel()
        print("2")
        channel.exchange_declare(exchange='STORE_CMD', type='topic')
        print("3")
        result = channel.queue_declare(exclusive=True)
        print("4")
        queue_name = result.method.queue
        print("5")
        def callback_rabbit(ch,method,properties,body):
            print("RICEVUTO MSG: RKEY:"+method.routing_key+" MSG: "+body+"\n")
        print("6")
        channel.queue_bind(exchange='STORE_CMD', queue=queue_name , routing_key='test.routing.key')
        print("7")
        channel.basic_consume(callback_rabbit,queue=queue_name,no_ack=True)
        print("8")
        channel.start_consuming()
    
    def start():
        t_msg = threading.Thread(target=receive_command)
        t_msg.start()
        t_msg.join(0)
        #self.receive_command()
    start()
    

    【讨论】:

    • 感谢分享。如果我们添加了第二个线程(单独的连接,单独的队列等,所以定义没有并发问题) - 你将如何优雅地处理应用程序关闭?例如,仅发出 control-C(键盘中断)是否有任何危害 - 或者这是否会导致两个线程接收相同的关闭序列。是否有任何finally 逻辑等来关闭父进程(主)所需的兔子连接?
    • 我查看了兔子邮件列表 (groups.google.com/forum/#!forum/rabbitmq-users),但尚未找到答案。你上面的代码看起来非常接近我需要的。只是在关闭部分寻找确认。
    • 最后一个问题 - 你对这个解决方案有什么想法(使用进程而不是线程) - stackoverflow.com/a/45142386/1882064 多进程的任何关闭细节?
    【解决方案2】:

    另一种方法是将线程方法channel.start_consuming 作为目标传递,然后将回调传递给consume 方法。 用法:consume(callback=your_method, queue=your_queue)

    import threading
    
    def consume(self, *args, **kwargs):
        if "channel" not in kwargs \
                or "callback" not in kwargs \
                or "queue" not in kwargs \
                or not callable(kwargs["callback"]):
            return None
    
        channel = kwargs["channel"]
        callback = kwargs["callback"]
        queue = kwargs["queue"]
        channel.basic_consume(callback, queue=queue, no_ack=True)
    
        t1 = threading.Thread(target=channel.start_consuming)
        t1.start()
        t1.join(0)
    

    【讨论】:

    • 这对我有用,可以让我的消费者作为 Flask 应用程序初始化的一部分运行。谢谢!
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多