spring cloud alibaba 完整实现(六)redis实现分布式锁
2026/9/15 10:33:54 网站建设 项目流程

前面有一章 Redis集群环境搭建 如果没有搭建环境的可以参考下,也是一些爬坑记录,如果操作有问题的欢迎留言讨论,在准备好环境后我们就来使用redis 实现一下分布式锁的一个业务情况。

前期准备

1. 先创建一个商品管理项目来作为本次功能实现的demo

2. 关于商品基本的增删改查代码就不做展示了,自行去写或者参考我前面给出来的代码

3. 把本次需要用的redis 及redisson的依赖引入

<!--redis --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-data-redis</artifactId> </dependency> <dependency> <groupId>org.apache.commons</groupId> <artifactId>commons-pool2</artifactId> </dependency> <!--redisson--> <dependency> <groupId>org.redisson</groupId> <artifactId>redisson-spring-boot-starter</artifactId> <version>3.11.0</version> </dependency>

4. redisson.yml 配置文件内容 自己调整一下集群地址,或者其他相关配置如密码等

clusterServersConfig: # 连接空闲超时,单位:毫秒 默认10000 idleConnectionTimeout: 10000 pingTimeout: 1000 # 同任何节点建立连接时的等待超时。时间单位是毫秒 默认10000 connectTimeout: 10000 # 等待节点回复命令的时间。该时间从命令发送成功时开始计时。默认3000 timeout: 3000 # 命令失败重试次数 retryAttempts: 3 # 命令重试发送时间间隔,单位:毫秒 retryInterval: 1500 # 重新连接时间间隔,单位:毫秒 reconnectionTimeout: 10000 # 执行失败最大次数 failedAttempts: 3 # 密码 # password: test1234 # 单个连接最大订阅数量 subscriptionsPerConnection: 5 clientName: null # loadBalancer 负载均衡算法类的选择 loadBalancer: !<org.redisson.connection.balancer.RoundRobinLoadBalancer> {} #从节点发布和订阅连接的最小空闲连接数 slaveSubscriptionConnectionMinimumIdleSize: 1 #从节点发布和订阅连接池大小 默认值50 slaveSubscriptionConnectionPoolSize: 50 # 从节点最小空闲连接数 默认值32 slaveConnectionMinimumIdleSize: 32 # 从节点连接池大小 默认64 slaveConnectionPoolSize: 64 # 主节点最小空闲连接数 默认32 masterConnectionMinimumIdleSize: 32 # 主节点连接池大小 默认64 masterConnectionPoolSize: 64 # 订阅操作的负载均衡模式 subscriptionMode: SLAVE # 只在从服务器读取 readMode: SLAVE # 集群地址 nodeAddresses: - "redis://114.116.51.164:6379" - "redis://114.116.51.164:6380" - "redis://114.116.51.164:6381" - "redis://114.116.51.164:6382" - "redis://114.116.51.164:6383" - "redis://114.116.51.164:6384" # 对Redis集群节点状态扫描的时间间隔。单位是毫秒。默认1000 scanInterval: 1000 #这个线程池数量被所有RTopic对象监听器,RRemoteService调用者和RExecutorService任务共同共享。默认2 threads: 0 #这个线程池数量是在一个Redisson实例内,被其创建的所有分布式数据类型和服务,以及底层客户端所一同共享的线程池里保存的线程数量。默认2 nettyThreads: 0 # 编码方式 默认org.redisson.codec.JsonJacksonCodec codec: !<org.redisson.codec.JsonJacksonCodec> {} #传输模式 transportMode: NIO # 分布式锁自动过期时间,防止死锁,默认30000 lockWatchdogTimeout: 30000 # 通过该参数来修改是否按订阅发布消息的接收顺序出来消息,如果选否将对消息实行并行处理,该参数只适用于订阅发布消息的情况, 默认true keepPubSubOrder: true # 用来指定高性能引擎的行为。由于该变量值的选用与使用场景息息相关(NORMAL除外)我们建议对每个参数值都进行尝试。 # #该参数仅限于Redisson PRO版本。 #performanceMode: HIGHER_THROUGHPUT

5. RedisConfig配置类

