消息中间件选型是后端开发中高频出现的架构决策,选错会导致吞吐不足、消息丢失或运维成本失控。主流消息中间件按定位分为三类:Kafka面向高吞吐日志流,RocketMQ面向业务消息与事务场景,RabbitMQ面向低延迟轻量任务。选型应围绕消息吞吐、可靠性与延迟三个核心指标展开,而不是盲目追随社区热度。
三种消息中间件的定位与适用场景
- Kafka:分区模型配合顺序写盘,单机吞吐可达百万级消息每秒,适合日志收集、用户行为追踪与大数据管道,但延迟相对高、消息丢失窗口大,不适合强一致业务。
- RocketMQ:由阿里巴巴开源,支持事务消息、延迟消息与消息重试,在订单、支付等业务场景中可靠性更好,吞吐低于Kafka但延迟更稳。
- RabbitMQ:基于Erlang的AMQP协议实现,功能完整、社区成熟,适合中低吞吐、要求低延迟与灵活路由的中小项目。
同业务中多套消息中间件并存也是常态:日志走Kafka,订单走RocketMQ,内部通知走RabbitMQ,各取所长。
高吞吐场景:Kafka的生产与消费配置
Kafka的吞吐瓶颈常出现在分区数与消费者组配置上。写入侧关注acks与批量参数:
# 生产端配置(Java)
Properties props = new Properties();
props.put("bootstrap.servers", "kafka1:9092,kafka2:9092");
props.put("acks", "all"); // 副本全部确认,保证不丢
props.put("linger.ms", 20); // 攒批发送,提升吞吐
props.put("batch.size", 65536);
props.put("compression.type", "lz4");
消费侧分区数决定了并行度,分区数应不小于消费者线程数,单个分区内消息有序,跨分区消费顺序由业务自行处理。
业务消息场景:RocketMQ的事务消息与延迟消息
订单创建与积分变动需要原子性,RocketMQ事务消息通过两阶段提交保证本地事务与消息发送的一致性,成为分布现场景常用方案:
// RocketMQ事务消息发送端
TransactionMQProducer producer = new TransactionMQProducer("tx-producer");
producer.setExecutorService(executor);
producer.start();
Message msg = new Message("ORDER_TOPIC", "create", orderJson.getBytes());
SendResult result = producer.sendMessageInTransaction(msg, orderId);
事务消息流程:先发送半消息,本地事务成功后commit,失败则rollback,配合回查接口处理网络异常,保证”先落库、再通知”的最终一致。
轻量场景:RabbitMQ的确认与死信配置
RabbitMQ适合延迟敏感的小任务,手动ACK配合死信队列是防丢的标准配置:
channel.basicConsume("task.queue", false, (consumerTag, delivery) -> {
try {
process(delivery.getBody());
channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
} catch (Exception e) {
// 超过重试次数进入死信队列
channel.basicReject(delivery.getEnvelope().getDeliveryTag(), false);
}
}, consumerTag -> {});
消费者异常时进入死信队列,由专门的补偿任务扫描重投,避免消息无限积压。
消息中间件选型的决策清单
- 先定指标:日吞吐、单条延迟、消息丢失容忍度,用数字定义需求。
- 再看运维:是否已有成熟的监控与集群管理能力,Kafka依赖ZooKeeper或KRaft,RocketMQ依赖NameServer。
- 最后看生态:连接器、客户端语言覆盖与团队熟悉程度。
选型完成后用压测脚本验证峰值下的延迟曲线,比看官方压测更贴近真实业务。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/xiao-xi-zhong-jian-jian-xuan-xing-shi-zhan-kafka-rocketmq/