消息队列是微服务解耦与异步化的核心基础设施。Spring 提供了统一的抽象层,但每种消息中间件在消息语义、事务保证与消费模型上各有侧重。本文基于实战场景对比集成方案。
1. 三剑客选型对比
| 维度 | RabbitMQ | Kafka | RocketMQ |
|---|---|---|---|
| 协议 | AMQP | 自定义协议 | 自定义协议 |
| 吞吐量 | 万级 | 百万级 | 十万级 |
| 延迟 | 微秒级 | 毫秒级 | 毫秒级 |
| 消费模型 | Push(推) | Pull(拉) | Push + Pull |
| 消息顺序 | 单队列有序 | Partition 内有序 | Queue 内有序 |
| 事务消息 | 本地事务 + Confirm | 事务生产者(幂等) | ✅ 原生支持 |
| 延迟消息 | 死信队列 + TTL | 不太适合 | ✅ 原生支持 18 级 |
| 消息追踪 | 插件 | 较复杂 | 较完善 |
| 管理界面 | 完善 | 需 Kowl/AKHQ | Console 完善 |
| 最佳场景 | 路由复杂、可靠投递 | 日志/流处理/大数据 | 金融级事务、定时任务 |
2. Spring AMQP + RabbitMQ
2.1 基础集成
spring:
rabbitmq:
host: localhost
port: 5672
username: guest
password: guest
listener:
simple:
concurrency: 5
max-concurrency: 20
acknowledge-mode: manual
prefetch: 10
retry:
enabled: true
max-attempts: 3
initial-interval: 1000ms
@Configuration
public class RabbitConfig {
@Bean
public Queue orderQueue() {
return QueueBuilder.durable("order.queue")
.withArgument("x-dead-letter-exchange", "order.dlx")
.withArgument("x-dead-letter-routing-key", "order.failed")
.build();
}
@Bean
public DirectExchange orderExchange() {
return new DirectExchange("order.exchange");
}
@Bean
public Binding orderBinding(Queue orderQueue, DirectExchange orderExchange) {
return BindingBuilder.bind(orderQueue).to(orderExchange).with("order.created");
}
}
@Component
public class OrderListener {
@RabbitListener(queues = "order.queue")
public void handleOrder(@Payload OrderMessage message,
Channel channel,
@Header(AmqpHeaders.DELIVERY_TAG) long tag) throws IOException {
try {
processOrder(message);
channel.basicAck(tag, false);
} catch (Exception e) {
channel.basicNack(tag, false, !message.isRedelivered()); // 重试一次后入死信
}
}
}
2.2 延迟消息(死信 + TTL)
@Bean
public Queue delayQueue() {
return QueueBuilder.durable("order.delay.queue")
.withArgument("x-message-ttl", 60000) // 60 秒后过期
.withArgument("x-dead-letter-exchange", "order.exchange")
.withArgument("x-dead-letter-routing-key", "order.timeout")
.build();
}
// 业务:订单 30 分钟未支付自动取消
rabbitTemplate.convertAndSend("order.delay.exchange", "order.delay", order, msg -> {
msg.getMessageProperties().setExpiration("1800000"); // 30 分钟 TTL
return msg;
});
RocketMQ 原生延迟:18 个固定级别(1s/5s/10s/30s/1m/…/2h)。
3. Spring Kafka
3.1 生产者配置
spring:
kafka:
bootstrap-servers: kafka1:9092,kafka2:9092
producer:
acks: all # 0/1/all
retries: 3
batch-size: 16384
linger-ms: 5
key-serializer: org.apache.kafka.common.serialization.StringSerializer
value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
properties:
enable.idempotence: true # 幂等生产者(防止重复发送)
max.in.flight.requests.per.connection: 5
@Service
public class KafkaOrderProducer {
@Autowired private KafkaTemplate<String, OrderEvent> kafkaTemplate;
public void sendOrderCreated(OrderEvent event) {
ProducerRecord<String, OrderEvent> record = new ProducerRecord<>(
"order-events",
event.getOrderId(), // 相同 key → 同一 partition → 顺序保证
event
);
ListenableFuture<SendResult<String, OrderEvent>> future =
kafkaTemplate.send(record);
future.addCallback(
result -> log.info("Sent: {}", result.getRecordMetadata().offset()),
failure -> log.error("Failed: {}", failure.getMessage())
);
}
}
3.2 消费者配置
spring:
kafka:
consumer:
group-id: order-service
auto-offset-reset: earliest
enable-auto-commit: false
max-poll-records: 50
isolation-level: read_committed // 事务消费者
listener:
ack-mode: manual_immediate
concurrency: 3 // 每个 topic 的并发消费者数
@Component
public class KafkaOrderConsumer {
@KafkaListener(topics = "order-events", groupId = "order-service")
public void consume(@Payload OrderEvent event,
@Header(KafkaHeaders.RECEIVED_PARTITION) int partition,
@Header(KafkaHeaders.OFFSET) long offset,
Acknowledgment ack) {
try {
// 幂等校验:orderId + eventType 已处理过?
if (idempotencyService.isProcessed(event.getId())) {
ack.acknowledge();
return;
}
processEvent(event);
idempotencyService.markProcessed(event.getId());
ack.acknowledge();
} catch (Exception e) {
// 不 ack,自动重试(配置 `retry`)或入死信 topic
throw e;
}
}
}
3.3 Kafka 事务
@Service
public class OrderService {
@Autowired private KafkaTemplate<String, Object> kafkaTemplate;
@Transactional("kafkaTransactionManager")
public void createOrderWithEvent(Order order) {
// 1. 保存订单到 DB
orderRepository.save(order);
// 2. 发送事件(同一事务中)
kafkaTemplate.send("order-events", new OrderCreatedEvent(order));
// 若 DB 回滚,消息也不会发送
}
}
4. RocketMQ Spring Starter
4.1 事务消息
RocketMQ 是原生支持事务消息最完善的方案。
@Service
public class OrderRocketService {
@Autowired private RocketMQTemplate rocketMQTemplate;
public void createOrder(Order order) {
Message<Order> message = MessageBuilder.withPayload(order)
.setHeader("KEYS", order.getId())
.build();
TransactionSendResult result = rocketMQTemplate.sendMessageInTransaction(
"order-producer-group",
"order-topic",
message,
order // 本地事务参数
);
if (!result.getLocalTransactionState().equals(LocalTransactionState.COMMIT_MESSAGE)) {
throw new OrderException("订单创建失败");
}
}
}
// 本地事务监听器
@RocketMQTransactionListener
public class OrderTransactionListener implements RocketMQLocalTransactionListener {
@Override
public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
try {
Order order = (Order) arg;
orderService.save(order); // 执行本地事务
return RocketMQLocalTransactionState.COMMIT_MESSAGE;
} catch (Exception e) {
return RocketMQLocalTransactionState.ROLLBACK_MESSAGE;
}
}
@Override
public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
// 事务回查:网络等原因导致 RocketMQ 未收到确认
String orderId = msg.getKeys();
if (orderService.exists(orderId)) {
return RocketMQLocalTransactionState.COMMIT_MESSAGE;
}
return RocketMQLocalTransactionState.ROLLBACK_MESSAGE;
}
}
5. 消费幂等性设计
5.1 幂等实现方案
消息唯一标识:msgId(中间件生成)或业务键(bizId)
方案 1:数据库唯一索引
CREATE TABLE idempotency (
idempotency_key VARCHAR(64) PRIMARY KEY,
created_at TIMESTAMP DEFAULT NOW()
);
INSERT IGNORE INTO idempotency (idempotency_key) VALUES ('order:123:created');
方案 2:Redis SETNX(带 TTL)
SET idempotency:order:123:created "1" NX EX 86400
方案 3:业务状态机防重
订单状态: PENDING → PAID → SHIPPED
PAID 状态下重复收到 PAYMENT_SUCCESS 消息 → 直接忽略
6. 死信队列与监控
// Spring Kafka 死信 topic
@KafkaListener(topics = "order-events", groupId = "order-service")
public void consumeWithDLQ(OrderEvent event) {
// 重试 3 次后自动转发到 order-events.DLT
}
// DLT 监听(人工/自动化处理)
@KafkaListener(topics = "order-events.DLT", groupId = "order-dlt-service")
public void handleDeadLetter(OrderEvent event,
@Header(KafkaHeaders.DLT_EXCEPTION_MESSAGE) String error) {
alertService.send("DLQ Alert", event, error);
// 人工介入或自动补偿逻辑
}
延伸阅读
- 事件驱动架构设计 — 消息队列在架构层面的设计模式
- Seata 分布式事务 — AT/TCC/Saga 模式与消息最终一致性
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。