package com.andy.protect.util; import com.fasterxml.jackson.annotation.JsonAutoDetect; import com.fasterxml.jackson.annotation.PropertyAccessor; import com.fasterxml.jackson.databind.ObjectMapper; import org.redisson.Redisson; import org.redisson.api.RedissonClient; import org.redisson.config.Config; import org.redisson.spring.data.connection.RedissonConnectionFactory; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.core.io.ClassPathResource; import org.springframework.data.redis.connection.RedisConnectionFactory; import org.springframework.data.redis.core.RedisTemplate; import org.springframework.data.redis.serializer.Jackson2JsonRedisSerializer; import org.springframework.data.redis.serializer.RedisSerializer; import org.springframework.data.redis.serializer.StringRedisSerializer; import java.io.IOException; @Configuration public class RedisConfig { @Bean public RedissonClient redisson() throws IOException { Config config = Config.fromYAML(new ClassPathResource("redisson.yml").getInputStream()); RedissonClient redisson = Redisson.create(config); return redisson; } @Bean public RedissonConnectionFactory redissonConnectionFactory(RedissonClient redisson) { return new RedissonConnectionFactory(redisson); } @Bean("redisTemplate") public RedisTemplate getRedisTemplate(RedisConnectionFactory redissonConnectionFactory) { RedisTemplate<Object, Object> redisTemplate = new RedisTemplate(); redisTemplate.setConnectionFactory(redissonConnectionFactory); redisTemplate.setValueSerializer(valueSerializer()); redisTemplate.setKeySerializer(keySerializer()); redisTemplate.setHashKeySerializer(keySerializer()); redisTemplate.setHashValueSerializer(valueSerializer()); return redisTemplate; } @Bean public RedisSerializer keySerializer() { return new StringRedisSerializer(); } @Bean public RedisSerializer valueSerializer() { Jackson2JsonRedisSerializer jackson2JsonRedisSerializer = new Jackson2JsonRedisSerializer(Object.class); ObjectMapper objectMapper = new ObjectMapper(); objectMapper.setVisibility(PropertyAccessor.ALL, JsonAutoDetect.Visibility.ANY); objectMapper.enableDefaultTyping(ObjectMapper.DefaultTyping.NON_FINAL); jackson2JsonRedisSerializer.setObjectMapper(objectMapper); return jackson2JsonRedisSerializer; } }

6. RedisUtil 工具类,redis集群里面也有这个,再贴一次吧

package com.andy.protect.util; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.data.redis.core.RedisTemplate; import org.springframework.stereotype.Component; import java.util.concurrent.TimeUnit; /** * redis使用简单工具类 */ @Component public class RedisUtil { @Autowired private RedisTemplate redisTemplate; public void set(String key, Object value) { redisTemplate.opsForValue().set(key, value); } public void set(String key, Object value, long timeout, TimeUnit unit) { redisTemplate.opsForValue().set(key, value, timeout, unit); } public boolean setIfAbsent(String key, Object value, long timeout, TimeUnit unit) { return redisTemplate.opsForValue().setIfAbsent(key, value, timeout, unit); } public <T> T get(String key, Class<?> T) { return (T) redisTemplate .opsForValue().get(key); } public void delete(String key){ redisTemplate.delete(key); } public String get(String key) { return (String) redisTemplate .opsForValue().get(key); } public Long decr(String key) { return redisTemplate .opsForValue().decrement(key); } public Long decr(String key, long delta) { return redisTemplate .opsForValue().decrement(key, delta); } public Long incr(String key) { return redisTemplate .opsForValue().increment(key); } public Long incr(String key, long delta) { return redisTemplate .opsForValue().increment(key, delta); } public void expire(String key, long time, TimeUnit unit) { redisTemplate.expire(key, time, unit); } }

到此,准备工作完成

接下来说明一下

解决思路及代码实现

思路我都直接写到代码中了,注释相对也算齐全,偷个懒就不再去分析了

