redisTemplate阻塞式处理消息队列

Posted 陈努力丶

tags:

篇首语:本文由小常识网(cha138.com)小编为大家整理,主要介绍了redisTemplate阻塞式处理消息队列相关的知识,希望对你有一定的参考价值。

用redis中的List可以实现队列,这样可以用来做消息处理和任务调度的队列

Redis 消息队列

redis五种数据结构

队列生产者

package cn.stylefeng.guns.knowledge.modular.knowledge.schedule;

import lombok.extern.slf4j.Slf4j;
import org.springframework.data.redis.core.RedisTemplate;

import java.util.Random;
import java.util.UUID;

/**
 * <p>
 * 队列生产者
 * </p>
 *
 * @SINCE 2021/11/30 21:03
 * @AUTHOR dispark
 * @Date: 2021/11/30 21:03
 */
@Slf4j
public class QueueProducer implements Runnable 

    /**
     * 生产者队列 key
     */
    public static final String QUEUE_PRODUCTER = "queue-producter";

    private RedisTemplate<String, Object> redisTemplate;

    public QueueProducer(RedisTemplate<String, Object> redisTemplate) 
        this.redisTemplate = redisTemplate;
    

    @Override
    public void run() 
        Random random = new Random();
        while (true) 
            try 
                Thread.sleep(random.nextInt(600) + 600);
                // 1.模拟生成一个任务
                UUID queueProducerId = UUID.randomUUID();
                // 2.将任务插入任务队列:queue-producter
                redisTemplate.opsForList().leftPush(QUEUE_PRODUCTER, queueProducerId.toString());
                log.info("生产一条数据 >>> ", queueProducerId.toString());
             catch (Exception e) 
                e.printStackTrace();
            
        
    


队列消费者

package cn.stylefeng.guns.knowledge.modular.knowledge.schedule;

import lombok.extern.slf4j.Slf4j;
import org.springframework.data.redis.core.RedisTemplate;

import java.util.Random;

/**
 * <p>
 * 队列消费者
 * </p>
 *
 * @SINCE 2021/11/30 21:14
 * @AUTHOR dispark
 * @Date: 2021/11/30 21:14
 */
@Slf4j
public class QueueConsumer implements Runnable 
    public static final String QUEUE_PRODUCTER = "queue-producter";
    public static final String TMP_QUEUE = "tmp-queue";

    private RedisTemplate<String, Object> redisTemplate;

    public QueueConsumer(RedisTemplate<String, Object> redisTemplate) 
        this.redisTemplate = redisTemplate;
    

    /**
     * 功能描述: 取值 - <brpop:阻塞式> - 推荐使用
     *
     * @author dispark
     * @date 2021/11/30 21:17
     */
    @Override
    public void run() 
        Random random = new Random();
        while (true) 
            // 1.从任务队列"queue-producter"中获取一个任务,并将该任务放入暂存队列"tmp-queue"
            Long ququeConsumerId = redisTemplate.opsForList().rightPush(QUEUE_PRODUCTER, TMP_QUEUE);
            // 2.处理任务----纯属业务逻辑,模拟一下:睡觉
            try 
                Thread.sleep(1000);
             catch (InterruptedException e) 
                e.printStackTrace();
            
            // 3.模拟成功和失败的偶然现象,模拟失败的情况,概率为2/13
            if (random.nextInt(13) % 7 == 0) 
                // 4.将本次处理失败的任务从暂存队列"tmp-queue"中,弹回任务队列"queue-producter"
                redisTemplate.opsForList().rightPush(TMP_QUEUE, QUEUE_PRODUCTER);
                log.info(ququeConsumerId + "处理失败,被弹回任务队列");
             else 
                // 5. 模拟成功的情况,将本次任务从暂存队列"tmp-queue"中清除
                redisTemplate.opsForList().rightPop(TMP_QUEUE);
                log.info(ququeConsumerId + "处理成功,被清除");
            
        
    


测试类

    @Test
    public void QueueThreadTotalEntry() throws Exception 
        // 1.启动一个生产者线程,模拟任务的产生
        new Thread(new QueueProducer(redisTemplate)).start();
        Thread.sleep(15000);
        // 2.启动一个线程者线程,模拟任务的处理
        new Thread(new QueueConsumer(redisTemplate)).start();
        // 3.主线程
        Thread.sleep(Long.MAX_VALUE);
    

并发情况下使用increment递增

线程一:

Long increment = redisTemplate.opsForValue().increment("increment", 1L);
            log.info("队列消费者 >> increment递增: ", increment);

线程二:

Long increment = redisTemplate.opsForValue().increment("increment", 1L);
            log.info("生产者队列 >> increment递增: ", increment);

借鉴

redis实现消息队列&发布/订阅模式使用
redisTemplate处理/获取redis消息队列

以上是关于redisTemplate阻塞式处理消息队列的主要内容,如果未能解决你的问题,请参考以下文章

Web在线聊天室(12) --- 收发消息(单例模式+阻塞式队列)

使用阻塞式队列处理大数据

JAVA队列的使用

java 队列的使用(转载)

如何使用Redis 做队列操作

对象锁,CPU时间片,阻塞队列