RabbitMQ中实现延时消息

平常项目中很多场景需要使用延时消息处理,例如订单超过多久没有支付需要取消等。如何在消息中间件RabbitMQ中实现该功能?下面描述下使用Dead Letter Exchange实现延时消息场景,当然会有别的其他实现方式。

1. 什么是Dead Letter Exchange?

RabbitMQ中通常消息被直接发送到队列中或者从Exchange中Route到队列上后,消息如果被消费者消费完毕并确认后消息就会从Broker中被删除。
如果存在以下三种情况,同时队列上设置了Dead Letter Exchange,消息会被转送到Dead Letter Exchange中。

  • 消息被拒绝(basicReject或者basicNack) requeue=false
  • 消息存活时间超过了TTL预设值(x-message-ttl)
  • 队列满了

Dead Letter Exchange像平常的Exchange一样,可以设置它的BuiltinExchangeType,也可以为它绑定队列。
这里我们可以通过设定Dead Letter Exchange,并为它绑定一个队列,然后定义Consumer消费这个队列,就可以达到处理延时消息的功能。

2. 代码实例

流程先

I. 定义消息生产者

    /***
     * 消息发送者
     */
    static class NormalEXSend {
        private Connection conn;
        private Channel chnl;

        public NormalEXSend(String tag) throws IOException, TimeoutException {
            ConnectionFactory connFact = initConnFac();
            conn = connFact.newConnection();
            chnl = conn.createChannel();

            // 定义正常工作Exchange
            chnl.exchangeDeclare(WORKER_EXCHANGE_NAME, BuiltinExchangeType.DIRECT);

            // 定义 dead letter exchange
            chnl.exchangeDeclare(DELAY_EXCHANGE_NAME, BuiltinExchangeType.DIRECT);
            Map<String, Object> args = new HashMap<>();
            args.put("x-message-ttl", 60000); // timeout 1min
            args.put("x-dead-letter-exchange", DELAY_EXCHANGE_NAME);
            args.put("x-dead-letter-routing-key", DEAD_ROUTING_KEY);

            // 定义正常工作Queue同时设置dead letter exchange
            chnl.queueDeclare(WORKER_QUEUE_NAME, false, false, false, args);

            // 绑定到正常工作Exchange
            chnl.queueBind(WORKER_QUEUE_NAME, WORKER_EXCHANGE_NAME, tag);
        }

        /**
         * 发送消息
         * @param key
         * @param msg
         * @throws IOException
         */
        public void send(String key, String msg) throws IOException {
            AMQP.BasicProperties props = MessageProperties.PERSISTENT_TEXT_PLAIN;
            // send a message to a exchange
            chnl.basicPublish(WORKER_EXCHANGE_NAME, key, props, msg.getBytes());
            System.out.println(String.format("[%s|%s|Sender] send 【%s】 to exchange:%s", Thread.currentThread().getName(), System.currentTimeMillis(), msg, WORKER_EXCHANGE_NAME));
        }
    }

II. 定义延时消息处理者

其中receive方法中consumerhandleDelivery方法参数properties可以获取到消息的death原因properties.getHeaders().get("x-first-death-reason"),可能值rejected | expired | maxlen。此处可以根据判断此值去处理由于超时而引起death的消息(就是我们想要处理的延时消息)。

    /**
     * 延时消息处理者
     */
    static class DelayEXRecv {
        private Connection conn;
        private Channel chnl;

        public DelayEXRecv() throws IOException, TimeoutException {
            ConnectionFactory connFact = initConnFac();
            conn = connFact.newConnection();
            chnl = conn.createChannel();
            // 定义延时消息队列
            chnl.queueDeclare(DELAY_QUEUE_NAME, false, false, false, null);

            // 绑定到延时消息Exchange
            chnl.queueBind(DELAY_QUEUE_NAME, DELAY_EXCHANGE_NAME, DEAD_ROUTING_KEY);
        }

        /**
         * 接收消息
         * @throws IOException
         */
        public void receive() throws IOException {
            chnl.basicQos(1);
            // no auto ack
            boolean autoAck = false;
            chnl.basicConsume(DELAY_QUEUE_NAME, autoAck, new DefaultConsumer(chnl) {
                public void handleDelivery(String consumerTag, Envelope envelope,
                                           AMQP.BasicProperties properties, byte[] body) throws IOException {
                    String message = new String(body, "UTF-8");
                    // 打印出延时原因 rejected | expired | maxlen
                    // 项目中可以根据原因处理目标消息
                    System.out.println(String.format("[%s|%s|Delay_Receiver] received the delay msg 【%s】 from EXCHANGE: %s, the delay reason is: %s", Thread.currentThread().getName(), System.currentTimeMillis(), message, envelope.getExchange(), properties.getHeaders().get("x-first-death-reason")));
                    // 确认消息
                    chnl.basicAck(envelope.getDeliveryTag(), false);
                }
            });
        }
    }

