rabbitmq(3):幂等与延迟消息实战
业务幂等性
幂等是一个数学概念,用函数表达式来描述是这样的:f(x) = f(f(x)),例如求绝对值函数。
在程序开发中,则是指同一个业务,执行一次或多次对业务状态的影响是一致的。例如:
- 根据id删除数据
- 查询数据
- 新增数据
但数据的更新往往不是幂等的,如果重复执行可能造成不一样的后果。比如:
- 取消订单,恢复库存的业务。如果多次恢复就会出现库存重复增加的情况
- 退款业务。重复退款对商家而言会有经济损失。
所以,我们要尽可能避免业务被重复执行。
然而在实际业务场景中,由于意外经常会出现业务被重复执行的情况,例如:
- 页面卡顿时频繁刷新导致表单重复提交
- 服务间调用的重试
- MQ消息的重复投递
我们在用户支付成功后会发送MQ消息到交易服务,修改订单状态为已支付,就可能出现消息重复投递的情况。如果消费者不做判断,很有可能导致消息被消费多次,出现业务故障。
举例:
- 假如用户刚刚支付完成,并且投递消息到交易服务,交易服务更改订单为已支付状态。
- 由于某种原因,例如网络故障导致生产者没有得到确认,隔了一段时间后重新投递给交易服务。
- 但是,在新投递的消息被消费之前,用户选择了退款,将订单状态改为了已退款状态。
- 退款完成后,新投递的消息才被消费,那么订单状态会被再次改为已支付。业务异常。
因此,我们必须想办法保证消息处理的幂等性。这里给出两种方案:
- 唯一消息ID
- 业务状态判断
唯一id
- 每一条消息都生成一个唯一的id,与消息一起投递给消费者。
- 消费者接收到消息后处理自己的业务,业务处理成功后将消息ID保存到数据库
- 如果下次又收到相同消息,去数据库查询判断是否存在,存在则为重复消息放弃处理。
SpringAMQP的MessageConverter自带了MessageID的功能
@Bean
public MessageConverter messageConverter() {
Jackson2JsonMessageConverter converter = new Jackson2JsonMessageConverter();
converter.setCreateMessageIds(true);
return converter;
}
构建的jackson消息转化器,设置了消息的id,保证了唯一性
但同时造成业务侵入、且有数据库的操作影响业务性能
业务判断
业务判断就是基于业务本身的逻辑或状态来判断是否是重复的请求或消息
当存在合适的业务判断条件时可以使用,但部分很难找到
如:
支付订单时,可以判断订单的状态在进行操作
@Override
public void markOrderPaySuccess(Long orderId) {
// 1.查询订单
Order old = getById(orderId);
// 2.判断订单状态
if (old == null || old.getStatus() != 1) {
// 订单不存在或者订单状态不是1,放弃处理
return;
}
// 3.尝试更新订单
Order order = new Order();
order.setId(orderId);
order.setStatus(2);
order.setPayTime(LocalDateTime.now());
updateById(order);
}
此做法只加了一层判断,存在线程安全问题
@Override
public void markOrderPaySuccess(Long orderId) {
// UPDATE `order` SET status = ? , pay_time = ? WHERE id = ? AND status = 1
lambdaUpdate()
.set(Order::getStatus, 2)
.set(Order::getPayTime, LocalDateTime.now())
.eq(Order::getId, orderId)
.eq(Order::getStatus, 1)
.update();
}
将业务逻辑以及判断合并在一条sql语句中
乐观锁思想兜底
支付服务与交易服务之间的订单状态一致性
-
首先,支付服务会正在用户支付成功以后利用MQ消息通知交易服务,完成订单状态同步。
-
其次,为了保证MQ消息的可靠性,我们采用了生产者确认机制、消费者确认、消费者失败重试等策略,确保消息投递的可靠性
-
最后,我们还在交易服务设置了定时任务,定期查询订单支付状态。
这样即便MQ通知失败,还可以利用定时任务作为兜底方案,确保订单支付状态的最终一致性。

使用mq的延迟消息,实现乐观锁:更新时检查“是否被别人改动过”,如果改动了,就放弃这次更新。
延迟消息
在一段时间以后才执行的任务,我们称之为延迟任务,
要实现延迟任务,最简单的方案是利用MQ的延迟消息。
在RabbitMQ中实现延迟消息也有两种方案:
- 死信交换机+TTL
- 延迟消息插件
死信交换机
概述:
当一个队列中的消息满足下列情况之一时,可以成为死信(dead letter):
- 消费者使用
basic.reject或basic.nack声明消费失败,并且消息的requeue参数设置为false - 消息是一个过期消息,超时无人消费
- 要投递的队列消息满了,无法投递
如果一个队列中的消息已经成为死信,并且这个队列通过**dead-letter-exchange**属性指定了一个交换机,
那么队列中的死信就会投递到这个交换机中,而这个交换机就称为死信交换机(Dead Letter Exchange)。
而此时有队列与死信交换机绑定,则最终死信就会被投递到这个队列中。

