【问题标题】:Python Paho client how to consume from RabbitMQ existing queuePython Paho 客户端如何从 RabbitMQ 现有队列中消费
【发布时间】:2019-07-20 15:08:00
【问题描述】:

我有一个 RabittMQ 队列 'test',我必须使用 Paho 客户端编写一个 python 消费者来仅使用来自这个 ('test') 队列的消息。

下面是我的消费者代码:

import paho.mqtt.client as mqtt
import json

server = "<machine_ip>"
port = 1883
username = "test"
password = "test"
#########################
client_id = "test_consumer"
topic = "test"

# The callback for when the client receives a CONNACK response from the server.
def on_connect(client, userdata, flags, rc):
    print("Connected with result code "+str(rc))
    client.subscribe(topic)

# The callback for when a PUBLISH message is received from the server.
def on_message(client, userdata, msg):
 data = str(msg.payload.decode());
 print(data)

client = mqtt.Client(client_id)
client.on_connect = on_connect
client.on_message = on_message
client.username_pw_set(username,password)
client.connect(server, port, 60, bind_address="")
client.loop_forever()

但是当我开始使用消费者时,它正在创建一个新队列,如附加的屏幕截图所示,这是我不想要的。有什么解决办法吗?

【问题讨论】:

  • 执行客户端时是否会创建一个新队列?

标签: python rabbitmq mqtt paho


【解决方案1】:

我不知道为什么 rabbitmq 为每个客户端(发布者或订阅者/消费者)创建一个队列,但是如果您使用下面的参数,所有客户端将使用相同的队列:

  1. 首先,您需要连接和订阅消费者: clean_session=假, 并使用 qos=1 订阅主题(例如“文本”)

  2. 您可以断开消费者(如果需要,模拟网络断开)

  3. 接下来,将您的发布者与: clean_session=假 并以 qos=1,retain=False 发布到同一主题('test')

【讨论】:

    【解决方案2】:
    1. 创建虚拟主机'benz'

    2. 使用路由键“eff41a11134d4fb9bd0325866b8a983d”创建队列“test

    3. 创建新用户

    注意 - 您必须为此创建一个新用户,RabittMQ 将无法使用以下生产者和消费者的默认用户名/密码。

    producer.py -

    import paho.mqtt.client as mqtt
    
    ip = '127.0.0.1'
    port = 1883
    vhost = "benz"
    routing_key= "eff41a11134d4fb9bd0325866b8a983d"
    username = 'username'
    password = 'password'
    
    client = mqtt.Client()
    client.username_pw_set(vhost + ':' + username, password)
    client.connect(ip, port, 60)
    
    while True:
        console_input = str(input())
        print (console_input)
        client.publish(routing_key, payload=console_input, qos=0)
    

    consumer.py

    import pika
    
    from demjson import decode
    
    ## RabittMQ configuration properties ##
    ip = '127.0.0.1'
    vhost = 'benz'
    username = 'username'
    password = 'password'
    queue_name = 'test'
    
    url = 'amqp://' + username + ':' + password + '@' + ip + ':5672/' + vhost
    params = pika.URLParameters(url)
    connection = pika.BlockingConnection(params)
    channel = connection.channel()
    channel.queue_declare(queue=queue_name, durable=True)
    
    def callback(ch, method, properties, body):
        data = str(body.decode());
        print(data)
    
    channel.basic_consume(queue_name, callback, auto_ack=True)
    channel.start_consuming()
    connection.close()
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2010-11-05
      • 2020-08-30
      相关资源
      最近更新 更多