← 返回 MQ 列表

核心架构与消息模型篇

核心架构与消息模型篇

一、消息模型

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 在这些维度的实现各有侧重,实际选型时需要结合业务场景的吞吐量、可靠性、实时性要求综合判断。