百度360必应搜狗淘宝本站头条
当前位置:网站首页 > 技术教程 > 正文

为了方便开发,我打算实现一个Redis 工具集

mhr18 2024-11-19 06:54 18 浏览 0 评论

前言

Redis 基本上是互联网公司必备的工具了,Redis的应用场景实在太多了,但是有很多相似的功能如果每个项目都要实现一遍就显得太麻烦了,所以为了方便,我打算开发一个基于 Redis 的工具集,尽量做到开箱即用。

目前实现功能

这个工具集并没有开发完成,实现了部分功能,如下图

简单介绍下已经实现的模块:

  • common : 整个项目公共模块,比如AOP工具等;
  • delay: Redis实现的延迟队列;
  • lock: Redis实现的分布式锁;
  • mq: Redis实现消息队列;
  • query: Redis实现分页模糊查询;
  • web: Redis实现web相关的功能; duplicate :防止重复提交;、

以上的这些模块都是已经实现的了,还有 社交、限流、幂等相关功能后面会陆续实现。

如何使用

  1. 引入 Maven 依赖(目前可以下载代码上传到自己的私服或者本地仓库,后面会推到 Maven 中央仓库)
  2. xml复制代码
  3. <dependency> <groupId>cn.org.wangchangjiu</groupId> <artifactId>redis-util-spring-boot-starter</artifactId> <version>1.0.0-SNAPSHOT</version> </dependency>
  4. 配置文件(application.yaml)开启各模块功能开关
  5. yaml复制代码
  6. redis: util: mq: enable: true delay: enable: true
  7. 实现消息发送者
  8. MQ消息发送:
  9. 延迟消息发送:
  10. 实现消息监听器
  11. MQ消息监听器:
  12. 延迟消息监听器:

MQ和delay实现细节

MQ实现细节

容器启动时,简单来说就是通过springboot自动装配,创建一些Bean,如下图:

值得注意的是,springboot3.X 自动装配方式有点变化,需要创建文件 META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports 文件,文件内容就直接写 自动配置类

RedisUtilAutoConfiguration 主自动装配类会 import 各个模块的自动装配类:

我们以 RedisStreamAutoConfiguration 为例:

该装配类生效需要显示打开,然后就是创建各种Bean。

最主要的Bean有:

  • RedisMessageConsumerManager: 该Bean实现了 BeanPostProcessor 接口,主要作用是,获取被注解 RedisMessageListener 修饰的方法,把信息封装在 RedisMessageConsumerContainer 对象里,方便后面反射调用。
  • StreamMessageListenerContainer:
    这个Bean主要是做 redis MQ 的配置,比如配置:一次最多获取多少条消息、没有消息时阻塞时间、执行任务的executor、错误处理器、以及消费组、是否自动ACK等配置,具体代码如下:
