基于redis和shedlock实现分布式锁(超简单)
一、背景 线上部署了两台服务器,通过nginx轮询的方式进行负载均衡。但是这样存在一个问题同一个用户的session共享问题。你或许会说,使用ipHash模式就可以解决session共享的问题,是的确实可以解决这个问题,但是同样也会带来另外一个问题,就是一台服务器很繁忙,另外一台服务器闲置的情况。所以为了避免服务器闲置的现象,我们采用了ip轮询和共享session存入redis的解决方案。不过今天要讲的主题不是这个,我们在单机测试的时候完全运行正常,但是部署了到正式的环境的时候出现了客户支付金额对应不上的情况。后来经过排查发现是并发问题,虽然单机做了并发处理,但是如果有多台服务器的时候就会出现服务器之间的并发问题。所以需要引入分布式锁,解决多服务器间的。 二、分布式锁的实现 1. jar包的引入 <dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-web</artifactId></dependency><dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-test</artifactId><scope>test</scope><exclusions><exclusion><groupId>org.junit.vintage</groupId><artifactId>junit-vintage-engine</artifactId></exclusion></exclusions></dependency><dependency><groupId>org.springframework.boot</groupId><artifactId>spring-boot-starter-data-redis</artifactId></dependency><dependency><groupId>net.javacrumbs.shedlock</groupId><artifactId>shedlock-provider-redis-spring</artifactId><version>2.3.0</version></dependency><dependency><groupId>org.apache.commons</groupId><artifactId>commons-pool2</artifactId><version>2.0</version></dependency><dependency><groupId>net.javacrumbs.shedlock</groupId><artifactId>shedlock-spring</artifactId><version>2.3.0</version></dependency><dependency><groupId>org.projectlombok</groupId><artifactId>lombok</artifactId></dependency><!--swagger--><!--https://mvnrepository.com/artifact/io.springfox/springfox-swagger-ui--><dependency><groupId>com.github.xiaoymin</groupId><artifactId>swagger-bootstrap-ui</artifactId><version>1.9.6</version></dependency><dependency><groupId>io.springfox</groupId><artifactId>springfox-swagger2</artifactId><version>2.9.2</version></dependency><dependency><groupId>org.aspectj</groupId><artifactId>aspectjweaver</artifactId><version>1.9.2</version></dependency> 2. redis的配置 配置文件 #redisredis.host=192.168.1.6redis.password=redis.port=6379redis.taskScheduler.poolSize=100redis.taskScheduler.defaultLockMaxDurationMinutes=10redis.default.timeout=10redisCache.expireTimeInMilliseconds=1200000 配置类 packagecom.example.redis_demo_limit.redis;importio.lettuce.core.ClientOptions;importio.lettuce.core.resource.ClientResources;importio.lettuce.core.resource.DefaultClientResources;importnet.javacrumbs.shedlock.core.LockProvider;importnet.javacrumbs.shedlock.provider.redis.spring.RedisLockProvider;importnet.javacrumbs.shedlock.spring.ScheduledLockConfiguration;importnet.javacrumbs.shedlock.spring.ScheduledLockConfigurationBuilder;importorg.apache.commons.pool2.impl.GenericObjectPoolConfig;importorg.springframework.beans.factory.annotation.Value;importorg.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;importorg.springframework.context.annotation.Bean;importorg.springframework.context.annotation.Configuration;importorg.springframework.context.annotation.Primary;importorg.springframework.data.redis.connection.RedisConnectionFactory;importorg.springframework.data.redis.connection.RedisPassword;importorg.springframework.data.redis.connection.RedisStandaloneConfiguration;importorg.springframework.data.redis.connection.lettuce.LettuceConnectionFactory;importorg.springframework.data.redis.connection.lettuce.LettucePoolingClientConfiguration;importorg.springframework.data.redis.core.RedisTemplate;importjava.time.Duration;@ConfigurationpublicclassRedisConfig{@Value("${redis.host}")privateStringredisHost;@Value("${redis.port}")privateintredisPort;@Value("${redis.password}")privateStringpassword;@Value("${redis.taskScheduler.poolSize}")privateinttasksPoolSize;@Value("${redis.taskScheduler.defaultLockMaxDurationMinutes}")privateintlockMaxDuration;@Bean(destroyMethod="shutdown")ClientResourcesclientResources(){returnDefaultClientResources.create();}@BeanpublicRedisStandaloneConfigurationredisStandaloneConfiguration(){RedisStandaloneConfigurationredisStandaloneConfiguration=newRedisStandaloneConfiguration(redisHost,redisPort);if(password!=null&&!password.trim().equals("")){RedisPasswordredisPassword=RedisPassword.of(password);redisStandaloneConfiguration.setPassword(redisPassword);}returnredisStandaloneConfiguration;}@BeanpublicClientOptionsclientOptions(){returnClientOptions.builder().disconnectedBehavior(ClientOptions.DisconnectedBehavior.REJECT_COMMANDS).autoReconnect(true).build();}@BeanLettucePoolingClientConfigurationlettucePoolConfig(ClientOptionsoptions,ClientResourcesdcr){returnLettucePoolingClientConfiguration.builder().poolConfig(newGenericObjectPoolConfig()).clientOptions(options).clientResources(dcr).build();}@BeanpublicRedisConnectionFactoryconnectionFactory(RedisStandaloneConfigurationredisStandaloneConfiguration,LettucePoolingClientConfigurationlettucePoolConfig){returnnewLettuceConnectionFactory(redisStandaloneConfiguration,lettucePoolConfig);}@Bean@ConditionalOnMissingBean(name="redisTemplate")@PrimarypublicRedisTemplate<Object,Object>redisTemplate(RedisConnectionFactoryredisConnectionFactory){RedisTemplate<Object,Object>template=newRedisTemplate<>();template.setConnectionFactory(redisConnectionFactory);returntemplate;}@BeanpublicLockProviderlockProvider(RedisConnectionFactoryconnectionFactory){returnnewRedisLockProvider(connectionFactory);}@BeanpublicScheduledLockConfigurationtaskSchedulerLocker(LockProviderlockProvider){returnScheduledLockConfigurationBuilder.withLockProvider(lockProvider).withPoolSize(tasksPoolSize).withDefaultLockAtMostFor(Duration.ofMinutes(lockMaxDuration)).build();}} 操作类 packagecom.example.redis_demo_limit.redis;publicinterfaceDataCacheRepository<T>{booleanadd(Stringcollection,Stringhkey,Tobject,Longtimeout);booleandelete(Stringcollection,Stringhkey);Tfind(Stringcollection,Stringhkey,Class<T>tClass);BooleanisAvailable();/***redis加锁**@paramkey*@paramsecond*@return*/Booleanlock(Stringkey,Stringvalue,Longsecond);ObjectgetValue(Stringkey);/***redis解锁**@paramkey*@return*/voidunLock(Stringkey);voidsetIfAbsent(Stringkey,longvalue,longttl);voidincrement(Stringkey);Longget(Stringkey);voidset(Stringkey,longvalue,longttl);voidset(Objectkey,Objectvalue,longttl);ObjectgetByKey(Stringkey);voidgetLock(Stringkey,StringclientID)throwsException;voidreleaseLock(Stringkey,StringclientID);booleanhasKey(Stringkey);} 实现类 packagecom.example.redis_demo_limit.redis;importcom.fasterxml.jackson.databind.ObjectMapper;importlombok.extern.slf4j.Slf4j;importorg.slf4j.Logger;importorg.slf4j.LoggerFactory;importorg.springframework.beans.factory.annotation.Autowired;importorg.springframework.beans.factory.annotation.Value;importorg.springframework.data.redis.core.RedisTemplate;importorg.springframework.data.redis.core.ValueOperations;importorg.springframework.data.redis.support.atomic.RedisAtomicLong;importorg.springframework.stereotype.Repository;importjava.time.Duration;importjava.util.TimeZone;importjava.util.concurrent.TimeUnit;@Slf4j@RepositorypublicclassCacheRepository<T>implementscom.example.redis_demo_limit.redis.DataCacheRepository<T>{privatestaticfinalObjectMapperOBJECT_MAPPER;privatestaticfinalTimeZoneDEFAULT_TIMEZONE=TimeZone.getTimeZone("UTC");static{OBJECT_MAPPER=newObjectMapper();OBJECT_MAPPER.setTimeZone(DEFAULT_TIMEZONE);}Loggerlogger=LoggerFactory.getLogger(CacheRepository.class);@AutowiredRedisTemplatetemplate;//andwe'reinbusiness@Value("${redis.default.timeout}00")LongdefaultTimeOut;publicbooleanaddPermentValue(Stringcollection,Stringhkey,Tobject){try{StringjsonObject=OBJECT_MAPPER.writeValueAsString(object);template.opsForHash().put(collection,hkey,jsonObject);returntrue;}catch(Exceptione){logger.error("Unabletoaddobjectofkey{}tocachecollection'{}':{}",hkey,collection,e.getMessage());returnfalse;}}@Overridepublicbooleanadd(Stringcollection,Stringhkey,Tobject,Longtimeout){LonglocalTimeout;if(timeout==null){localTimeout=defaultTimeOut;}else{localTimeout=timeout;}try{StringjsonObject=OBJECT_MAPPER.writeValueAsString(object);template.opsForHash().put(collection,hkey,jsonObject);template.expire(collection,localTimeout,TimeUnit.SECONDS);returntrue;}catch(Exceptione){logger.error("Unabletoaddobjectofkey{}tocachecollection'{}':{}",hkey,collection,e.getMessage());returnfalse;}}@Overridepublicbooleandelete(Stringcollection,Stringhkey){try{template.opsForHash().delete(collection,hkey);returntrue;}catch(Exceptione){logger.error("Unabletodeleteentry{}fromcachecollection'{}':{}",hkey,collection,e.getMessage());returnfalse;}}@OverridepublicTfind(Stringcollection,Stringhkey,Class<T>tClass){try{StringjsonObj=String.valueOf(template.opsForHash().get(collection,hkey));returnOBJECT_MAPPER.readValue(jsonObj,tClass);}catch(Exceptione){if(e.getMessage()==null){logger.error("Entry'{}'doesnotexistincache",hkey);}else{logger.error("Unabletofindentry'{}'incachecollection'{}':{}",hkey,collection,e.getMessage());}returnnull;}}@OverridepublicBooleanisAvailable(){try{returntemplate.getConnectionFactory().getConnection().ping()!=null;}catch(Exceptione){logger.warn("Redisserverisnotavailableatthemoment.");}returnfalse;}@OverridepublicBooleanlock(Stringkey,Stringvalue,Longsecond){Booleanabsent=template.opsForValue().setIfAbsent(key,value,second,TimeUnit.SECONDS);returnabsent;}@OverridepublicObjectgetValue(Stringkey){returntemplate.opsForValue().get(key);}@OverridepublicvoidunLock(Stringkey){template.delete(key);}@Overridepublicvoidincrement(Stringkey){RedisAtomicLongcounter=newRedisAtomicLong(key,template.getConnectionFactory());counter.incrementAndGet();}@OverridepublicvoidsetIfAbsent(Stringkey,longvalue,longttl){ValueOperations<String,Object>ops=template.opsForValue();ops.setIfAbsent(key,value,Duration.ofSeconds(ttl));}@OverridepublicLongget(Stringkey){RedisAtomicLongcounter=newRedisAtomicLong(key,template.getConnectionFactory());returncounter.get();}@Overridepublicvoidset(Stringkey,longvalue,longttl){RedisAtomicLongcounter=newRedisAtomicLong(key,template.getConnectionFactory());counter.set(value);counter.expire(ttl,TimeUnit.SECONDS);}@Overridepublicvoidset(Objectkey,Objectvalue,longttl){template.opsForValue().set(key,value,ttl,TimeUnit.SECONDS);}@OverridepublicObjectgetByKey(Stringkey){returntemplate.opsForValue().get(key);}@OverridepublicvoidgetLock(Stringkey,StringclientID)throwsException{Booleanlock=false;//重试3次,每间隔1秒重试1次for(intj=0;j<=3;j++){lock=lock(key,clientID,10L);if(lock){log.info("获得锁》》》"+key);break;}try{Thread.sleep(5000);}catch(InterruptedExceptione){log.error("线程休眠异常",e);break;}}//重试3次依然没有获取到锁,那么返回服务器繁忙,请稍后重试if(!lock){thrownewException("服务繁忙");}}@OverridepublicvoidreleaseLock(Stringkey,StringclientID){if(clientID.equals(getByKey(key))){unLock(key);}}@OverridepublicbooleanhasKey(Stringkey){returntemplate.hasKey(key);}} 三、使用方法 importcom.example.redis_demo_limit.annotation.LimitedAccess;importcom.example.redis_demo_limit.redis.DataCacheRepository;importlombok.extern.slf4j.Slf4j;importorg.springframework.web.bind.annotation.PostMapping;importorg.springframework.web.bind.annotation.RequestMapping;importorg.springframework.web.bind.annotation.RestController;importjavax.annotation.Resource;importjava.util.UUID;@Slf4j@RestController@RequestMapping("/redis")publicclassRedisController{privatestaticfinalStringKEY="key";@ResourceprivateDataCacheRepositorydataCacheRepository;@LimitedAccess(frequency=1,second=1)@PostMapping("/add")publicStringadd(Stringstr){dataCacheRepository.set("str","addsuccess",200L);return"success";}//分布式锁使用示例@PostMapping("/pay")publicStringpay(StringuserName,Integeraccount){StringclientID=UUID.randomUUID().toString();//设置锁的过期时间,避免死锁Booleanlock=dataCacheRepository.lock(userName,clientID,6000L);if(!lock){log.info("未获取到锁{}",userName);return"程序繁忙,请稍后再试!";}try{//等待5s,方便测试Thread.sleep(5000);if(dataCacheRepository.hasKey(KEY)){LongaLong=dataCacheRepository.get(KEY);dataCacheRepository.set(KEY,aLong+account,-1);returnaccount+aLong+"";}else{dataCacheRepository.set(KEY,account,-1);returnaccount+"";}}catch(InterruptedExceptione){log.error(e.getMessage(),e);return"程序运行异常,请联系管理员!";}finally{if(clientID.equals(dataCacheRepository.getByKey(userName))){log.info("finally删除锁{}",userName);dataCacheRepository.unLock(userName);}}}} 四、结果测试 先在8082点击pay 再在8081点击pay 从结果来看,分布式锁起到了作用。很简单吧,还不来试一下! 本文分享自微信公众号 - 全民java空间(ljl236915692)。如有侵权,请联系 support@oschina.cn 删除。本文参与“OSC源创计划”,欢迎正在阅读的你也加入,一起分享。