package com.andy.protect.service.impl; import com.alibaba.fastjson.JSON; import com.andy.protect.dao.ProtectInfoDao; import com.andy.protect.entity.ProtectInfo; import com.andy.protect.service.ProtectService; import com.andy.protect.util.RedisUtil; import jodd.util.StringUtil; import org.redisson.api.RLock; import org.redisson.api.RReadWriteLock; import org.redisson.api.RedissonClient; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; import java.util.Random; import java.util.concurrent.TimeUnit; /** * 分布式锁实现热点数据查询问题 * * 解决思路: * 前置准备,将redis依赖及redisson(加锁解决redis相关问题)依赖加上 * 1.为了防止高并发请求,所以第一步应该是在查询时先查询redis缓存信息,没有时再到数据库查询 * a.当查询到了缓存数据,则更新缓存过期时间,并返回 * b.没有查询到缓存,查询数据库,并且在查询成功后,将数据添加到缓存中 * 2.当数据有更新情况发生,缓存应该也对应需要更新,优先选择直接删除缓存,下次查询的时候再自动重建到redis,避免并发写覆盖 * 3.为了防止缓存雪崩(大量数据同时失效)问题,所以我们要给过期时间上加一个随机数,错开大批量key的过期时刻 * 4.当缓存没有预热,大量请求打入,此时就会有缓存击穿风险,所以需要互斥锁来保证当没有缓存时,只有一个请求来进行缓存重建,其余请求等待或降级 * 5.用读写锁来保证缓存与数据库数据不一致问题,update用写锁独占执行,get重建数据时用读锁共享执行,降低锁粒度 * 6.所有锁都配置显式过期时间兜底,避免服务宕机锁永久滞留导致死锁 * 7.禁止静默吞异常,所有异常至少打印日志,方便线上排查缓存写失败导致的不一致问题 */ @Service public class ProtectServiceImpl implements ProtectService { // 日志对象,用于记录缓存操作异常,避免线上问题无迹可查 private static final Logger log = LoggerFactory.getLogger(ProtectServiceImpl.class); // 默认正常业务数据缓存时间:1天,单位秒 static final long CACHE_TIME = 60 * 60 * 24; // 定义商品缓存的key前缀,后续拼接业务id组成完整缓存key static final String PROTECT_KEY = "protect:cache:"; // 互斥锁的名称前缀:热点缓存重建专用锁,保证同一id下只有一个请求去查库重建缓存 public static final String LOCK_PRODUCT_HOT_CACHE_CREATE_PREFIX = "lock:product:hot_cache_create:"; // 读写锁的名称前缀:用于数据更新和缓存重建阶段,保证数据库和缓存的数据一致性 public static final String LOCK_PRODUCT_UPDATE = "lock:product:update:"; // 空值缓存哨兵标识:存入Redis的空值占位符,用于防御缓存穿透,避免不存在的id反复穿透到数据库 static final String EMPTY_CACHE = "{}"; // 空值缓存默认最大过期时间:90秒,不存在的短时间内不会再反复打库 private static final long NULL_CACHE_MAX_SECONDS = 90; // 缓存重建互斥锁最大持有时间:30秒,兜底避免锁永久不释放 private static final long CACHE_REBUILD_LOCK_SECONDS = 30; // 写锁最大持有时间:30秒,兜底避免写锁永久不释放 private static final long WRITE_LOCK_EXPIRE_SECONDS = 30; // 缓存重建锁最大等待获取时间:200毫秒,高并发场景下避免大量请求长时间排队打满线程池 private static final long TRY_LOCK_WAIT_MILLISECONDS = 200; @Autowired private ProtectInfoDao protectInfoDao; @Autowired private RedisUtil redisUtil; /** * redisson注入:使用RedissonClient接口而非具体类,符合面向接口编程规范,避免多实例注入歧义 */ @Autowired private RedissonClient redissonClient; @Override public ProtectInfo add(ProtectInfo protectInfo) { // 先执行数据库新增操作,保证库侧数据写入成功 ProtectInfo add = protectInfoDao.add(protectInfo); try { // 新增成功后,同步写入缓存,让新数据直接预热到缓存中 redisUtil.set(PROTECT_KEY + add.getId(), JSON.toJSONString(add), getCacheTime(), TimeUnit.SECONDS); } catch (Exception e) { // 缓存写入失败不影响主流程,但必须打印错误日志,便于后续排查缓存不一致问题 log.error("新增数据写入缓存失败, id={}", add.getId(), e); } return add; } @Override public boolean delete(ProtectInfo protectInfo) { // 先执行数据库删除操作,保证库侧数据删除成功 boolean delete = protectInfoDao.delete(protectInfo); try { if (delete) { // 数据库删除成功后,同步删除对应缓存,避免后续请求读到已删除的旧缓存数据 redisUtil.delete(PROTECT_KEY + protectInfo.getId()); } } catch (Exception e) { // 缓存删除失败不影响主流程,打印错误日志便于后续人工清理脏缓存 log.error("删除数据清理缓存失败, id={}", protectInfo.getId(), e); } return delete; } @Override public boolean update(ProtectInfo protectInfo) { boolean update = false; // 获取当前商品id对应的读写锁,用于保证更新操作的独占性 RReadWriteLock readWriteLock = redissonClient.getReadWriteLock(LOCK_PRODUCT_UPDATE + protectInfo.getId()); RLock writeLock = readWriteLock.writeLock(); // 加写锁并设置显式过期时间兜底,避免服务宕机导致锁永久滞留 writeLock.lock(WRITE_LOCK_EXPIRE_SECONDS, TimeUnit.SECONDS); try { // 先更新数据库,保证库侧数据为最新状态 update = protectInfoDao.update(protectInfo); if (update) { // 优先选择删除缓存而非直接set新值,后续读请求miss时会自动重建最新缓存,天然规避并发写缓存覆盖的问题 redisUtil.delete(PROTECT_KEY + protectInfo.getId()); } } finally { // 无论更新成功失败,最终都要释放写锁,避免锁资源泄露 writeLock.unlock(); } return update; } @Override public ProtectInfo get(int id) { ProtectInfo protectInfo; // 拼接当前查询id对应的完整缓存key String protectKey = PROTECT_KEY + id; // 第一阶段:快速路径,不加锁直接查询缓存,99%的正常请求会直接从这里返回,性能最高 protectInfo = getProtectInfoByCache(protectKey); if (protectInfo != null) { return protectInfo; } /**-------上面是缓存快速查询操作,下面是缓存miss后的数据库查询+重建阶段---------*/ // 获取当前商品id对应的热点缓存重建互斥锁,保证同一时刻只有一个请求去查库重建缓存 RLock rebuildLock = redissonClient.getLock(LOCK_PRODUCT_HOT_CACHE_CREATE_PREFIX + id); boolean lockAcquired = false; try { // 使用tryLock替代无限阻塞lock,设置最大等待时间,高并发下避免大量请求无限排队打满服务线程池 lockAcquired = rebuildLock.tryLock(TRY_LOCK_WAIT_MILLISECONDS, CACHE_REBUILD_LOCK_SECONDS, TimeUnit.MILLISECONDS); if (!lockAcquired) { // 没有抢到重建锁的请求,直接降级走数据库查询,避免大量请求在锁上长时间等待 return protectInfoDao.get(id); } // double-check二次校验缓存,抢到锁后再次查询缓存,防止前面已经有线程完成了缓存重建,后续抢到锁的线程重复查库 protectInfo = getProtectInfoByCache(protectKey); if (protectInfo != null) { return protectInfo; } /*********使用读锁来保证数据库与缓存不一致问题***********/ // 获取当前商品id对应的读锁,共享锁不互斥其他读请求,只和更新的写锁互斥,保证重建缓存过程中不会被写操作打断 RReadWriteLock readWriteLock = redissonClient.getReadWriteLock(LOCK_PRODUCT_UPDATE + id); RLock readLock = readWriteLock.readLock(); readLock.lock(); try { // 缓存中没有数据,查询数据库 protectInfo = protectInfoDao.get(id); if (protectInfo == null) { // 查询不到数据,将空值哨兵存入缓存,短时间内拦截相同不存在id的请求,防止缓存穿透 redisUtil.set(PROTECT_KEY + id, EMPTY_CACHE, nullCacheTime(), TimeUnit.SECONDS); return null; } // 查询到有效数据,写入缓存,完成缓存重建 redisUtil.set(PROTECT_KEY + protectInfo.getId(), JSON.toJSONString(protectInfo), getCacheTime(), TimeUnit.SECONDS); } catch (Exception e) { // 数据库查询或缓存写入异常,打印日志,直接降级返回库侧结果,避免整个请求链路不可用 log.error("缓存重建阶段查询库或写缓存失败, id={}", id, e); // 异常场景下直接查库兜底,保证用户请求能拿到结果 return protectInfoDao.get(id); } finally { // 无论重建成功失败,最终都要释放读锁,避免锁资源泄露 readLock.unlock(); } } catch (Exception e) { log.error("热点缓存重建流程异常, id={}", id, e); // 缓存重建流程异常,直接降级查库兜底,保证高可用 return protectInfoDao.get(id); } finally { // 只有当前线程成功抢到锁的情况下,才去释放锁,避免Redisson释放未持有锁抛出IllegalMonitorStateException if (lockAcquired) { rebuildLock.unlock(); } } return protectInfo; } /** * 返回业务数据缓存时间:基础1天 + 0-4小时随机偏移 * 错开大量key的过期时刻,防止大批量热点key同时失效引发缓存雪崩 * @return 最终缓存过期秒数 */ public long getCacheTime() { // 生成0到4的随机整数,乘以小时换算为秒,叠加到基础缓存时间上 int randomOffsetHour = new Random().nextInt(5); return CACHE_TIME + randomOffsetHour * 60 * 60; } /** * 返回空值缓存的过期时间:60秒 + 0-29秒随机偏移 * 短时间内拦截不存在的id反复穿透到数据库,同时也不会让空值在Redis中永久滞留 * @return 空值缓存过期秒数 */ public long nullCacheTime() { return 60 + new Random().nextInt(30); } /** * 获取缓存数据 * 有有效业务数据就返回ProtectInfo对象 * 命中空值哨兵直接返回null,交由上层逻辑处理数据不存在的场景 * 缓存中无数据就直接返回null */ public ProtectInfo getProtectInfoByCache(String key) { ProtectInfo product = null; String productStr = redisUtil.get(key); // 判断缓存值是否非空非空串 if (!StringUtil.isEmpty(productStr)) { // 查询到空值哨兵,直接返回null,不返回空对象,避免上层调用方取属性触发NPE if (EMPTY_CACHE.equals(productStr)) { return null; } // 正常解析缓存中的JSON字符串为业务对象 product = JSON.parseObject(productStr, ProtectInfo.class); // 只有当当前缓存剩余过期时间小于总过期时间的1/3时,才执行续期 // 避免每次读到缓存都续期,导致热点缓存被无限续期、旧数据永远无法自然过期 long ttl = redisUtil.getExpire(key, TimeUnit.SECONDS); if (ttl != -1 && ttl < CACHE_TIME / 3) { redisUtil.expire(key, getCacheTime(), TimeUnit.SECONDS); } } return product; } }

这个代码基本也对redis缓存的常见问题做了处理,这个有必要说明一下

1.缓存数据库不一致,指更新了数据库,缓存数据可能还是老数据

通过更新时删除和redisson读写锁来保证,但是没办法100% 还可以使用zookeeper来做,不过redis效率相对更高,CP AP总要牺牲一个

2.缓存击穿(缓存失效)大量数据失效,击穿缓存,打到数据库

使用过期时间加随机数避免大量数据同时失效

3.缓存重建缓存还没有数据,海量请求同时进入,每个人都在查询并保存缓存

加锁只允许一个请求重建缓存 还可以提前往缓存插入数据(比如秒杀)

4.缓存穿透击穿和穿透的区别就是穿透是连数据库也没有数据(黑客攻击)不停请求数据库

可以将查不到的数据也放到redis,也可以使用布隆过滤器

5.缓存雪崩和缓存击穿类似的,区别是击穿一般指一条数据高并发请求,雪崩是不同的数据大批量过期

也可通过随机时间,热点数据不过期,或者分布到不同服务器等方案解决

总结

关于redis实现分布式锁的实例就到这儿了,工作中使用可能没有必要那么全,程序都是要去适应业务的,业务量达不到,像加锁或者其他操作这些都没有必要,反而影响性能,根据实际业务去完成适合的程序才是一个优秀的程序员(毕竟要考虑成本)。你写的很多,用牛刀去杀鸡,后面基础弱的又维护不了,完全没有必要的结果

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询