消息队列是分布式系统异步解耦的核心基础设施。RocketMQ 适合金融级可靠消息,Kafka 适合高吞吐日志流处理,Pulsar 则是云原生时代的统一消息平台。
1. 存储模型对比
| 特性 | Kafka | RocketMQ | Pulsar |
|---|---|---|---|
| 存储 | 日志分段文件 | CommitLog + ConsumeQueue | BookKeeper (分布式日志) |
| 消费模型 | Pull | Push + Pull | Push + 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 | 流处理集成好 |
| 延迟消息 | RocketMQ | 18 级延迟原生支持 |
6. 消息队列最佳实践
- 生产端:异步发送、批量发送、失败重试、限流保护
- 消费端:幂等设计、消费确认、死信队列、消费速率监控
- 运维端:监控堆积延迟、磁盘水位、消费组 Rebalance
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。