消息中间件选型不是选”哪个更好”,是选”哪个更适合场景”
RabbitMQ和Kafka经常被拿来做二选一比较,但它们的架构模型完全不同。RabbitMQ基于Exchange-Queue模型,适合复杂路由和点对点消息投递;Kafka基于Topic-Partition模型,适合高吞吐日志流和事件溯源。选错中间件的代价是后期架构大改,远比初期多花两天调研严重得多。
架构模型差异与场景匹配
RabbitMQ:消息被路由到Queue,消费者从Queue取消息,消费后ACK删除。消息生命周期短,消费即消失。适合订单状态流转、任务分发、RPC调用。
Kafka:消息追加写入Partition日志文件,消费者按Offset消费。消息持久保留,可重复消费。适合日志收集、事件溯源、数据管道。
场景决策要点:
- 吞吐量需求 <10万/秒:RabbitMQ适合,Kafka可以但大材小用
- 吞吐量需求 >50万/秒:Kafka适合,RabbitMQ不适合
- 消息路由逻辑复杂(topic/fanout/headers):RabbitMQ适合
- 消息需要重复消费/回溯:Kafka适合
- 强依赖消息顺序性:单Queue/单Partition可保证
- 消费者数量经常变化:RabbitMQ适合,Kafka需要Rebalance
Spring Boot集成RabbitMQ生产级配置
@Configuration
public class RabbitMQConfig {
@Bean
public CustomExchange delayExchange() {
Map<String, Object> args = new HashMap<>();
args.put("x-delayed-type", "direct");
return new CustomExchange("order.delay.exchange",
"x-delayed-message", true, false, args);
}
@Bean
public Queue orderTimeoutQueue() {
return QueueBuilder.durable("order.timeout.queue")
.withArgument("x-dead-letter-exchange", "order.dlx.exchange")
.withArgument("x-dead-letter-routing-key", "order.dead")
.build();
}
@Bean
public SimpleRabbitListenerContainerFactory containerFactory(
ConnectionFactory connectionFactory) {
SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory();
factory.setConnectionFactory(connectionFactory);
factory.setAcknowledgeMode(AcknowledgeMode.MANUAL);
factory.setPrefetchCount(20);
factory.setConcurrentConsumers(3);
factory.setMaxConcurrentConsumers(10);
return factory;
}
}
// 生产者发送延迟消息
@Service
public class OrderMessageProducer {
@Autowired
private RabbitTemplate rabbitTemplate;
public void sendOrderTimeoutCheck(OrderEvent event, int delayMs) {
rabbitTemplate.convertAndSend("order.delay.exchange",
"order.timeout", event,
msg -> {
msg.getMessageProperties().setDelay(delayMs);
msg.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);
return msg;
});
}
}
// 消费者手动ACK
@RabbitListener(queues = "order.timeout.queue")
public void handleOrderTimeout(Message message, Channel channel) throws IOException {
long tag = message.getMessageProperties().getDeliveryTag();
try {
OrderEvent event = JSON.parseObject(new String(message.getBody()), OrderEvent.class);
orderService.cancelOrder(event.getOrderId());
channel.basicAck(tag, false);
} catch (Exception e) {
channel.basicNack(tag, false, false);
}
}
Spring Boot集成Kafka生产级配置
// application.yml配置
spring:
kafka:
bootstrap-servers: kafka1:9092,kafka2:9092,kafka3:9092
producer:
acks: all
retries: 3
batch-size: 16384
linger-ms: 5
consumer:
group-id: order-service-group
auto-offset-reset: earliest
max-poll-records: 500
enable-auto-commit: false
// 生产者
@Service
public class EventProducer {
@Autowired
private KafkaTemplate<String, Object> kafkaTemplate;
public void sendOrderCreatedEvent(OrderCreatedEvent event) {
kafkaTemplate.send("order-events", event.getOrderId(), event)
.addCallback(
result -> log.info("消息发送成功: partition={}, offset={}",
result.getRecordMetadata().partition(),
result.getRecordMetadata().offset()),
ex -> log.error("消息发送失败: {}", ex.getMessage())
);
}
}
// 消费者手动提交Offset
@KafkaListener(topics = "order-events", groupId = "inventory-service-group")
public void handleOrderEvent(ConsumerRecord<String, OrderCreatedEvent> record,
Acknowledgment ack) {
try {
OrderCreatedEvent event = record.value();
inventoryService.reserveStock(event.getItems());
ack.acknowledge();
} catch (Exception e) {
log.error("库存预留失败, partition={}, offset={}",
record.partition(), record.offset(), e.getMessage());
}
}
高并发场景下的关键调优参数
RabbitMQ调优:
- prefetch_count设20-50,太大会导致消费者内存压力,太小影响吞吐
- 消息持久化(delivery-mode=2)降低2-3倍吞吐,非关键消息可用非持久化
- 镜像队列(ha-policy)会额外增加网络开销,3节点集群吞吐下降约30%
Kafka调优:
- Partition数 = 目标吞吐 / 单Partition吞吐,过度分区增加Rebalance时间和Controller压力
- linger.ms设5-10ms,batch.size设16KB-64KB,吞吐和延迟的平衡点
- 消费者max.poll.records和max.poll.interval.ms配合调优,避免Rebalance
选型前画清楚消息流图,标注每个节点的吞吐、持久化、路由需求,匹配结果自然就出来了。不要在选型阶段纠结”谁更好”,要在上线后用监控数据验证”选的对不对”。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/xiao-xi-zhong-jian-jian-gao-bing-fa-xuan-xing-shi-zhan/