MQ 核心概念:生产-消费模型 + 主题/队列 + ACK 机制
一句话:消息队列(MQ)是异步通信的中间件,通过生产-消费模型实现系统解耦。核心概念包括生产者、消费者、主题/队列、ACK机制等。
1. MQ基础
1.1 什么是消息队列?
graph LR A[生产者] -->|发送消息| B[消息队列] B -->|消费消息| C[消费者] style B fill:#e1f5fe
消息队列:异步通信的中间件,允许生产者发送消息到队列,消费者从队列接收消息。
1.2 核心概念
| 概念 | 说明 |
|---|---|
| Producer | 生产者,发送消息 |
| Consumer | 消费者,接收消息 |
| Queue | 队列,存储消息 |
| Topic | 主题,消息分类 |
| ACK | 确认机制,确认消息已处理 |
| Broker | 消息代理,管理消息存储和转发 |
2. 生产-消费模型
2.1 点对点模型
graph LR A[生产者] -->|消息| B[队列] B -->|消息| C[消费者1] style B fill:#e8f5e8
点对点:一个消息只被一个消费者处理。
2.2 发布-订阅模型
graph TD A[生产者] -->|消息| B[主题] B -->|消息| C[消费者1] B -->|消息| D[消费者2] B -->|消息| E[消费者3] style B fill:#e1f5fe
发布-订阅:一个消息被多个消费者处理。
3. 主题/队列
3.1 主题(Topic)
# 生产者发送消息到主题
producer.send("user-events", message)
# 多个消费者订阅同一主题
consumer1.subscribe("user-events")
consumer2.subscribe("user-events")
3.2 队列(Queue)
# 生产者发送消息到队列
producer.send("task-queue", message)
# 只有一个消费者处理消息
consumer.receive("task-queue")
4. ACK机制
4.1 什么是ACK?
graph LR A[消费者] -->|处理消息| B[消息队列] B -->|ACK| A style A fill:#e8f5e8
ACK:确认机制,消费者处理完消息后发送确认,告诉消息队列可以删除消息。
4.2 自动ACK vs 手动ACK
# 自动ACK:收到消息自动确认
consumer.receive("queue", auto_ack=True)
# 手动ACK:处理完手动确认
message = consumer.receive("queue", auto_ack=False)
# 处理消息
process_message(message)
# 手动确认
message.ack()
5. 实际案例
5.1 Kafka示例
from kafka import KafkaProducer, KafkaConsumer
import json
# 生产者
producer = KafkaProducer(
bootstrap_servers=['localhost:9092'],
value_serializer=lambda v: json.dumps(v).encode('utf-8')
)
# 发送消息
producer.send('user-events', {'user_id': 123, 'action': 'login'})
# 消费者
consumer = KafkaConsumer(
'user-events',
bootstrap_servers=['localhost:9092'],
value_deserializer=lambda m: json.loads(m.decode('utf-8')),
auto_offset_reset='earliest',
enable_auto_commit=True
)
# 消费消息
for message in consumer:
print(message.value)
5.2 RabbitMQ示例
import pika
import json
# 连接
connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()
# 声明队列
channel.queue_declare(queue='task_queue', durable=True)
# 生产者
channel.basic_publish(
exchange='',
routing_key='task_queue',
body=json.dumps({'task': 'process_data'}),
properties=pika.BasicProperties(
delivery_mode=2, # 消息持久化
)
)
# 消费者
def callback(ch, method, properties, body):
print(f"收到消息: {body}")
# 处理消息
ch.basic_ack(delivery_tag=method.delivery_tag)
channel.basic_consume(queue='task_queue', on_message_callback=callback)
channel.start_consuming()
6. 常见坑点
1. 消息丢失
# 解决:使用持久化消息和ACK机制
properties = pika.BasicProperties(delivery_mode=2) # 持久化
2. 重复消费
# 解决:实现幂等性处理
def process_message(message):
if is_processed(message['id']):
return # 已处理,跳过
# 处理消息
mark_as_processed(message['id'])
3. 消息积压
# 或者使用消息分片
核心要点
# Kafka
producer = KafkaProducer(bootstrap_servers=['localhost:9092'])
producer.send('topic', {'key': 'value'})
consumer = KafkaConsumer('topic', bootstrap_servers=['localhost:9092'])
for message in consumer:
process(message.value)
# RabbitMQ
channel.queue_declare(queue='queue', durable=True)
channel.basic_publish(exchange='', routing_key='queue', body='message')
channel.basic_consume(queue='queue', on_message_callback=callback)
速记卡(面试闪卡)
Q1:一句话讲清「MQ 核心概念:生产-消费模型 + 主题/队列 + ACK 机制」到底是什么?
A:消息队列是异步通信中间件,用生产者-消费者模型解耦系统,核心含主题/队列与 ACK 确认。
Q2:生产-消费模型怎么理解? —— 怎么理解?
A:像快递柜:生产者塞件、消费者取件,两边互不等待、解耦提速。这是生产者-消费者模型(producer-consumer model)。中间那个存转发的叫 Broker。
Q3:主题 vs 队列怎么理解? —— 怎么理解?
A:队列像私信一人收(点对点,一个消息只一个消费者);主题像群发广播多人收(发布-订阅,pub/sub)。这是点对点 vs 发布订阅(Queue vs Topic)。
Q4:ACK 机制怎么理解? —— 怎么理解?
A:消费者吃完饭举手说”好了”(ACK)才撤盘,否则消息一直留着、防丢。这是确认机制(ACK, acknowledgement)。分自动 ACK 和手动 ACK 两种。
Q5:坑点与实战怎么理解? —— 怎么理解?
A:消息会丢(要持久化+ACK)、会重复(要做幂等)、会积压(加消费者)。Kafka/RabbitMQ 都绕不开这三坑。这是消息可靠性三件套(loss / duplicate / backlog)。
Q6:核心速记主线有哪些?
-
模型:Producer 发、Broker 存转发、Consumer 收
-
模式:Queue 点对点 vs Topic 发布订阅
-
可靠:ACK 确认 + 持久化防丢失,幂等防重复
-
实战:Kafka/RabbitMQ,警惕积压与重复消费
口诀
A:消息队列解耦合,
生产消费两头忙;
队列私聊主题广播,
ACK举手才收场。
相关链接
-
📋 目录:00-消息队列实战
-
📚 学习清单:技术学习路线图 > 消息队列实战