09. 消息队列深度解析

RocketMQ、Kafka、Pulsar 原理对比:存储模型、事务消息、顺序消费、延迟消息与性能调优

消息队列是分布式系统异步解耦的核心基础设施。RocketMQ 适合金融级可靠消息,Kafka 适合高吞吐日志流处理,Pulsar 则是云原生时代的统一消息平台。

1. 存储模型对比

特性KafkaRocketMQPulsar
存储日志分段文件CommitLog + ConsumeQueueBookKeeper (分布式日志)
消费模型PullPush + PullPush + Pull
顺序保证Partition 内有序Queue 内有序Partition 内有序
延迟消息TimeWheel (有限)18 级延迟原生支持
事务消息事务生产者✅ 原生支持✅ 原生支持
多租户
Geo 复制MirrorMaker不完全✅ 内置
吞吐百万级/秒十万级/秒百万级/秒

2. RocketMQ 核心架构

Producer → NameServer (路由发现)
              ↓
         Broker Master/Slave
              ↓
Consumer ← NameServer

2.1 存储结构

CommitLog        # 所有消息顺序写入(1G 一个文件)
├── ConsumeQueue  # 消费队列索引(定长 20byte/条目)
│   ├── TopicA/Queue0
│   ├── TopicA/Queue1
│   └── TopicB/Queue0
└── IndexFile     # 按 Key/Time 索引

2.2 事务消息

// 发送半消息
TransactionMQProducer producer = new TransactionMQProducer("tx_group");
producer.setTransactionListener(new TransactionListener() {
    @Override
    public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
        try {
            // 执行本地事务
            orderService.createOrder((Order) arg);
            return LocalTransactionState.COMMIT_MESSAGE;
        } catch (Exception e) {
            return LocalTransactionState.ROLLBACK_MESSAGE;
        }
    }
    
    @Override
    public LocalTransactionState checkLocalTransaction(MessageExt msg) {
        // 本地事务回查
        Order order = orderService.getById(msg.getKeys());
        return order != null ? COMMIT_MESSAGE : ROLLBACK_MESSAGE;
    }
});

2.3 顺序消费

// 全局顺序:所有消息发到同一个 Queue
SendResult result = producer.send(msg, new MessageQueueSelector() {
    @Override
    public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
        return mqs.get(0);  // 固定发送到 Queue 0
    }
}, null);

// 局部顺序:按 Sharding Key 选择 Queue(如订单 ID)
SendResult result = producer.send(msg, (mqs, m, arg) -> {
    int index = ((String) arg).hashCode() % mqs.size();
    return mqs.get(Math.abs(index));
}, orderId);

// 消费者:单线程消费保证顺序
consumer.registerMessageListener(new MessageListenerOrderly() {
    @Override
    public ConsumeOrderlyStatus consumeMessage(List<MessageExt> msgs, ConsumeOrderlyContext ctx) {
        for (MessageExt msg : msgs) {
            process(msg);
        }
        return ConsumeOrderlyStatus.SUCCESS;
    }
});

3. Kafka 核心原理

3.1 ISR 机制

Partition 0:
  Leader: Broker 1 (offset: 0-1000)
  ISR: [Broker 1, Broker 2]  # 同步副本集合
  OSR: [Broker 3]            # 滞后副本

副本滞后判定: replica.lag.time.max.ms=10000

3.2 生产者配置

Properties props = new Properties();
props.put("bootstrap.servers", "kafka1:9092,kafka2:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");

// 幂等生产者(自动去重)
props.put("enable.idempotence", "true");
props.put("acks", "all");
props.put("retries", Integer.MAX_VALUE);

// 事务生产者
props.put("transactional.id", "prod-1");

KafkaProducer<String, String> producer = new KafkaProducer<>(props);

// 事务发送
producer.initTransactions();
try {
    producer.beginTransaction();
    producer.send(new ProducerRecord<>("topic", "key", "value"));
    producer.sendOffsetsToTransaction(consumer.position(consumer.assignment()), consumer.groupMetadata());
    producer.commitTransaction();
} catch (Exception e) {
    producer.abortTransaction();
}

3.3 消费者分区分配

策略算法特点
Range按 Topic 范围分配默认,可能不均衡
RoundRobin轮询分配较均衡
Sticky最小化重新分配再均衡影响小
Cooperative增量再均衡Kafka 2.4+ 推荐

4. Pulsar:云原生消息队列

Pulsar 的核心创新是 存储与计算分离

  • Broker:无状态,处理生产和消费
  • BookKeeper:分布式日志存储
  • ZooKeeper/etcd:元数据管理
Producer → Broker → Bookie (Ledger)
              ↓
Consumer ← Broker

4.1 多租户与命名空间

# 创建租户
pulsar-admin tenants create my-org

# 创建命名空间
pulsar-admin namespaces create my-org/finance

# 设置策略
pulsar-admin namespaces set-retention my-org/finance \
  --size 10G --time 3d

4.2 统一消息模型

Pulsar 统一了队列(Queue)和流(Stream):

  • Exclusive/Failover:队列模式,单消费者
  • Shared:队列模式,多消费者竞争
  • Key_Shared:按 Key 路由到固定消费者(顺序保证 + 并行)

5. 选型建议

场景推荐理由
金融交易RocketMQ强一致、事务消息完善
日志采集/大数据Kafka吞吐最高、生态成熟
云原生/多租户Pulsar存储计算分离、弹性扩缩
实时计算Kafka/Pulsar流处理集成好
延迟消息RocketMQ18 级延迟原生支持

6. 消息队列最佳实践

  1. 生产端:异步发送、批量发送、失败重试、限流保护
  2. 消费端:幂等设计、消费确认、死信队列、消费速率监控
  3. 运维端:监控堆积延迟、磁盘水位、消费组 Rebalance

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「distributed-systems」更多文章

  1. 分布式高可用架构模式:多活、容灾、降级与 K8s 编排高可用
  2. 分布式链路追踪实战:OpenTelemetry、Jaeger 与 W3C Trace Context
  3. 分布式缓存深度策略:Redis Cluster、一致性哈希与多级缓存架构