Skip to content
第 201 / 250 章中间件⏱ 12 分钟阅读

第 201 章:MQ 选型对比与应用场景

学习目标

  • 掌握 Kafka / RabbitMQ / RocketMQ / Redis Stream 对比
  • 根据场景选型
  • 理解消息中间件核心指标
  • 避免常见踩坑

一、四大消息中间件

特性KafkaRabbitMQRocketMQRedis Stream
语言Scala/JavaErlangJavaC
吞吐量百万/s万/s十万/s万/s
延迟ms 级us 级ms 级ms 级
持久化磁盘内存+磁盘磁盘内存+磁盘
消息回溯
事务
死信
延迟消息❌(需插件)
优先级队列
集群
运维难度

二、消息中间件核心指标

2.1 四大维度

2.2 可靠性级别

级别保证适用
最多一次可能丢日志
至少一次可能重复订单
恰好一次不丢不重金融(难)

三、典型场景选型

3.1 选型决策树

3.2 业务消息(订单、交易)

推荐:RocketMQ 或 RabbitMQ

原因:

  • ✅ 事务消息(订单创建 + 库存扣减)
  • ✅ 延迟消息(订单超时)
  • ✅ 死信队列(失败重试)
  • ✅ 万级吞吐够用

3.3 日志 / 事件流(用户行为、IoT)

推荐:Kafka

原因:

  • ✅ 百万级吞吐
  • ✅ 消息回溯(可重放)
  • ✅ 持久化(磁盘顺序写)
  • ✅ 流处理生态(Kafka Streams)

3.4 微服务解耦(简单)

推荐:Redis Stream 或 RabbitMQ

原因:

  • ✅ 部署简单
  • ✅ 路由灵活(RabbitMQ)
  • ✅ 成本低

3.5 金融交易(强事务)

推荐:RocketMQ

原因:

  • ✅ 原生事务消息
  • ✅ 阿里双 11 验证
  • ✅ 高可用架构

四、消息可靠性设计

4.1 生产者不丢

java
// Kafka
props.put("acks", "all");
props.put("enable.idempotence", true);

// RabbitMQ
spring.rabbitmq.publisher-confirm-type=correlated

// RocketMQ
producer.setRetryTimesWhenSendFailed(3);
producer.send(msg, new SendCallback() { ... });

4.2 Broker 不丢

properties
# Kafka 多副本
replication.factor=3
min.insync.replicas=2

# RabbitMQ 镜像队列
ha-mode: all
ha-sync-mode: automatic

4.3 消费者不丢

java
// 业务成功 + 提交 offset
inventoryService.deduct(event);
ack.acknowledge();

4.4 事务消息

java
// RocketMQ 事务消息
TransactionMQProducer producer = new TransactionMQProducer("group");
producer.setTransactionListener(new TransactionListener() {
    @Override
    public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
        // 本地事务
        try {
            orderService.create(msg);
            return LocalTransactionState.COMMIT_MESSAGE;
        } catch (Exception e) {
            return LocalTransactionState.ROLLBACK_MESSAGE;
        }
    }

    @Override
    public LocalTransactionState checkLocalTransaction(MessageExt msg) {
        // 回查(防止本地事务结果未确定)
        return orderService.exists(msg.getKeys())
            ? LocalTransactionState.COMMIT_MESSAGE
            : LocalTransactionState.ROLLBACK_MESSAGE;
    }
});

五、消息幂等设计

5.1 唯一 ID 方案

java
// 生产
Message msg = new Message(topic, body);
msg.setKeys(orderId);   // 唯一键
producer.send(msg);

// 消费
String orderId = msg.getKeys();
if (redisTemplate.opsForValue().setIfAbsent("mq:order:" + orderId, "1", 24, TimeUnit.HOURS)) {
    // 首次消费
    inventoryService.deduct(orderId);
} else {
    log.info("重复消息: {}", orderId);
}

5.2 数据库唯一约束

sql
CREATE TABLE order_event (
    id BIGINT PRIMARY KEY,
    order_id BIGINT,
    UNIQUE KEY uk_order (order_id)
);

六、消息积压治理

6.1 应急

bash
# 1. 临时扩容消费者
kafka-topics --alter --topic orders --partitions 30

# 2. 直接读底层 log(应急)
kafka-console-consumer --bootstrap-server localhost:9092 \
  --topic orders --partition 0 --offset earliest

6.2 长期方案

  • 异步消费
  • 批量消费
  • 业务降级
  • 限流保护

七、消息中间件监控

指标监控方式
发送成功率应用埋点
消费 LagKafka Lag / RabbitMQ 队列深度
消费失败率应用埋点
延迟end-to-end 时间戳

八、本章小结

场景推荐
业务消息RocketMQ / RabbitMQ
日志事件Kafka
微服务解耦RabbitMQ / Redis Stream
简单队列Redis Stream
金融交易RocketMQ

动手练习

  1. 对比测试 Kafka / RabbitMQ / Redis Stream 性能
  2. 用事务消息实现下单 + 库存
  3. 实现消费幂等(Redis 去重)
  4. 监控 Lag 并配置告警

推荐阅读


下一章:第 202 章:Elasticsearch 入门

本站基于 VitePress 构建 · 由 Codebook 团队维护