提交日志信息,发送给所有的消费者。
使用扇形交换机 实现
[root@jinkang-e2elog rabbitmq]# cat emit_log.py #!/usr/bin/env python import pika import sys connection = pika.BlockingConnection(pika.ConnectionParameters( host=‘localhost‘)) channel = connection.channel() channel.exchange_declare(exchange=‘logs‘, exchange_type=‘fanout‘) message = ‘ ‘.join(sys.argv[1:]) or "info: Hello World!" channel.basic_publish(exchange=‘logs‘, routing_key=‘‘, body=message) print " [x] Sent %r" % (message,) connection.close()
[root@jinkang-e2elog rabbitmq]# cat recive_log.py #!/usr/bin/env python import pika connection = pika.BlockingConnection( pika.ConnectionParameters(host=‘localhost‘)) channel = connection.channel() channel.exchange_declare(exchange=‘logs‘, exchange_type=‘fanout‘) result = channel.queue_declare(queue=‘‘, 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( queue=queue_name, on_message_callback=callback, auto_ack=True) channel.start_consuming()
原文:https://www.cnblogs.com/jkklearn/p/13786756.html