less复制代码@Bean(initMethod = "start", destroyMethod = "stop")
@DependsOn("redisMessageConsumerManager")
@ConditionalOnMissingBean
public StreamMessageListenerContainer<String, MapRecord<String, String, String>> streamMessageListenerContainer(@Autowired RedisMessageConsumerManager redisMessageConsumerManager,
                                                                                                                @Autowired RedisConnectionFactory redisConnectionFactory,
                                                                                                                @Autowired ErrorHandler errorHandler) {
    MyRedisStreamProperties.Options options = myRedisStreamProperties.getOptions();
    StreamMessageListenerContainer.StreamMessageListenerContainerOptions<String, MapRecord<String, String, String>> containerOptions =
            StreamMessageListenerContainer.StreamMessageListenerContainerOptions
                    .builder()
                    // 一次最多获取多少条消息
                    .batchSize(options.getBatchSize())
                    // 运行 Stream 的 poll task
                    .executor(getStreamMessageListenerExecutor())
                    // Stream 中没有消息时,阻塞多长时间,需要比 `spring.redis.timeout` 的时间小
                    .pollTimeout(options.getPollTimeout())
                    // 获取消息的过程或获取到消息给具体的消息者处理的过程中,发生了异常的处理
                    .errorHandler(errorHandler)
                    .build();

    StreamMessageListenerContainer<String, MapRecord<String, String, String>> streamMessageListenerContainer =
            StreamMessageListenerContainer.create(redisConnectionFactory, containerOptions);

    // 获取 被 RedisMessageListener 注解修饰的 bean
    Map<String, RedisMessageConsumerContainer> consumerContainerGroups =
            redisMessageConsumerManager.getConsumerContainerGroups();

    // 循环遍历,创建 消费组
    consumerContainerGroups.forEach((groupQueue, redisMessageConsumerContainer) -> {
        String[] groupQueues = groupQueue.split("#");

        // 创建消费组
        createGroups(groupQueues);

        RedisMessageListener redisMessageListener = redisMessageConsumerContainer.getRedisMessageListener();
        if(!redisMessageListener.useGroup()){
            // 独立消费 不使用组
            streamMessageListenerContainer.receive(StreamOffset.fromStart(groupQueues[1]), new DefaultGroupStreamListener(redisMessageConsumerContainer));
        } else {
            // 消费组 消费
            if(redisMessageListener.autoAck()){
                // 自动ACK
                streamMessageListenerContainer.receiveAutoAck(Consumer.from(groupQueues[0], "consumer:" + UUID.randomUUID()),
                        StreamOffset.create(groupQueues[1], ReadOffset.lastConsumed()), new DefaultGroupStreamListener(redisMessageConsumerContainer));
            } else {
                // 手动 ACK
                streamMessageListenerContainer.receive(Consumer.from(groupQueues[0], "consumer:" + UUID.randomUUID()),
                        StreamOffset.create(groupQueues[1], ReadOffset.lastConsumed()), new DefaultGroupStreamListener(redisMessageConsumerContainer));
            }
        }
    });
    return streamMessageListenerContainer;
}

/**
 *  创建消费组
 * @param groupQueues
 */
private void createGroups(String[] groupQueues) {
    // 判断是否存在队列Key
    if (stringRedisTemplate.hasKey(groupQueues[1])) {
        // 获取消费组 没有则创建
        StreamInfo.XInfoGroups groups = stringRedisTemplate.opsForStream().groups(groupQueues[1]);
        if (groups.isEmpty()) {
            stringRedisTemplate.opsForStream().createGroup(groupQueues[1], groupQueues[0]);
        } else {
            AtomicBoolean exists= new AtomicBoolean(false);
            groups.forEach(xInfoGroup -> {
                if (xInfoGroup.groupName().equals(groupQueues[0])){
                    exists.set(true);
                }
            });
            if(!exists.get()){
                stringRedisTemplate.opsForStream().createGroup(groupQueues[1], groupQueues[0]);
            }
        }
    } else {
        stringRedisTemplate.opsForStream().createGroup(groupQueues[1], groupQueues[0]);
    }
}

// todo 后面这个线程池也可以交由用户配置
private Executor getStreamMessageListenerExecutor() {
    AtomicInteger index = new AtomicInteger(1);
    int processors = Runtime.getRuntime().availableProcessors();
    ThreadPoolExecutor executor = new ThreadPoolExecutor(processors, processors, 0, TimeUnit.SECONDS,
            new LinkedBlockingDeque<>(), r -> {
        Thread thread = new Thread(r);
        thread.setName("async-stream-consumer-" + index.getAndIncrement());
        thread.setDaemon(true);
        return thread;
    });
    return executor;
}

发送消息流程:

redis 延迟队列的实现原理和这个差不多,主要是 redission延迟队列 + 自定义注解 + 反射,代码都差不多,我之前写过一篇基于 redission + kafka的延迟队列实现方式类似,只是把kafka那部分除去了,具体可以看那篇博文,地址: 基于 Redisson 和 Kafka 的延迟队列设计方案

后记

其他模块的设计细节后面再说,欢迎大家使用,也请大家多多提建议以及好的功能点,我都可以整合上去。

相关推荐

【推荐】一个开源免费、AI 驱动的智能数据管理系统,支持多数据库

如果您对源码&技术感兴趣,请点赞+收藏+转发+关注,大家的支持是我分享最大的动力!!!.前言在当今数据驱动的时代,高效、智能地管理数据已成为企业和个人不可或缺的能力。为了满足这一需求,我们推出了这款开...

