1. 程式人生 > >RabbitMQ訊息通訊,一個生產者和多個消費者,廣播式訊息通訊

RabbitMQ訊息通訊,一個生產者和多個消費者,廣播式訊息通訊

上一則我們說到了一個對多個的RabbitMQ訊息佇列通訊的實現方法,生產者傳送的訊息只能被一個消費者接收並處理,上則請閱讀:http://blog.csdn.net/u012631731/article/details/78450389
本則說的是廣播式的訊息通訊方法實現,所有的消費者都可以收到生產者傳送的訊息


還是直接上程式碼吧,有描述直接在程式碼裡面註釋:
client.py
#!/usr/bin/env python
import pika
import sys
#不解釋
connection = pika.BlockingConnection(pika.ConnectionParameters(host='localhost'))
channel = connection.channel()
#這裡是設定訊息佇列的一些屬性,之前兩則文章都沒有設定exchange這個屬性,直接賦值為空字串
#這裡需要描述的是,並不是生產者直接傳送訊息到訊息佇列裡面,生產者把訊息傳送給了exchange,
#通過exchange再把訊息傳送給某一個或多個訊息佇列裡面(queue),這裡沒有建立訊息佇列,因為這個事例是表達廣播式訊息佇列,
#有一點需要說明一下,如果exchange屬性設定為空,RabbitMQ就會採用預設的屬性設定exchange,然後需要設定
#routing_key屬性,指明所要傳送的訊息佇列名稱,這裡設定exchange的名稱為logs,型別屬性為fanout,還有另外三個屬性direct, topic, headers,後續再說明
channel.exchange_declare(exchange='logs',
                         exchange_type='fanout')
message = ' '.join(sys.argv[1:]) or "info: Hello World!"
#指定exchange的名稱
channel.basic_publish(exchange='logs',
                      routing_key='',
                      body=message)
print(" [x] Sent %r" % message)
connection.close()


server.py
#!/usr/bin/env python
import pika


connection = pika.BlockingConnection(pika.ConnectionParameters(host='localhost'))
channel = connection.channel()
#建立exchange的名稱為logs,指定型別為fanout
channel.exchange_declare(exchange='logs',
                         exchange_type='fanout')
#刪除隨機建立的訊息佇列
result = channel.queue_declare(exclusive=True)
queue_name = result.method.queue
channel.queue_bind(exchange='logs',
                   queue=queue_name)
print(' [*] Waiting for logs. To exit press CTRL+C')
def callback(ch, method, properties, body):
    print(" [x] %r" % body)
channel.basic_consume(callback,
                      queue=queue_name,
                      no_ack=True)
channel.start_consuming()


執行多個python server.py 和一個python client.py看看效果吧
更多資訊請查詢RabbitMQ官網:http://www.rabbitmq.com