RabbitMQ 问题与修复
RabbitMQ 核心组成 1. 消息(Message) 消息是 RabbitMQ 中的基本单位,通常包含数据(有效负载)和一些元数据(如路由键、标头等)。 2. 生产者(Producer) 生产者是…
RabbitMQ 核心组成
- 消息(Message)
消息是 RabbitMQ 中的基本单位,通常包含数据(有效负载)和一些元数据(如路由键、标头等)。
- 生产者(Producer)
生产者是发送消息到 RabbitMQ 的应用程序或服务。它将消息发布到交换机。
- 消费者(Consumer)
消费者是接收和处理来自 RabbitMQ 消息的应用程序或服务。它从队列中读取消息并进行处理。
- 交换机(Exchange)
交换机是 RabbitMQ 中的核心组件,负责接收来自生产者的消息并根据路由规则将其转发到一个或多个队列。交换机的类型有多种:
- 直接交换机(Direct Exchange):根据消息的路由键将消息发送到绑定的队列。
- 主题交换机(Topic Exchange):根据路由模式将消息发送到一个或多个队列,支持使用通配符。
- 扇出交换机(Fanout Exchange):将接收到的每条消息广播到所有绑定的队列。
- 头交换机(Headers Exchange):基于消息头中的属性进行路由。
- 队列(Queue)
队列是存储消息的缓冲区,消费者从队列中获取消息进行处理。消息在队列中的顺序通常是 FIFO(先进先出)。
- 绑定(Binding)
绑定是交换机与队列之间的关系。通过绑定,交换机知道将消息路由到哪些队列。绑定可以指定路由键或匹配模式。
- 路由键(Routing Key)
路由键是生产者在发送消息时指定的标识符,交换机使用它来决定将消息发送到哪个队列。
- 确认(Acknowledgment)
消费者在处理完消息后需要发送确认,告诉 RabbitMQ 消息已经成功处理。如果没有确认,RabbitMQ 会将消息重新放回队列,供其他消费者处理。
- 持久性(Persistence)
RabbitMQ 提供消息的持久化功能,即使 RabbitMQ 服务器崩溃,消息也不会丢失。可以通过将队列和消息标记为持久化来实现。
- 管理界面(Management Interface)
RabbitMQ 提供了一个用户友好的 Web 界面,用于监控和管理消息队列、交换机、绑定和消费者等。
RabbitMQ 与 Python 通信
1@app.get("/publish")
2async def publish():
3 # 连接到RabbitMQ服务器
4 connection = pika.BlockingConnection(pika.ConnectionParameters('localhost', 5672))
5 # 获取通道
6 channel = connection.channel()
7 # 创建自定义交换机
8 channel.exchange_declare(exchange='order.direct.exchange', exchange_type='direct')
9 # 声明队列
10 channel.queue_declare(queue='order-queue')
11 # 将队列绑定到交换机上,使用 routing_key='order.created.key'
12 channel.queue_bind(exchange='order.direct.exchange', queue='order-queue', routing_key='order.created.key')
13 # 发送消息到交换机
14 channel.basic_publish(exchange='order.direct.exchange',
15 routing_key='order.created.key',
16 body='Hello RabbitMQ!')
17 print("消息已发送:Hello RabbitMQ!")
18 # 关闭连接
19 connection.close()
20 return True
1@app.get("/consume")
2async def consume():
3 # 创建并启动独立的消费线程
4 consumer_thread = Thread(target=consume_messages)
5 consumer_thread.start()
6 return {"message": "已启动消费线程"}
7
8def consume_messages():
9 # 连接到RabbitMQ服务器
10 connection = pika.BlockingConnection(pika.ConnectionParameters('localhost', 5672))
11 channel = connection.channel()
12 # 声明队列
13 channel.queue_declare(queue='order-queue')
14
15 # 定义回调函数,用于处理接收到的消息
16 def callback(ch, method, properties, body):
17 print(f"接收到消息:{body.decode()}")
18
19 # 启动消费者,并指定回调函数
20 channel.basic_consume(queue='order-queue', on_message_callback=callback, auto_ack=True)
21 # 持续监听并消费消息
22 channel.start_consuming()
23 # 停止消费
24 # channel.stop_consuming()
25 # 删除没有使用且为空的队列
26 channel.queue_delete(queue='order-queue', if_unused=True, if_empty=True)
RabbitMQ 与 Java 通信
1@RestController
2@RequestMapping("/api/products")
3public class MessageProducer {
4
5 @Autowired
6 private RabbitTemplate rabbitTemplate;
7
8 private final String exchange = "my_exchange";
9 private final String routingKey = "my_routing_key";
10
11 @GetMapping("/sendMessage")
12 public String sendMessage() {
13 String message = "Hello World";
14 rabbitTemplate.convertAndSend(exchange, routingKey, message);
15 return "SUCCESS";
16 }
17}
1@Service
2public class MessageConsumer {
3
4 @RabbitListener(queues = "my_queue")
5 public void receiveMessage(String message) {
6 System.out.println("Received: " + message);
7 }
8}
1@Configuration
2public class RabbitConfig {
3 // 使用 Java 配置类来声明交换机、队列和绑定关系
4 @Bean
5 public TopicExchange exchange() {
6 return new TopicExchange("my_exchange");
7 }
8
9 @Bean
10 public Queue queue() {
11 return new Queue("my_queue");
12 }
13
14 @Bean
15 public Binding binding(Queue queue, TopicExchange exchange) {
16 return BindingBuilder.bind(queue).to(exchange).with("my_routing_key");
17 }
18}