Pure Storage推出统一数据管理云平台及新闪存阵列

PureStorage公司今日推出企业数据云(EnterpriseDataCloud),称其为组织在混合环境中存储、管理和使用数据方式的全面架构升级。该公司表示,EDC使组织能够在本地、云端和混...

对Java学习的10条建议(对java课程的建议)

不少Java的初学者一开始都是信心满满准备迎接挑战,但是经过一段时间的学习之后,多少都会碰到各种挫败,以下北风网就总结一些对于初学者非常有用的建议,希望能够给他们解决现实中的问题。Java编程的准备:...

SQLShift 重大更新:Oracle→PostgreSQL 存储过程转换功能上线!

官网:https://sqlshift.cn/6月,SQLShift迎来重大版本更新!作为国内首个支持Oracle->OceanBase存储过程智能转换的工具,SQLShift在过去一...

JDK21有没有什么稳定、简单又强势的特性?

佳未阿里云开发者2025年03月05日08:30浙江阿里妹导读这篇文章主要介绍了Java虚拟线程的发展及其在AJDK中的实现和优化。阅前声明:本文介绍的内容基于AJDK21.0.5[1]以及以上...

「松勤软件测试」网站总出现404 bug?总结8个原因,不信解决不了

在进行网站测试的时候,有没有碰到过网站崩溃,打不开,出现404错误等各种现象,如果你碰到了,那么恭喜你,你的网站出问题了,是什么原因导致网站出问题呢,根据松勤软件测试的总结如下:01数据库中的表空间不...

Java面试题及答案最全总结(2025版)

大家好,我是Java面试陪考员最近很多小伙伴在忙着找工作,给大家整理了一份非常全面的Java面试题及答案。涉及的内容非常全面,包含:Spring、MySQL、JVM、Redis、Linux、Sprin...

数据库日常运维工作内容(数据库日常运维 工作内容)

#数据库日常运维工作包括哪些内容?#数据库日常运维工作是一个涵盖多个层面的综合性任务,以下是详细的分类和内容说明:一、数据库运维核心工作监控与告警性能监控:实时监控CPU、内存、I/O、连接数、锁等待...

分布式之系统底层原理(上)(底层分布式技术)

作者:allanpan,腾讯IEG高级后台工程师导言分布式事务是分布式系统必不可少的组成部分,基本上只要实现一个分布式系统就逃不开对分布式事务的支持。本文从分布式事务这个概念切入,尝试对分布式事务...

oracle 死锁了怎么办?kill 进程 直接上干货

1、查看死锁是否存在selectusername,lockwait,status,machine,programfromv$sessionwheresidin(selectsession...

SpringBoot 各种分页查询方式详解(全网最全)

一、分页查询基础概念与原理1.1什么是分页查询分页查询是指将大量数据分割成多个小块(页)进行展示的技术,它是现代Web应用中必不可少的功能。想象一下你去图书馆找书,如果所有书都堆在一张桌子上,你很难...

《战场兄弟》全事件攻略 一般事件合同事件红装及隐藏职业攻略

《战场兄弟》全事件攻略,一般事件合同事件红装及隐藏职业攻略。《战场兄弟》事件奖励,事件条件。《战场兄弟》是OverhypeStudios制作发行的一款由xcom和桌游为灵感来源,以中世纪、低魔奇幻为...

LoadRunner(loadrunner录制不到脚本)

一、核心组件与工作流程LoadRunner性能测试工具-并发测试-正版软件下载-使用教程-价格-官方代理商的架构围绕三大核心组件构建,形成完整测试闭环:VirtualUserGenerator(...

Redis数据类型介绍(redis 数据类型)

介绍Redis支持五种数据类型:String(字符串),Hash(哈希),List(列表),Set(集合)及Zset(sortedset:有序集合)。1、字符串类型概述1.1、数据类型Redis支持...

RMAN备份监控及优化总结(rman备份原理)

今天主要介绍一下如何对RMAN备份监控及优化,这里就不讲rman备份的一些原理了,仅供参考。一、监控RMAN备份1、确定备份源与备份设备的最大速度从磁盘读的速度和磁带写的带度、备份的速度不可能超出这两...

取消回复欢迎 发表评论: