redis实现分布式锁)
前面有一章 Redis集群环境搭建 如果没有搭建环境的可以参考下也是一些爬坑记录如果操作有问题的欢迎留言讨论在准备好环境后我们就来使用redis 实现一下分布式锁的一个业务情况。前期准备1. 先创建一个商品管理项目来作为本次功能实现的demo2. 关于商品基本的增删改查代码就不做展示了自行去写或者参考我前面给出来的代码3. 把本次需要用的redis 及redisson的依赖引入!--redis -- dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-data-redis/artifactId /dependency dependency groupIdorg.apache.commons/groupId artifactIdcommons-pool2/artifactId /dependency !--redisson-- dependency groupIdorg.redisson/groupId artifactIdredisson-spring-boot-starter/artifactId version3.11.0/version /dependency4. 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_THROUGHPUT5. 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) { RedisTemplateObject, 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实现分布式锁的实例就到这儿了工作中使用可能没有必要那么全程序都是要去适应业务的业务量达不到像加锁或者其他操作这些都没有必要反而影响性能根据实际业务去完成适合的程序才是一个优秀的程序员毕竟要考虑成本。你写的很多用牛刀去杀鸡后面基础弱的又维护不了完全没有必要的结果