III. 试验一把

    private static final String WORKER_EXCHANGE_NAME = "exchange.worker";
    private static final String DELAY_EXCHANGE_NAME = "exchange.delay";
    private static final String WORKER_QUEUE_NAME = "queue.worker";
    private static final String DELAY_QUEUE_NAME = "queue.delay";
    private static final String DEAD_ROUTING_KEY = "dead.key.message";

    public static void main(String[] args) {
        ExecutorService exec = Executors.newFixedThreadPool(2);
        exec.execute(new Runnable() {
            @Override
            public void run() {
                try {
                    String key = "worker";
                    NormalEXSend sender = new NormalEXSend(key);
                    for (int i =0; i < 5; i++) {
                        sender.send(key, String.format("YaYYY, one message, No.:%s!", i));
                        Thread.sleep(3000);
                    }
                } catch (IOException | TimeoutException | InterruptedException e) {
                    e.printStackTrace();
                }
            }
        });

        exec.execute(new Runnable() {
            @Override
            public void run() {
                try {
                    DelayEXRecv receiver = new DelayEXRecv();
                    receiver.receive();
                    System.out.println(String.format("[%s|%s|Delay_Receiver] Starting the Delay Msg Receiver process...", Thread.currentThread().getName(), System.currentTimeMillis()));
                } catch (IOException | TimeoutException e) {
                    e.printStackTrace();
                }

            }
        });

        exec.shutdown();
    }

IV. 打印结果

[pool-1-thread-2|1515750089010|Delay_Receiver] Starting the Delay Msg Receiver process...
[pool-1-thread-1|1515750089020|Sender] send 【YaYYY, one message, No.:0!】 to exchange:exchange.worker
[pool-1-thread-1|1515750092020|Sender] send 【YaYYY, one message, No.:1!】 to exchange:exchange.worker
[pool-1-thread-1|1515750095020|Sender] send 【YaYYY, one message, No.:2!】 to exchange:exchange.worker
[pool-1-thread-1|1515750098021|Sender] send 【YaYYY, one message, No.:3!】 to exchange:exchange.worker
[pool-1-thread-1|1515750101022|Sender] send 【YaYYY, one message, No.:4!】 to exchange:exchange.worker
[pool-2-thread-4|1515750149038|Delay_Receiver] received the delay msg 【YaYYY, one message, No.:0!】 from EXCHANGE: exchange.delay, the delay reason is: expired
[pool-2-thread-5|1515750152035|Delay_Receiver] received the delay msg 【YaYYY, one message, No.:1!】 from EXCHANGE: exchange.delay, the delay reason is: expired
[pool-2-thread-6|1515750155035|Delay_Receiver] received the delay msg 【YaYYY, one message, No.:2!】 from EXCHANGE: exchange.delay, the delay reason is: expired
[pool-2-thread-7|1515750158036|Delay_Receiver] received the delay msg 【YaYYY, one message, No.:3!】 from EXCHANGE: exchange.delay, the delay reason is: expired
[pool-2-thread-8|1515750161036|Delay_Receiver] received the delay msg 【YaYYY, one message, No.:4!】 from EXCHANGE: exchange.delay, the delay reason is: expired

可以看出消息是在制定延时的1min后才被获取消费。
Yayy, 至此结束。

参考:http://www.rabbitmq.com/dlx.html

