2021-03-28_Rabbitmq延迟队列学习

Rabbitmq延迟队列学习

1概述

1.1延迟队列

延迟队列,即消息进入队列后不会立即被消费,只有到达指定时间后,才会被消费。

应用场景:

下单后,30分钟未支付,取消订单,回滚库存。

30分钟内字符,则从MQ消息队列中修改消息状态half-->可发送。

30分钟到, 未支付,则取消订单,回滚库存,删除消息。

在RabbitMQ中并未提供延迟队列功能。
但是可以使用:TTL+死信队列 组合实现延迟队列的效果。

2代码示例

2.1生产端

package com.kikop.workpattern.topic.producer;

import com.kikop.common.utils.DateUtil;
import com.kikop.common.utils.RabbitMQUtils;
import com.kikop.config.RabbitMQConfig;
import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.BuiltinExchangeType;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;

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

/**
 * @author kikop
 * @version 1.0
 * @project Name: myrabbitmqdemo
 * @file Name: WeatherBureau
 * @desc 功能描述 消费消息
 * @date 2020/9/6
 * @time 18:46
 * @by IDE: IntelliJ IDEA
 */
public class DelayedQueueProducer {

    /**
     * 单独为某条消息设置超时
     *
     * @throws IOException
     * @throws TimeoutException
     */
    private static void createDelayQueueTest() throws IOException, TimeoutException {

        Connection connection = RabbitMQUtils.getConnection();
        final Channel channel = connection.createChannel();

        // 定义死信交换机
        boolean durable_dlx_exchange = true;
        channel.exchangeDeclare(RabbitMQConfig.EXCHANGE_DLX_TOPIC, BuiltinExchangeType.TOPIC, durable_dlx_exchange);


        // 1.过期队列及属性设置
        Map<String, Object> arguments = new HashMap<String, Object>();

        // 1.1.设置整个队列过期时间
        // arguments.put("x-message-ttl", 10000);

        // 1.2.设置队列的私信参数(队列过期的后处理,移到死信交换机中)
        arguments.put("x-dead-letter-exchange", RabbitMQConfig.EXCHANGE_DLX_TOPIC);
        arguments.put("x-dead-letter-routing-key", "rk_mydlx");

        // 1.3.创建持久化队列
        boolean durable_queue = true;
        channel.queueDeclare(RabbitMQConfig.QUEUE_DELAY_TOPIC_PRIVATE, durable_queue, false, false, arguments);

        // 2.创建持久化转发代理
        boolean durable_exchange = true;
        channel.exchangeDeclare(RabbitMQConfig.EXCHANGE_DELAY_TOPIC_PRIVATE, BuiltinExchangeType.TOPIC, durable_exchange);
        String routingKey = "mydelay";
        channel.queueBind(RabbitMQConfig.QUEUE_DELAY_TOPIC_PRIVATE, RabbitMQConfig.EXCHANGE_DELAY_TOPIC_PRIVATE, routingKey);

        // 3.设置消息持久化
        String message = "购票超时,订单编号:000001,锁定的库存编号为:000110,即将执行回滚,(增加库存,作废订单)操作!";
        int durable_msg = 2; // 设置消息是否持久化,1: 非持久化 2:持久化
        AMQP.BasicProperties.Builder properties = new AMQP.BasicProperties().builder();
        properties.deliveryMode(durable_msg);  // 设置消息是否持久化,1: 非持久化 2:持久化

        // 4.发送 mock half消息
        properties.expiration(3 * 60 * 1000 + ""); // 过期时间3分钟
        System.out.println(DateUtil.getCurrentDateStr());
        channel.basicPublish(RabbitMQConfig.EXCHANGE_DELAY_TOPIC_PRIVATE, routingKey,
                properties.build(), message.getBytes("UTF-8"));

        channel.close();
        connection.close();
    }


    public static void main(String[] args) throws Exception {

        createDelayQueueTest();
        return;
    }

}

2.2消费端

package com.kikop.workpattern.topic.consumer;

import com.kikop.common.utils.DateUtil;
import com.kikop.common.utils.RabbitMQUtils;
import com.kikop.config.RabbitMQConfig;
import com.rabbitmq.client.*;

import java.io.IOException;
import java.util.concurrent.TimeoutException;

/**
 * @author kikop
 * @version 1.0
 * @project Name: myrabbitmqdemo
 * @file Name: BiaDu
 * @desc 功能描述 消费消息
 * @date 2020/9/6
 * @time 18:46
 * @by IDE: IntelliJ IDEA
 */
public class DelayedQueueConsumer {

    public static void main(String[] args) throws IOException, TimeoutException {

        delayConsumerTest();
        return;
    }


    private static void delayConsumerTest() throws IOException, TimeoutException {
        Connection connection = RabbitMQUtils.getConnection();
        final Channel channel = connection.createChannel();


        // begin_配置通道队列(可选,可在生产端完成)

        // 定义通道的非持久化队列
        channel.queueDeclare(RabbitMQConfig.QUEUE_DLX_TOPIC, false, false, false, null);

        // 定义需要消费的私信队列 QUEUE_DLX_TOPIC,从关联的死信交换机中抓取
        // 增加现有的一个转发代理,代理 EXCHANGE_DLX_TOPIC 必须先存着
        //  queueBind用于将队列与交换机绑定
        // 参数1:队列名
        // 参数2:交互机名
        // 参数三:路由key
        channel.queueBind(RabbitMQConfig.QUEUE_DLX_TOPIC, RabbitMQConfig.EXCHANGE_DLX_TOPIC, "rk_mydlx");
        // end_配置通道队列(可选)

        channel.basicQos(1);

        // 只消费死信队列里面的消息
        channel.basicConsume(RabbitMQConfig.QUEUE_DLX_TOPIC, false, new DefaultConsumer(channel) {
            @Override
            public void handleDelivery(String consumerTag, Envelope envelope, AMQP.BasicProperties properties, byte[] body) throws IOException {
                System.out.println(DateUtil.getCurrentDateStr() + "@信息:" + new String(body));

                channel.basicAck(envelope.getDeliveryTag(), false);
            }
        });
    }

}

2.3测试

// 9.延迟队列(TTL+DLX)

// 实际上是私有的
public static final String EXCHANGE_DELAY_TOPIC_PRIVATE = "ex_delay_topic_private";
public static final String QUEUE_DELAY_TOPIC_PRIVATE = "queue_delay_topic_private";


public static final String EXCHANGE_DLX_TOPIC = "ex_dlx_topic";
public static final String QUEUE_DLX_TOPIC = "queue_dlx_topic";

生产端:2021-03-28 23:02:47

消费端:2021-03-28 23:05:47@信息:购票超时,订单编号:000001,锁定的库存编号为:000110,即将执行回滚,(增加库存,作废订单)操作!

参考

1RabbitMQ消息最终一致性解决方案

https://blog.csdn.net/weixin_43466542/article/details/101676944

2.RabbitMq(十一) 死信交换机DLX介绍及使用

https://blog.csdn.net/liuhenghui5201/article/details/107538775

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

推荐阅读更多精彩内容