(ImplementationScheduled Tasks)

场景: 比如未付款订单,超过一定时间后,系统自动取消订单并释放占有物品。 常用Solution: spring的 schedule Scheduled Tasks轮询Database 缺点: 消耗系统内存、增加了Database的压力、存在较大的时间误差 解决:rabbitmq的MessageTTL和死信Exchange结合

Message的TTL(Time To Live)

• Message的TTL就是Message的存活时间。 • RabbitMQ可以对Queue和Message分别SettingsTTL。 • 对QueueSettings就是Queue没有消费者连着的保留时间,也可以对每一个单独的Message做单独的 Settings。超过了这个时间,我们认为这个Message就死了,称之为死信。 • 如果QueueSettings了,Message也Settings了,那么会取小的。所以一个Message如果被路由到不同的队 列中,这个Message死亡的时间有可能不一样(不同的QueueSettings)。这里单讲单个Message的 TTL,因为它才是Implementation延迟任务的关键。可以通过SettingsMessage的expiration字段或者x- message-ttl属性来Settings时间,两者是一样的效果。

Dead Letter Exchanges(DLX)

• 一个Message在满足如下条件下,会进死信路由,记住这里是路由而不是Queue, 一个路由可以对应很多Queue。(什么是死信) • 一个Message被Consumer拒收了,并且rejectMethod的Parameter里requeue是false。也就是说不 会被再次放在Queue里,被其他消费者Usage。(basic.reject/ basic.nack)requeue=false • 上面的Message的TTL到了,Message过期了。 • Queue的长度限制满了。排在前面的Message会被丢弃或者扔到死信路由上 • Dead Letter Exchange其实就是一种普通的exchange,和Create其他 exchange没有两样。只是在某一个SettingsDead Letter Exchange的Queue中有 Message过期了,会自动触发Message的转发,发送到Dead Letter Exchange中去。 • 我们既可以控制Message在一段时间后变成死信,又可以控制变成死信的Message 被路由到某一个指定的交换机,结合二者,其实就可以Implementation一个Delayed Queue

• 手动ack&异常Message统一放在一个Queue处理建议的两种方式 • catch异常后,手动发送到指定Queue,然后Usagechannel给rabbitmq确认Message已消费 • 给Queue绑定死信Queue,Usagenack(requque为false)确认Message消费失败

Delayed QueueImplementation

SpringBoot中使用延时队列

• 1、Queue、Exchange、Binding可以@Bean进去 • 2、监听Message的Method可以有三种Parameter(不分数量,顺序) • Object content, Message message, Channel channel • 3、channel可以用来拒绝Message,否则自动ack;

代码实现

发送Message到MQ:

import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.ResponseBody;
import org.springframework.web.bind.annotation.RestController;

import java.util.Date;
import java.util.UUID;

/**
 * @description: TODO
 * @author: <a href="mailto:batis@foxmail.com">清风</a>
 * @date: 2022/3/13 15:41
 * @version: 1.0
 */
@RestController
public class RabbitMQControoler {

    @Autowired
    RabbitTemplate rabbitTemplate;

    @ResponseBody
    @GetMapping("/createOrder")
    public String sendMessage(){
        OrderEntity entity = new OrderEntity();
        entity.setOrderEn(UUID.randomUUID().toString());
        entity.setModifyTime(new Date());
        //发送消息到MQ
        rabbitTemplate.convertAndSend("order-event-exchange", "order.create.order",entity);
        return "OK";
    }
}

通过延迟Queue,将延迟的订单,发送到另一个Queue,通过监听过期延迟的Queue,接收延迟的订单。

import com.rabbitmq.client.Channel;
import org.springframework.amqp.core.*;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

import java.io.IOException;
import java.util.HashMap;
import java.util.Map;

/**
 * @description: 延迟队列
 * @author: <a href="mailto:batis@foxmail.com">清风</a>
 * @date: 2022/3/13 13:27
 * @version: 1.0
 */
@Configuration
public class MyRabbitMQConfig {

    @RabbitListener(queues="order.release.order.queue")
    public void listener(OrderEntity entity, Channel channel, Message message) throws IOException {
        System.out.println("收到过期的订单信息:准备关闭订单"+entity.getOrderSn());
        channel.basicAck(message.getMessageProperties().getDeliveryTag(),false);
    }


    /**
     * 容器中的Binding、Queue、Exchange都会自动创建(RabbitMQ没有的情况)
     * @return
     */