Producer
↓ routingKey=red
normal.direct
↓ bindingKey=red
normal.queue
↓ 死信
RabbitMQ自动重新发布
↓ routingKey=error
dlx.direct
↓ bindingKey=error
dlx.queue
作用:
- 收集因处理失败而被拒绝的消息
- 收集因队列满了而被拒绝的消息
- 收集因TTL(有效期)到期的消息
主要过程:
生产者将设置了 TTL 的消息发送到普通交换机 normal.direct,由其路由到普通队列 normal.queue。
normal.queue 通过参数 x-dead-letter-exchange 指定了死信交换机 dlx.direct,并通过 x-dead-letter-routing-key 指定死信转发时使用的路由键。
当消息在 normal.queue 中超过 TTL 后,会变成死信,
RabbitMQ Broker 会自动将其转发到死信交换机 dlx.direct。
随后,死信交换机根据指定的 routingKey 将消息路由到 dlx.queue,最终由消费者消费,从而实现延迟消费的效果。
使用:
创建normal交换机及其队列,出现死信便会交给死信交换器:
@Bean
public DirectExchange normalExchange() {
return new DirectExchange("normal.direct", true, false);
}
@Bean
public Queue normalQueue() {
return QueueBuilder
.durable("normal.direct")
.deadLetterExchange("dlx.direct")
.deadLetterRoutingKey("error")
.build();
}
@Bean
public Binding normalExchangeBinding() {
return BindingBuilder.bind(normalQueue()).to(normalExchange()).with("normal");
}
死信交换器及其队列进行监听:
@RabbitListener(
bindings=@QueueBinding(
value=@Queue(name="dlx.queue",durable = "true"),
exchange=@Exchange(name="dlx.direct",type = ExchangeTypes.DIRECT),
key={"error"}
)
)
public void listenDlxQueue(String message){
log.info("消费者监听到消息:{}",message);
}
发送消息:
@Test
public void testSendDelayMessage() {
rabbitTemplate.convertAndSend
("normal.direct","normal","hello",message ->
{
message.getMessageProperties().setExpiration("10000");
return message;
});
}
通过定义message自身的属性Expiration,设置了ttl
这里使用了 new MessagePostProcessor()进行处理,与下面的方法类似,但此方法内只有一个函数,使用lambda更为方便
@Test
void testPublisherDelayMessage() {
// 1.创建消息
String message = "hello, delayed message";
// 2.发送消息,利用消息后置处理器添加消息头
rabbitTemplate.convertAndSend("delay.direct", "delay", message, new MessagePostProcessor() {
@Override
public Message postProcessMessage(Message message) throws AmqpException {
// 添加延迟消息属性
message.getMessageProperties().setDelay(5000);
return message;
}
});
}
延迟信息插件
使用官方提供的延迟信息插件,实现了对交换器的改造,直接变为死信交换器
地址:Scheduling Messages with RabbitMQ | RabbitMQ
使用:
基于注解方式:
@RabbitListener(bindings = @QueueBinding(
value = @Queue(name = "delay.queue", durable = "true"),
exchange = @Exchange(name = "delay.direct", delayed = "true"),
key = "delay"
))
public void listenDelayMessage(String msg){
log.info("接收到delay.queue的延迟消息:{}", msg);
}
基于bean方式:
@Bean
public DirectExchange delayExchange(){
return ExchangeBuilder
.directExchange("delay.direct") // 指定交换机类型和名称
.delayed() // 设置delay的属性为true
.durable(true) // 持久化
.build();
}
@Bean
public Queue delayedQueue(){
return new Queue("delay.queue");
}
@Bean
public Binding delayQueueBinding(){
return BindingBuilder.bind(delayedQueue()).to(delayExchange()).with("delay");
}
发送消息:
@Test
void testPublisherDelayMessage() {
// 1.创建消息
String message = "hello, delayed message";
// 2.发送消息,利用消息后置处理器添加消息头
rabbitTemplate.convertAndSend("delay.direct", "delay", message, new MessagePostProcessor() {
@Override
public Message postProcessMessage(Message message) throws AmqpException {
// 添加延迟消息属性
message.getMessageProperties().setDelay(5000);
return message;
}
});
}
注:
延迟消息插件内部会维护一个本地数据库表,同时使用Elang Timers功能实现计时。如果消息的延迟时间设置较长,可能会导致堆积的延迟消息非常多,会带来较大的CPU开销,同时延迟消息的时间会存在误差。
因此,不建议设置延迟时间过长的延迟消息。
更多推荐




所有评论(0)