©著作权归作者所有,转载或内容合作请联系作者
  • 序言:七十年代末,一起剥皮案震惊了整个滨河市,随后出现的几起案子,更是在滨河造成了极大的恐慌,老刑警刘岩,带你破解...
    沈念sama阅读 199,711评论 5 468
  • 序言:滨河连续发生了三起死亡事件,死亡现场离奇诡异,居然都是意外死亡,警方通过查阅死者的电脑和手机,发现死者居然都...
    沈念sama阅读 83,932评论 2 376
  • 文/潘晓璐 我一进店门,熙熙楼的掌柜王于贵愁眉苦脸地迎上来,“玉大人,你说我怎么就摊上这事。” “怎么了?”我有些...
    开封第一讲书人阅读 146,770评论 0 330
  • 文/不坏的土叔 我叫张陵,是天一观的道长。 经常有香客问我,道长,这世上最难降的妖魔是什么? 我笑而不...
    开封第一讲书人阅读 53,799评论 1 271
  • 正文 为了忘掉前任,我火速办了婚礼,结果婚礼上,老公的妹妹穿的比我还像新娘。我一直安慰自己,他们只是感情好,可当我...
    茶点故事阅读 62,697评论 5 359
  • 文/花漫 我一把揭开白布。 她就那样静静地躺着,像睡着了一般。 火红的嫁衣衬着肌肤如雪。 梳的纹丝不乱的头发上,一...
    开封第一讲书人阅读 48,069评论 1 276
  • 那天,我揣着相机与录音,去河边找鬼。 笑死,一个胖子当着我的面吹牛,可吹牛的内容都是我干的。 我是一名探鬼主播,决...
    沈念sama阅读 37,535评论 3 390
  • 文/苍兰香墨 我猛地睁开眼,长吁一口气:“原来是场噩梦啊……” “哼!你这毒妇竟也来了?” 一声冷哼从身侧响起,我...
    开封第一讲书人阅读 36,200评论 0 254
  • 序言:老挝万荣一对情侣失踪,失踪者是张志新(化名)和其女友刘颖,没想到半个月后,有当地人在树林里发现了一具尸体,经...
    沈念sama阅读 40,353评论 1 294
  • 正文 独居荒郊野岭守林人离奇死亡,尸身上长有42处带血的脓包…… 初始之章·张勋 以下内容为张勋视角 年9月15日...
    茶点故事阅读 35,290评论 2 317
  • 正文 我和宋清朗相恋三年,在试婚纱的时候发现自己被绿了。 大学时的朋友给我发了我未婚夫和他白月光在一起吃饭的照片。...
    茶点故事阅读 37,331评论 1 329
  • 序言:一个原本活蹦乱跳的男人离奇死亡,死状恐怖,灵堂内的尸体忽然破棺而出,到底是诈尸还是另有隐情,我是刑警宁泽,带...
    沈念sama阅读 33,020评论 3 315
  • 正文 年R本政府宣布,位于F岛的核电站,受9级特大地震影响,放射性物质发生泄漏。R本人自食恶果不足惜,却给世界环境...
    茶点故事阅读 38,610评论 3 303
  • 文/蒙蒙 一、第九天 我趴在偏房一处隐蔽的房顶上张望。 院中可真热闹,春花似锦、人声如沸。这庄子的主人今日做“春日...
    开封第一讲书人阅读 29,694评论 0 19
  • 文/苍兰香墨 我抬头看了看天上的太阳。三九已至,却和暖如春,着一层夹袄步出监牢的瞬间,已是汗流浃背。 一阵脚步声响...
    开封第一讲书人阅读 30,927评论 1 255
  • 我被黑心中介骗来泰国打工, 没想到刚下飞机就差点儿被人妖公主榨干…… 1. 我叫王不留,地道东北人。 一个月前我还...
    沈念sama阅读 42,330评论 2 346
  • 正文 我出身青楼,却偏偏与公主长得像,于是被迫代替她去往敌国和亲。 传闻我的和亲对象是个残疾皇子,可洞房花烛夜当晚...
    茶点故事阅读 41,904评论 2 341

推荐阅读更多精彩内容