IT袋

当前位置:主页 > 经验教程 > 建站编程 >

关于Redisson延迟队列的一些思考

关于Redisson延迟队列的一些思考(2)

时间:2024-01-29 01:29:29 来源:IT袋 作者:马勇
导读:关于Redisson延迟队列的一些思考,Redisson会怎么做延迟队列 首先提供一段简单的Redisson延迟队列的Demo代码给大家学习使用。 maven依赖: dependency> groupId>org.redisson/groupId> artifactId>redisson/artif

关于Redisson延迟队列的一些思考

Redisson会怎么做延迟队列

首先提供一段简单的Redisson延迟队列的Demo代码给大家学习使用。

maven依赖:

   <dependency>
      <groupId>org.redisson</groupId>
      <artifactId>redisson</artifactId>
      <version>3.16.8</version>
    </dependency>

延迟队列投递方:


package org.idea.redission.framework.delay.queue;
import com.alibaba.fastjson2.JSON;
import org.redisson.Redisson;
import org.redisson.api.RBlockingQueue;
import org.redisson.api.RedissonClient;
import org.redisson.config.Config;
/**
 * @Author idea
 * @Date: Created in 10:21 2024/1/28
 * @Description
 */
public class DelayQueueConsumerMain {
    public static void main(String[] args) {
        Config config = new Config();
        config.useSingleServer().setAddress("redis://cloud.db:8801").setPassword("pwd");
        RedissonClient client = Redisson.create(config);
        RBlockingQueue<String> blockingQueue = client.getBlockingQueue("delay_queue");
        new Thread(new Runnable() {
            @Override
            public void run() {
                System.out.println("开始拉去延迟消息");
                while (true) {
                    try {
                        String item = blockingQueue.take();
                        long currentTime = System.currentTimeMillis();
                        MessageModel messageModel = JSON.parseObject(item, MessageModel.class);
                        System.out.println("获取延迟消息,投递时间" + (currentTime - messageModel.getPushTime()) + "ms前,content" + messageModel.getContent());
                    } catch (InterruptedException e) {
                        e.printStackTrace();
                    }
                }
            }
        }).start();
    }
}

延迟队列消费方:

package org.idea.redission.framework.delay.queue;
import com.alibaba.fastjson2.JSON;
import lombok.extern.slf4j.Slf4j;
import org.redisson.Redisson;
import org.redisson.api.RBlockingQueue;
import org.redisson.api.RDelayedQueue;
import org.redisson.api.RedissonClient;
import org.redisson.config.Config;
import java.util.concurrent.TimeUnit;
/**
 * @Author idea
 * @Date: Created in 10:25 2024/1/28
 * @Description
 */
@Slf4j
public class DelayQueueProducerMain {
    public static void main(String[] args) throws InterruptedException {
        Config config = new Config();
        config.useSingleServer().setAddress("redis://cloud.db:8801").setPassword("pwd");
        RedissonClient client = Redisson.create(config);
        RBlockingQueue<String> blockingQueue = client.getBlockingQueue("delay_queue");
        RDelayedQueue<String> delayedQueue = client.getDelayedQueue(blockingQueue);
        int i = 0;
        while (true) {
            i++;
            TimeUnit.SECONDS.sleep(3);
            MessageModel messageModel = new MessageModel();
            messageModel.setContent("test-content-" + i);
            messageModel.setPushTime(System.currentTimeMillis());
            delayedQueue.offer(JSON.toJSONString(messageModel), 3, TimeUnit.SECONDS);
            System.out.println("投递第" + i + "条消息进延迟队列");
        }
    }
}

相关阅读