RabbitMQ: 消息代理中间件
RabbitMQ是一个开源的消息代理和队列服务器,用于通过整个企业中的分布式系统传递消息,它支持多种消息传递协议,并且可以用于跨多种应用和多种不同的操作系统平台。
以下是一些RabbitMQ的常见用法和代码示例:
- 消息队列:
import pika
# 连接到RabbitMQ服务器
connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()
# 声明队列
channel.queue_declare(queue='hello')
# 定义回调函数来处理消息
def callback(ch, method, properties, body):
print(f" Received {body}")
# 开始监听队列,并处理消息
channel.basic_consume(queue='hello', on_message_callback=callback, auto_ack=True)
print(' [*] Waiting for messages. To exit press CTRL+C')
channel.start_consuming()
- 发布/订阅模式:
import pika
# 连接到RabbitMQ服务器
connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()
# 声明交换器
channel.exchange_declare(exchange='logs', exchange_type='fanout')
# 回调函数来处理消息
def callback(ch, method, properties, body):
print(f" Received {body}")
# 启动监听,并处理消息
channel.basic_consume(queue=queue_name, on_message_callback=callback, auto_ack=True)
print(' [*] Waiting for messages. To exit press CTRL+C')
channel.start_consuming()
- 路由模式:
import pika
# 连接到RabbitMQ服务器
connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()
# 声明交换器
channel.exchange_declare(exchange='direct_logs', exchange_type='direct')
# 回调函数来处理消息
def callback(ch, method, properties, body):
print(f" Received {body}")
# 定义队列
queue_name = channel.queue_declare(exclusive=True).method.queue
# 绑定交换器和队列
severities = ['error', 'info', 'warning']
for severity in severities:
channel.queue_bind(exchange='direct_logs', queue=queue_name, routing_key=severity)
# 启动监听,并处理消息
channel.basic_consume(queue=queue_name, on_message_callback=callback, auto_ack=True)
print(' [*] Waiting for messages. To exit press CTRL+C')
channel.start_consuming()
- RPC(远程过程调用):
import pika
import uuid
# 连接到RabbitMQ服务器
connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()
# 声明一个回调函数来处理RPC响应
def on_response(ch, method, properties, body):
if properties.correlation_id == correlation_id:
print(f" Received {body}")
# 声明一个回调函数来处理RPC请求
def on_request(ch, method, properties, body):
print(f" Received {body}")
# 处理请求...
response = b"Response to the request"
评论已关闭