核心架构与消息模型篇
一、消息模型
1. 点对点模型(P2P, Point-to-Point)
核心特征:一条消息只能被一个消费者消费,消费后消息从队列中删除。
┌──────────┐
Producer ────────▶│ Queue │───▶ Consumer A(消费成功)
└──────────┘
│
└────▶ Consumer B(收不到,已被 A 消费)
适用场景:
- 任务分发(一个任务只分配给一个 worker)
- 订单处理(一条订单只由一个服务处理)
代码示例 — 基于 ActiveMQ 的 P2P:
java
// 生产者
public class OrderProducer {
public static void main(String[] args) {
ConnectionFactory factory = new ActiveMQConnectionFactory("tcp://localhost:61616");
Connection conn = factory.createConnection();
conn.start();
Session session = conn.createSession(false, Session.AUTO_ACKNOWLEDGE);
Queue queue = session.createQueue("order.queue");
MessageProducer producer = session.createProducer(queue);
for (int i = 0; i < 10; i++) {
TextMessage msg = session.createTextMessage("订单 #" + i);
producer.send(msg);
}
producer.close();
session.close();
conn.close();
}
}
// 消费者 A
public class OrderConsumerA {
public static void main(String[] args) throws Exception {
ConnectionFactory factory = new ActiveMQConnectionFactory("tcp://localhost:61616");
Connection conn = factory.createConnection();
conn.start();
Session session = conn.createSession(false, Session.AUTO_ACKNOWLEDGE);
Queue queue = session.createQueue("order.queue");
MessageConsumer consumer = session.createConsumer(queue);
consumer.setMessageListener(msg -> {
TextMessage text = (TextMessage) msg;
System.out.println("消费者 A 收到: " + text.getText());
// 假设确认处理
});
}
}
// 消费者 B(与 A 同一队列)
// 效果:A 和 B 轮询收到消息,每条消息只被其中一个消费
2. 发布订阅模型(Pub/Sub)
核心特征:一条消息广播到所有订阅者,每个订阅者都能收到全量消息。
┌──────────────┐
│ │───▶ Subscriber A(收到全部消息)
Publisher ───────▶│ Topic │───▶ Subscriber B(收到全部消息)
│ │───▶ Subscriber C(收到全部消息)
└──────────────┘
适用场景:
- 事件广播(订单状态变更通知所有关心系统)
- 数据同步(配置变更通知所有节点)
- 日志收集(所有服务发送日志到统一 topic)
代码示例 — RocketMQ Pub/Sub:
java
// 生产者 — 向 Topic 发送消息
public class OrderEventPublisher {
public static void main(String[] args) throws Exception {
DefaultMQProducer producer = new DefaultMQProducer("order_event_group");
producer.setNamesrvAddr("localhost:9876");
producer.start();
for (int i = 0; i < 100; i++) {
OrderEvent event = new OrderEvent(i, "PAID", System.currentTimeMillis());
Message msg = new Message("order_topic",
"order_event",
event.toJson().getBytes(StandardCharsets.UTF_8));
producer.send(msg);
}
producer.shutdown();
}
}
// 消费者 A — 库存服务(订阅 order_topic)
public class InventoryConsumer {
public static void main(String[] args) throws Exception {
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("inventory_group");
consumer.setNamesrvAddr("localhost:9876");
consumer.subscribe("order_topic", "*"); // 订阅所有 tag
consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
for (MessageExt msg : msgs) {
OrderEvent event = OrderEvent.fromJson(new String(msg.getBody()));
System.out.println("库存服务处理: " + event);
// 扣减库存逻辑
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
});
consumer.start();
}
}
// 消费者 B — 积分服务(也订阅 order_topic)
// 效果:库存和积分服务各自独立消费同一条消息
P2P vs Pub/Sub 对比
| 对比维度 | 点对点(P2P) | 发布订阅(Pub/Sub) |
|---|---|---|
| 消费关系 | 1:1 | 1:N |
| 消息去向 | 一个消费者独占 | 所有订阅者各一份 |
| 消息消费后 | 从队列删除 | 每个订阅者独立消费 |
| 典型实现 | Queue(JMS) | Topic(JMS) |
| 适合场景 | 任务分发、负载均衡 | 事件广播、数据同步 |
二、集群架构
1. 单机模式
架构:一个 Broker 节点,所有消息存储在单机。
Producer ──▶ Broker ──▶ Consumer
优缺点:
- ✅ 部署简单,适合开发测试
- ❌ 单点故障,无高可用
- ❌ 性能受单机限制
2. 主从/副本模式(Master-Slave)
架构:一个 Master 节点负责读写,一个或多个 Slave 节点同步数据,故障时切换。
┌──────────┐
│ Master │ ←── 生产者读写
└────┬─────┘
│ 同步复制
┌────▼─────┐
│ Slave │ ←── 故障时接管
└──────────┘
RocketMQ 主从配置示例:
properties
# broker.properties — Master 配置
brokerClusterName=DefaultCluster
brokerName=broker-a
brokerId=0 # 0 表示 Master
namesrvAddr=localhost:9876
brokerRole=SYNC_MASTER # 同步复制模式
flushDiskType=ASYNC_FLUSH
# broker.properties — Slave 配置
brokerClusterName=DefaultCluster
brokerName=broker-a
brokerId=1 # > 0 表示 Slave
namesrvAddr=localhost:9876
brokerRole=SLAVE
flushDiskType=ASYNC_FLUSH
三种复制模式:
| 模式 | 说明 | 优缺点 |
|---|---|---|
| 同步复制(SYNC_MASTER) | Master 写入后等待 Slave 确认才返回 | 数据最安全,延迟略高 |
| 异步复制(ASYNC_MASTER) | Master 写入即返回,Slave 异步同步 | 性能好,极端情况可能丢少量数据 |
| 半同步复制 | 至少一个 Slave 确认即可 | 折中方案 |
3. 集群分片模式(分布式集群)
架构:数据分片存储到多个 Broker,每个 Broker 负责一部分数据,整体提升吞吐量和容量。
┌──────────────┐
│ NameServer │ ←── 路由发现
└──────┬───────┘
│
┌─────────────────┼──────────────────┐
│ │ │
┌─────▼─────┐ ┌──────▼──────┐ ┌──────▼──────┐
│ Broker A │ │ Broker B │ │ Broker C │
│ 分区 0, 3 │ │ 分区 1, 4 │ │ 分区 2, 5 │
│ Master │ │ Master │ │ Master │
└─────┬─────┘ └──────┬──────┘ └──────┬──────┘
│ │ │
┌─────▼─────┐ ┌──────▼──────┐ ┌──────▼──────┐
│ Broker A' │ │ Broker B' │ │ Broker C' │
│ Slave │ │ Slave │ │ Slave │
└───────────┘ └─────────────┘ └─────────────┘
Kafka 分区架构:
java
// 生产者 — 指定 key 路由到特定分区
public class PartitionProducer {
public static void main(String[] args) {
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092,localhost:9093,localhost:9094");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
for (int i = 0; i < 100; i++) {
String orderId = "ORDER_" + i;
// 同一个 orderId 的消息进入同一分区,保证顺序
ProducerRecord<String, String> record = new ProducerRecord<>(
"order_topic", // topic
orderId, // key → 决定分区
"订单数据: " + i // value
);
producer.send(record);
}
producer.close();
}
}
三、通信模式
1. 推模式(Push)
Broker 主动将消息推送给消费者。
┌──────────┐
Producer ────▶│ Broker │───▶ Consumer(被动接收)
└──────────┘
优点:
- 实时性高,消息到达后立即推送
- 消费者实现简单(只需注册回调)
缺点:
- Broker 压力大(需要维护推送状态)
- 消费者处理慢时容易被打满
RocketMQ Push 示例(本质是长轮询伪装成 Push):
java
// RocketMQ Push Consumer — 注册监听器即可
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("consumer_group");
consumer.setNamesrvAddr("localhost:9876");
consumer.subscribe("topic", "*");
consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
// Broker 推送消息到这里,实际底层是长轮询
for (MessageExt msg : msgs) {
process(msg);
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
});
consumer.start();
2. 拉模式(Pull)
消费者主动轮询 Broker 拉取消息。
┌──────────┐
Producer ────▶│ Broker │◀─── Consumer(主动轮询)
└──────────┘
优点:
- 消费者控制拉取节奏,不会被打满
- Broker 压力小
缺点:
- 实时性不如 Push(需要轮询)
- 消费者实现复杂(需要管理拉取偏移量)
Kafka Pull 示例:
java
// Kafka 原生 Pull 模式
public class KafkaPullConsumer {
public static void main(String[] args) {
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("group.id", "consumer_group");
props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
props.put("enable.auto.commit", "false"); // 手动提交偏移量
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Arrays.asList("order_topic"));
while (true) {
// 主动拉取,超时 1000ms
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
for (ConsumerRecord<String, String> record : records) {
System.out.printf("分区=%d, 偏移量=%d, 值=%s%n",
record.partition(), record.offset(), record.value());
}
// 手动提交偏移量
consumer.commitSync();
}
}
}
推模式 vs 拉模式对比
| 对比维度 | 推模式(Push) | 拉模式(Pull) |
|---|---|---|
| 实时性 | 高 | 取决于轮询间隔 |
| 消费者控制力 | 弱(被动接收) | 强(自己控制速率) |
| Broker 压力 | 大 | 小 |
| 实现复杂度 | 简单 | 中等 |
| 典型代表 | RocketMQ Push | Kafka / RabbitMQ |
| 适用场景 | 实时性要求高 | 消费能力强/需背压控制 |
四、RocketMQ 完整集群部署实战
部署架构
┌───────────────────┐
│ NameServer │
│ (路由注册中心) │
└───────────────────┘
│
┌───────────────────┼────────────────────┐
│ │ │
┌─────▼───────┐ ┌──────▼───────┐ ┌───────▼──────┐
│ Broker-a │ │ Broker-b │ │ Broker-c │
│ Master:0 │ │ Master:0 │ │ Master:0 │
│ Slave:1 │ │ Slave:1 │ │ Slave:1 │
└─────────────┘ └──────────────┘ └──────────────┘
启动命令
bash
# 1. 启动 NameServer
nohup sh bin/mqnamesrv > namesrv.log 2>&1 &
# 2. 启动 Broker(Master)
nohup sh bin/mqbroker -c conf/2m-2s-sync/broker-a.properties > broker-a.log 2>&1 &
# 3. 启动 Broker(Slave)
nohup sh bin/mqbroker -c conf/2m-2s-sync/broker-a-s.properties > broker-a-s.log 2>&1 &
生产环境推荐配置
bash
# JVM 参数调优
JAVA_OPT="${JAVA_OPT} -server -Xms8g -Xmx8g -Xmn4g"
JAVA_OPT="${JAVA_OPT} -XX:+UseG1GC -XX:G1HeapRegionSize=16m"
JAVA_OPT="${JAVA_OPT} -XX:+DisableExplicitGC"
JAVA_OPT="${JAVA_OPT} -Drocketmq.client.logUseSlf4j=true"
# Broker 核心配置
brokerClusterName=OnlineCluster
brokerName=broker-a
brokerId=0
namesrvAddr=192.168.1.1:9876;192.168.1.2:9876
defaultTopicQueueNums=8
autoCreateTopicEnable=false
autoCreateSubscriptionGroup=false
listenPort=10911
brokerRole=SYNC_MASTER
flushDiskType=ASYNC_FLUSH
# 文件保留 72 小时
fileReservedTime=72
# 单个文件大小 1GB
mapedFileSizeCommitLog=1073741824
总结:理解消息模型(P2P vs Pub/Sub)、集群架构(主从 vs 分片)、通信模式(Push vs Pull)是掌握消息队列的基础。不同 MQ 在这些维度的实现各有侧重,实际选型时需要结合业务场景的吞吐量、可靠性、实时性要求综合判断。