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举手才收场。

相关链接