    //死信队列
    @Bean
    public Queue orderDelayQueue(){
        Map<String, Object> arguments = new HashMap<>();
        arguments.put("x-dead-Letter-exchange", "order-event-exchange");//交换机
        arguments.put("x-dead-letter-routing-key", "order.release.order");//路由键
        arguments.put("x-message-ttl", 60000);//过期时间

        //public Queue(String name, boolean durable, boolean exclusive, boolean autoDelete, @Nullable Map<String, Object> arguments)
        Queue queue = new Queue("order.delay.queue",true,false,false,arguments);
        return queue;

    }

    //最后过期的队列
    @Bean
    public Queue orderReleaseQueue(){
        Queue queue = new Queue("order.release.order.queue",true,false,false);
        return queue;
    }

    //交换机
    @Bean
    public Exchange orderEventExchange(){
        return new TopicExchange("order-event-exchange", true, false);
    }

    //绑定交换机和死信队列
    @Bean
    public Binding orderCreateOrderBingding(){
        //public Binding(String destination, Binding.DestinationType destinationType, String exchange, String routingKey, @Nullable Map<String, Object> arguments)
        return new Binding("order-event-exchange",
                Binding.DestinationType.QUEUE,
                "order-event-exchange",
                "order.create.order",
                null);
    }

    //绑定交换机和最后过期的队列
    @Bean
    public Binding orderReleaseOrderBingding(){
        return new Binding("order-event-exchange",
                Binding.DestinationType.QUEUE,
                "order.release.order.queue",
                "order.release.order",
                null);
    }
}

如何保证消息可靠性-消息丢失

• 1、MessageLost • Message发送出去,由于网络Problem没有抵达Server • 做好Fault ToleranceMethod(try-catch),发送Message可能会网络失败,失败后要有重试机 制,可记录到Database,采用定期扫描重发的方式 • 做好Logging记录,每个Message状态是否都被Server收到都应该记录 • 做好定期重发,如果Message没有发送成功,定期去Database扫描未成功的Message进 行重发 • Message抵达Broker,Broker要将Message写入磁盘(Persistence)才算成功。此时Broker尚 未Persistence完成,宕机。 • publisher也必须加入确认CallbackMechanism,确认成功的Message,修改DatabaseMessage状态。 • 自动ACK的状态下。消费者收到Message,但没来得及Message然后宕机 • 一定开启手动ACK,消费成功才移除,失败或者没来得及处理就noAck并重 新入队

如何保证消息可靠性-消息重复

• 2、Message重复 • Message消费成功,Transaction已经提交,ack时,机器宕机。导致没有ack成功,Broker的Message 重新由unack变为ready,Concurrency送给其他消费者 • Message消费失败,由于重试Mechanism,自动又将Message发送出去 • 成功消费,ack时宕机,Message由unack变为ready,Broker又重新发送 • 消费者的业务消费Interface应该Design为Idempotency的。比如扣Library存有 工作单的状态标志 • Usage防重表(redis/mysql),发送Message每一个都有业务的唯 一标识,处理过就不用处理 • rabbitMQ的每一个Message都有redelivered字段,可以Get是否 是被重新投递过来的,而不是第一次投递过来的

• 2、Message重复 • Message消费成功,Transaction已经提交,ack时,机器宕机。导致没有ack成功,Broker的Message 重新由unack变为ready,Concurrency送给其他消费者 • Message消费失败,由于重试Mechanism,自动又将Message发送出去 • 成功消费,ack时宕机,Message由unack变为ready,Broker又重新发送 • 消费者的业务消费Interface应该Design为Idempotency的。比如扣Library存有 工作单的状态标志 • Usage防重表(redis/mysql),发送Message每一个都有业务的唯 一标识,处理过就不用处理 • rabbitMQ的每一个Message都有redelivered字段,可以Get是否

是被重新投递过来的,而不是第一次投递过来的

• 3、Message积压 • 消费者宕机积压 • 消费者消费能力不足积压 • 发送者发送流量太大 • 上线更多的消费者,进行正常消费 • 上线专门的Queue消费Service,将Message先批量取出来,记录Database,离线慢慢处理

CodeUsage,请参考RabbitMQ Operation Mechanism Installation请参考dockerInstallationRabbitMQ LearningRabbitMQ请参看Introduction to RabbitMQ