From cce90eaef6119c8720cb12a7d9d1da4ae86dc942 Mon Sep 17 00:00:00 2001 From: unseok kim Date: Mon, 6 Jun 2022 06:40:26 -0400 Subject: [PATCH] GH-3805: Introduce RedisLockRegistry.RedisLockType mode Fixes https://github.com/spring-projects/spring-integration/issues/3805 The Redis Pub-Sub doesn't work in all the environment, therefore there has to be a choice to use old busy-spin algorithm * Change to select between spinLock method and pub-sub method * Make spinLock as a default one to let the `RedisLockRegistry` work everywhere * Fix javadoc, convention, lazy init * Fix javadoc, convention * Code clean up and docs for `RedisLockType` feature **Cherry-pick to 5.5.x** --- .../redis/util/RedisLockRegistry.java | 393 ++++++++++++------ .../redis/util/RedisLockRegistryTests.java | 54 +++ src/reference/asciidoc/redis.adoc | 9 + 3 files changed, 332 insertions(+), 124 deletions(-) diff --git a/spring-integration-redis/src/main/java/org/springframework/integration/redis/util/RedisLockRegistry.java b/spring-integration-redis/src/main/java/org/springframework/integration/redis/util/RedisLockRegistry.java index b5836c33b2..a2b4ff95b1 100644 --- a/spring-integration-redis/src/main/java/org/springframework/integration/redis/util/RedisLockRegistry.java +++ b/spring-integration-redis/src/main/java/org/springframework/integration/redis/util/RedisLockRegistry.java @@ -34,6 +34,7 @@ import java.util.concurrent.TimeoutException; import java.util.concurrent.locks.Condition; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; +import java.util.function.Function; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; @@ -92,32 +93,6 @@ public final class RedisLockRegistry implements ExpirableLockRegistry, Disposabl private static final int DEFAULT_CAPACITY = 100_000; - private static final String OBTAIN_LOCK_SCRIPT = - "local lockClientId = redis.call('GET', KEYS[1])\n" + - "if lockClientId == ARGV[1] then\n" + - " redis.call('PEXPIRE', KEYS[1], ARGV[2])\n" + - " return true\n" + - "elseif not lockClientId then\n" + - " redis.call('SET', KEYS[1], ARGV[1], 'PX', ARGV[2])\n" + - " return true\n" + - "end\n" + - "return false"; - - private static final String UNLINK_UNLOCK_SCRIPT = - "if (redis.call('unlink', KEYS[1]) == 1) then " + - "redis.call('publish', ARGV[1], KEYS[1]) " + - "return true " + - "end " + - "return false"; - - private static final String DELETE_UNLOCK_SCRIPT = - "if (redis.call('del', KEYS[1]) == 1) then " + - "redis.call('publish', ARGV[1], KEYS[1]) " + - "return true " + - "end " + - "return false"; - - private final Map locks = new LinkedHashMap(16, 0.75F, true) { @@ -136,20 +111,12 @@ public final class RedisLockRegistry implements ExpirableLockRegistry, Disposabl private final StringRedisTemplate redisTemplate; - private final RedisScript obtainLockScript; - - private final RedisScript unLinkUnLockScript; - - private final RedisScript deleteUnLockScript; - - private final RedisUnLockNotifyMessageListener unlockNotifyMessageListener; - - private final RedisMessageListenerContainer redisMessageListenerContainer; - private final long expireAfter; private int cacheCapacity = DEFAULT_CAPACITY; + private RedisLockType redisLockType = RedisLockType.SPIN_LOCK; + /** * An {@link ExecutorService} to call {@link StringRedisTemplate#delete} in * the separate thread when the current one is interrupted. @@ -164,8 +131,19 @@ public final class RedisLockRegistry implements ExpirableLockRegistry, Disposabl private boolean executorExplicitlySet; private volatile boolean unlinkAvailable = true; + private volatile boolean isRunningRedisMessageListenerContainer = false; + /** + * It is set via lazy initialization when it is a {@link RedisLockType#PUB_SUB_LOCK}. + */ + private volatile RedisPubSubLock.RedisUnLockNotifyMessageListener unlockNotifyMessageListener; + + /** + * It is set via lazy initialization when it is a {@link RedisLockType#PUB_SUB_LOCK}. + */ + private volatile RedisMessageListenerContainer redisMessageListenerContainer; + /** * Constructs a lock registry with the default (60 second) lock expiration. * @param connectionFactory The connection factory. @@ -185,18 +163,18 @@ public final class RedisLockRegistry implements ExpirableLockRegistry, Disposabl Assert.notNull(connectionFactory, "'connectionFactory' cannot be null"); Assert.notNull(registryKey, "'registryKey' cannot be null"); this.redisTemplate = new StringRedisTemplate(connectionFactory); - this.obtainLockScript = new DefaultRedisScript<>(OBTAIN_LOCK_SCRIPT, Boolean.class); - this.unLinkUnLockScript = new DefaultRedisScript<>(UNLINK_UNLOCK_SCRIPT, Boolean.class); - this.deleteUnLockScript = new DefaultRedisScript<>(DELETE_UNLOCK_SCRIPT, Boolean.class); this.registryKey = registryKey; this.expireAfter = expireAfter; this.unLockChannelKey = registryKey + "-channel"; - this.unlockNotifyMessageListener = new RedisUnLockNotifyMessageListener(); - this.redisMessageListenerContainer = new RedisMessageListenerContainer(); - setupUnlockMessageListener(connectionFactory); } private void setupUnlockMessageListener(RedisConnectionFactory connectionFactory) { + Assert.isNull(RedisLockRegistry.this.redisMessageListenerContainer, + "'redisMessageListenerContainer' must not have been re-initialized."); + Assert.isNull(RedisLockRegistry.this.unlockNotifyMessageListener, + "'unlockNotifyMessageListener' must not have been re-initialized."); + RedisLockRegistry.this.redisMessageListenerContainer = new RedisMessageListenerContainer(); + RedisLockRegistry.this.unlockNotifyMessageListener = new RedisPubSubLock.RedisUnLockNotifyMessageListener(); final Topic topic = new ChannelTopic(this.unLockChannelKey); this.redisMessageListenerContainer.setConnectionFactory(connectionFactory); this.redisMessageListenerContainer.setTaskExecutor(this.executor); @@ -226,12 +204,27 @@ public final class RedisLockRegistry implements ExpirableLockRegistry, Disposabl this.cacheCapacity = cacheCapacity; } + + /** + * Set {@link RedisLockType} mode to work in. + * By default, the {@link RedisLockType#SPIN_LOCK} is used - works in all the environment. + * The {@link RedisLockType#PUB_SUB_LOCK} is a preferred mode when not in Master/Replica connections - + * less network chatter. + * Set the type of unlockType, Select the lock method. + * @param redisLockType the {@link RedisLockType} to work in. + * @since 5.5.13 + */ + public void setRedisLockType(RedisLockType redisLockType) { + Assert.notNull(redisLockType, "'redisLockType' cannot be null"); + this.redisLockType = redisLockType; + } + @Override public Lock obtain(Object lockKey) { Assert.isInstanceOf(String.class, lockKey); String path = (String) lockKey; synchronized (this.locks) { - return this.locks.computeIfAbsent(path, RedisLock::new); + return this.locks.computeIfAbsent(path, getRedisLockConstructor(this.redisLockType)); } } @@ -260,9 +253,50 @@ public final class RedisLockRegistry implements ExpirableLockRegistry, Disposabl } } - private final class RedisLock implements Lock { + /** + * The mode in which this registry is going to work with locks. + */ + public enum RedisLockType { - private final String lockKey; + /** + * The lock is acquired by periodically(100ms) checking whether the lock can be acquired. + */ + SPIN_LOCK, + + /** + * The lock is acquired by redis pub-sub subscription. + */ + PUB_SUB_LOCK + } + + private Function getRedisLockConstructor(RedisLockType redisLockType) { + switch (redisLockType) { + case SPIN_LOCK: + return RedisSpinLock::new; + case PUB_SUB_LOCK: + return RedisPubSubLock::new; + default: + throw new IllegalArgumentException(); + } + } + + private abstract class RedisLock implements Lock { + + private static final String OBTAIN_LOCK_SCRIPT = + "local lockClientId = redis.call('GET', KEYS[1]) " + + "if lockClientId == ARGV[1] then " + + " redis.call('PEXPIRE', KEYS[1], ARGV[2]) " + + " return true " + + "elseif not lockClientId then " + + " redis.call('SET', KEYS[1], ARGV[1], 'PX', ARGV[2]) " + + " return true " + + "end " + + "return false"; + + protected static final RedisScript + OBTAIN_LOCK_REDIS_SCRIPT = new DefaultRedisScript<>(OBTAIN_LOCK_SCRIPT, Boolean.class); + + protected final String lockKey; private final ReentrantLock localLock = new ReentrantLock(); @@ -280,12 +314,31 @@ public final class RedisLockRegistry implements ExpirableLockRegistry, Disposabl return this.lockedAt; } + /** + * Attempt to acquire a lock in redis. + * @param time the maximum time(milliseconds) to wait for the lock, -1 infinity + * @return true if the lock was acquired and false if the waiting time elapsed before the lock was acquired + * @throws InterruptedException – + * if the current thread is interrupted while acquiring the lock (and interruption of lock acquisition is supported) + */ + protected abstract boolean tryRedisLockInner(long time) throws ExecutionException, InterruptedException; + + /** + * Unlock the lock using the unlink method in redis. + */ + protected abstract void removeLockKeyInnerUnlink(); + + /** + * Unlock the lock using the delete method in redis. + */ + protected abstract void removeLockKeyInnerDelete(); + @Override - public void lock() { + public final void lock() { this.localLock.lock(); while (true) { try { - if (subscribeLock(-1L)) { + if (tryRedisLock(-1L)) { return; } } @@ -308,11 +361,11 @@ public final class RedisLockRegistry implements ExpirableLockRegistry, Disposabl } @Override - public void lockInterruptibly() throws InterruptedException { + public final void lockInterruptibly() throws InterruptedException { this.localLock.lockInterruptibly(); while (true) { try { - if (subscribeLock(-1L)) { + if (tryRedisLock(-1L)) { return; } } @@ -329,7 +382,7 @@ public final class RedisLockRegistry implements ExpirableLockRegistry, Disposabl } @Override - public boolean tryLock() { + public final boolean tryLock() { try { return tryLock(0, TimeUnit.MILLISECONDS); } @@ -340,13 +393,13 @@ public final class RedisLockRegistry implements ExpirableLockRegistry, Disposabl } @Override - public boolean tryLock(long time, TimeUnit unit) throws InterruptedException { + public final boolean tryLock(long time, TimeUnit unit) throws InterruptedException { if (!this.localLock.tryLock(time, unit)) { return false; } try { long waitTime = TimeUnit.MILLISECONDS.convert(time, unit); - boolean acquired = subscribeLock(waitTime); + boolean acquired = tryRedisLock(waitTime); if (!acquired) { this.localLock.unlock(); } @@ -359,58 +412,23 @@ public final class RedisLockRegistry implements ExpirableLockRegistry, Disposabl return false; } - private boolean subscribeLock(long time) throws ExecutionException, InterruptedException { - final long expiredTime = System.currentTimeMillis() + time; - if (obtainLock()) { - return true; - } - - if (!(RedisLockRegistry.this.isRunningRedisMessageListenerContainer - && RedisLockRegistry.this.redisMessageListenerContainer.isRunning())) { - runRedisMessageListenerContainer(); - } - while (time == -1 || expiredTime >= System.currentTimeMillis()) { - try { - Future future = - RedisLockRegistry.this.unlockNotifyMessageListener.subscribeLock(this.lockKey); - //DCL - if (obtainLock()) { - return true; - } - try { - //if short expireAfter key expire for ttl, no receive unlock msg - long waitTime = time >= 0 ? time : RedisLockRegistry.this.expireAfter; - future.get(waitTime, TimeUnit.MILLISECONDS); - } - catch (TimeoutException ignore) { - } - if (obtainLock()) { - return true; - } - } - finally { - RedisLockRegistry.this.unlockNotifyMessageListener.unSubscribeLock(this.lockKey); - } - } - return false; - } - - private boolean obtainLock() { - Boolean success = - RedisLockRegistry.this.redisTemplate.execute(RedisLockRegistry.this.obtainLockScript, - Collections.singletonList(this.lockKey), RedisLockRegistry.this.clientId, - String.valueOf(RedisLockRegistry.this.expireAfter)); - - boolean result = Boolean.TRUE.equals(success); - + private boolean tryRedisLock(long time) throws ExecutionException, InterruptedException { + final boolean result = tryRedisLockInner(time); if (result) { this.lockedAt = System.currentTimeMillis(); } return result; } + protected final Boolean obtainLock() { + return RedisLockRegistry.this.redisTemplate + .execute(OBTAIN_LOCK_REDIS_SCRIPT, Collections.singletonList(this.lockKey), + RedisLockRegistry.this.clientId, + String.valueOf(RedisLockRegistry.this.expireAfter)); + } + @Override - public void unlock() { + public final void unlock() { if (!this.localLock.isHeldByCurrentThread()) { throw new IllegalStateException("You do not own lock at " + this.lockKey); } @@ -446,9 +464,7 @@ public final class RedisLockRegistry implements ExpirableLockRegistry, Disposabl private void removeLockKey() { if (RedisLockRegistry.this.unlinkAvailable) { try { - RedisLockRegistry.this.redisTemplate.execute( - RedisLockRegistry.this.unLinkUnLockScript, Collections.singletonList(this.lockKey), - RedisLockRegistry.this.unLockChannelKey); + removeLockKeyInnerUnlink(); return; } catch (Exception ex) { @@ -463,24 +479,22 @@ public final class RedisLockRegistry implements ExpirableLockRegistry, Disposabl } } } - - RedisLockRegistry.this.redisTemplate.execute( - RedisLockRegistry.this.deleteUnLockScript, Collections.singletonList(this.lockKey), - RedisLockRegistry.this.unLockChannelKey); + removeLockKeyInnerDelete(); } @Override - public Condition newCondition() { + public final Condition newCondition() { throw new UnsupportedOperationException("Conditions are not supported"); } - public boolean isAcquiredInThisProcess() { + public final boolean isAcquiredInThisProcess() { return RedisLockRegistry.this.clientId.equals( RedisLockRegistry.this.redisTemplate.boundValueOps(this.lockKey).get()); } @Override public String toString() { + SimpleDateFormat dateFormat = new SimpleDateFormat("yyyy-MM-dd@HH:mm:ss.SSS"); return "RedisLock [lockKey=" + this.lockKey + ",lockedAt=" + dateFormat.format(new Date(this.lockedAt)) @@ -524,41 +538,172 @@ public final class RedisLockRegistry implements ExpirableLockRegistry, Disposabl return RedisLockRegistry.this; } + } + + private final class RedisPubSubLock extends RedisLock { + + private static final String UNLINK_UNLOCK_SCRIPT = + "if (redis.call('unlink', KEYS[1]) == 1) then " + + "redis.call('publish', ARGV[1], KEYS[1]) " + + "return true " + + "end " + + "return false"; + + private static final String DELETE_UNLOCK_SCRIPT = + "if (redis.call('del', KEYS[1]) == 1) then " + + "redis.call('publish', ARGV[1], KEYS[1]) " + + "return true " + + "end " + + "return false"; + + private static final RedisScript + UNLINK_UNLOCK_REDIS_SCRIPT = new DefaultRedisScript<>(UNLINK_UNLOCK_SCRIPT, Boolean.class); + + private static final RedisScript + DELETE_UNLOCK_REDIS_SCRIPT = new DefaultRedisScript<>(DELETE_UNLOCK_SCRIPT, Boolean.class); + + private RedisPubSubLock(String path) { + super(path); + } + + @Override + protected boolean tryRedisLockInner(long time) throws ExecutionException, InterruptedException { + return subscribeLock(time); + } + + @Override + protected void removeLockKeyInnerUnlink() { + RedisLockRegistry.this.redisTemplate.execute( + UNLINK_UNLOCK_REDIS_SCRIPT, Collections.singletonList(this.lockKey), + RedisLockRegistry.this.unLockChannelKey); + } + + @Override + protected void removeLockKeyInnerDelete() { + RedisLockRegistry.this.redisTemplate.execute( + DELETE_UNLOCK_REDIS_SCRIPT, Collections.singletonList(this.lockKey), + RedisLockRegistry.this.unLockChannelKey); + + } + + private boolean subscribeLock(long time) throws ExecutionException, InterruptedException { + final long expiredTime = System.currentTimeMillis() + time; + if (obtainLock()) { + return true; + } + + if (!(RedisLockRegistry.this.isRunningRedisMessageListenerContainer + && RedisLockRegistry.this.redisMessageListenerContainer != null + && RedisLockRegistry.this.redisMessageListenerContainer.isRunning())) { + + runRedisMessageListenerContainer(); + } + while (time == -1 || expiredTime >= System.currentTimeMillis()) { + try { + Future future = + RedisLockRegistry.this.unlockNotifyMessageListener.subscribeLock(this.lockKey); + //DCL + if (obtainLock()) { + return true; + } + try { + //if short expireAfter key expire for ttl, no receive unlock msg + long waitTime = time >= 0 ? time : RedisLockRegistry.this.expireAfter; + future.get(waitTime, TimeUnit.MILLISECONDS); + } + catch (TimeoutException ignore) { + } + if (obtainLock()) { + return true; + } + } + finally { + RedisLockRegistry.this.unlockNotifyMessageListener.unSubscribeLock(this.lockKey); + } + } + return false; + } + private void runRedisMessageListenerContainer() { - synchronized (RedisLockRegistry.this.redisMessageListenerContainer) { + synchronized (RedisLockRegistry.this.locks) { if (!(RedisLockRegistry.this.isRunningRedisMessageListenerContainer + && RedisLockRegistry.this.redisMessageListenerContainer != null && RedisLockRegistry.this.redisMessageListenerContainer.isRunning())) { - RedisLockRegistry.this.redisMessageListenerContainer.afterPropertiesSet(); + + if (RedisLockRegistry.this.redisMessageListenerContainer == null) { + setupUnlockMessageListener(RedisLockRegistry.this.redisTemplate.getConnectionFactory()); + RedisLockRegistry.this.redisMessageListenerContainer.afterPropertiesSet(); + } + RedisLockRegistry.this.redisMessageListenerContainer.start(); RedisLockRegistry.this.isRunningRedisMessageListenerContainer = true; } } } + + private static final class RedisUnLockNotifyMessageListener implements MessageListener { + + private final Map> notifyMap = new ConcurrentHashMap<>(); + + @Override + public void onMessage(Message message, byte[] pattern) { + final String lockKey = new String(message.getBody()); + unlockNotify(lockKey); + } + + public Future subscribeLock(String lockKey) { + return this.notifyMap.computeIfAbsent(lockKey, key -> new SettableListenableFuture<>()); + } + + public void unSubscribeLock(String localLock) { + this.notifyMap.remove(localLock); + } + + private void unlockNotify(String lockKey) { + this.notifyMap.computeIfPresent(lockKey, (key, lockFuture) -> { + lockFuture.set(key); + return lockFuture; + }); + } + + } + } - private static final class RedisUnLockNotifyMessageListener implements MessageListener { - private final Map> notifyMap = new ConcurrentHashMap<>(); + private final class RedisSpinLock extends RedisLock { + + private RedisSpinLock(String path) { + super(path); + } @Override - public void onMessage(Message message, byte[] pattern) { - final String lockKey = new String(message.getBody()); - unlockNotify(lockKey); + protected boolean tryRedisLockInner(long time) throws InterruptedException { + long now = System.currentTimeMillis(); + if (time == -1L) { + while (!obtainLock()) { + Thread.sleep(100); //NOSONAR + } + return true; + } + else { + long expire = now + TimeUnit.MILLISECONDS.convert(time, TimeUnit.MILLISECONDS); + boolean acquired; + while (!(acquired = obtainLock()) && System.currentTimeMillis() < expire) { //NOSONAR + Thread.sleep(100); //NOSONAR + } + return acquired; + } } - public Future subscribeLock(String lockKey) { - return this.notifyMap.computeIfAbsent(lockKey, key -> new SettableListenableFuture<>()); + @Override + protected void removeLockKeyInnerUnlink() { + RedisLockRegistry.this.redisTemplate.unlink(this.lockKey); } - public void unSubscribeLock(String localLock) { - this.notifyMap.remove(localLock); - } - - private void unlockNotify(String lockKey) { - this.notifyMap.computeIfPresent(lockKey, (key, lockFuture) -> { - lockFuture.set(key); - return lockFuture; - }); + @Override + protected void removeLockKeyInnerDelete() { + RedisLockRegistry.this.redisTemplate.delete(this.lockKey); } } diff --git a/spring-integration-redis/src/test/java/org/springframework/integration/redis/util/RedisLockRegistryTests.java b/spring-integration-redis/src/test/java/org/springframework/integration/redis/util/RedisLockRegistryTests.java index c992c9a42d..713121d2fd 100644 --- a/spring-integration-redis/src/test/java/org/springframework/integration/redis/util/RedisLockRegistryTests.java +++ b/spring-integration-redis/src/test/java/org/springframework/integration/redis/util/RedisLockRegistryTests.java @@ -22,6 +22,7 @@ import static org.mockito.ArgumentMatchers.any; import static org.mockito.BDDMockito.willReturn; import static org.mockito.Mockito.mock; +import java.util.Collection; import java.util.List; import java.util.Map; import java.util.Properties; @@ -47,6 +48,9 @@ import org.apache.commons.logging.LogFactory; import org.junit.After; import org.junit.Before; import org.junit.Test; +import org.junit.runner.RunWith; +import org.junit.runners.Parameterized; +import org.junit.runners.Parameterized.Parameters; import org.springframework.data.redis.connection.RedisConnectionFactory; import org.springframework.data.redis.core.RedisCallback; @@ -54,6 +58,7 @@ import org.springframework.data.redis.core.RedisOperations; import org.springframework.data.redis.core.StringRedisTemplate; import org.springframework.integration.redis.rules.RedisAvailable; import org.springframework.integration.redis.rules.RedisAvailableTests; +import org.springframework.integration.redis.util.RedisLockRegistry.RedisLockType; import org.springframework.integration.test.util.TestUtils; /** @@ -66,14 +71,26 @@ import org.springframework.integration.test.util.TestUtils; * @since 4.0 * */ +@RunWith(Parameterized.class) public class RedisLockRegistryTests extends RedisAvailableTests { + private final RedisLockType testRedisLockType; + + public RedisLockRegistryTests(RedisLockType redisLockType) { + this.testRedisLockType = redisLockType; + } + private final Log logger = LogFactory.getLog(getClass()); private final String registryKey = UUID.randomUUID().toString(); private final String registryKey2 = UUID.randomUUID().toString(); + @Parameters + public static Collection getRedisLockTypeParameters() { + return List.of(RedisLockType.values()); + } + @Before @After public void setupShutDown() { @@ -90,6 +107,7 @@ public class RedisLockRegistryTests extends RedisAvailableTests { @RedisAvailable public void testLock() { RedisLockRegistry registry = new RedisLockRegistry(getConnectionFactoryForTest(), this.registryKey); + registry.setRedisLockType(testRedisLockType); for (int i = 0; i < 10; i++) { Lock lock = registry.obtain("foo"); lock.lock(); @@ -108,6 +126,7 @@ public class RedisLockRegistryTests extends RedisAvailableTests { @RedisAvailable public void testLockInterruptibly() throws Exception { RedisLockRegistry registry = new RedisLockRegistry(getConnectionFactoryForTest(), this.registryKey); + registry.setRedisLockType(testRedisLockType); for (int i = 0; i < 10; i++) { Lock lock = registry.obtain("foo"); lock.lockInterruptibly(); @@ -126,6 +145,7 @@ public class RedisLockRegistryTests extends RedisAvailableTests { @RedisAvailable public void testReentrantLock() { RedisLockRegistry registry = new RedisLockRegistry(getConnectionFactoryForTest(), this.registryKey); + registry.setRedisLockType(testRedisLockType); for (int i = 0; i < 10; i++) { Lock lock1 = registry.obtain("foo"); lock1.lock(); @@ -152,6 +172,7 @@ public class RedisLockRegistryTests extends RedisAvailableTests { @RedisAvailable public void testReentrantLockInterruptibly() throws Exception { RedisLockRegistry registry = new RedisLockRegistry(getConnectionFactoryForTest(), this.registryKey); + registry.setRedisLockType(testRedisLockType); for (int i = 0; i < 10; i++) { Lock lock1 = registry.obtain("foo"); lock1.lockInterruptibly(); @@ -178,6 +199,7 @@ public class RedisLockRegistryTests extends RedisAvailableTests { @RedisAvailable public void testTwoLocks() throws Exception { RedisLockRegistry registry = new RedisLockRegistry(getConnectionFactoryForTest(), this.registryKey); + registry.setRedisLockType(testRedisLockType); for (int i = 0; i < 10; i++) { Lock lock1 = registry.obtain("foo"); lock1.lockInterruptibly(); @@ -204,6 +226,7 @@ public class RedisLockRegistryTests extends RedisAvailableTests { @RedisAvailable public void testTwoThreadsSecondFailsToGetLock() throws Exception { final RedisLockRegistry registry = new RedisLockRegistry(getConnectionFactoryForTest(), this.registryKey); + registry.setRedisLockType(testRedisLockType); final Lock lock1 = registry.obtain("foo"); lock1.lockInterruptibly(); final AtomicBoolean locked = new AtomicBoolean(); @@ -234,6 +257,7 @@ public class RedisLockRegistryTests extends RedisAvailableTests { @RedisAvailable public void testTwoThreads() throws Exception { RedisLockRegistry registry = new RedisLockRegistry(getConnectionFactoryForTest(), this.registryKey); + registry.setRedisLockType(testRedisLockType); Lock lock1 = registry.obtain("foo"); AtomicBoolean locked = new AtomicBoolean(); CountDownLatch latch1 = new CountDownLatch(1); @@ -272,7 +296,9 @@ public class RedisLockRegistryTests extends RedisAvailableTests { @RedisAvailable public void testTwoThreadsDifferentRegistries() throws Exception { RedisLockRegistry registry1 = new RedisLockRegistry(getConnectionFactoryForTest(), this.registryKey); + registry1.setRedisLockType(testRedisLockType); RedisLockRegistry registry2 = new RedisLockRegistry(getConnectionFactoryForTest(), this.registryKey); + registry2.setRedisLockType(testRedisLockType); Lock lock1 = registry1.obtain("foo"); AtomicBoolean locked = new AtomicBoolean(); CountDownLatch latch1 = new CountDownLatch(1); @@ -319,6 +345,7 @@ public class RedisLockRegistryTests extends RedisAvailableTests { @RedisAvailable public void testTwoThreadsWrongOneUnlocks() throws Exception { RedisLockRegistry registry = new RedisLockRegistry(getConnectionFactoryForTest(), this.registryKey); + registry.setRedisLockType(testRedisLockType); Lock lock = registry.obtain("foo"); lock.lockInterruptibly(); AtomicBoolean locked = new AtomicBoolean(); @@ -347,7 +374,9 @@ public class RedisLockRegistryTests extends RedisAvailableTests { @RedisAvailable public void testExpireTwoRegistries() throws Exception { RedisLockRegistry registry1 = new RedisLockRegistry(getConnectionFactoryForTest(), this.registryKey, 100); + registry1.setRedisLockType(testRedisLockType); RedisLockRegistry registry2 = new RedisLockRegistry(getConnectionFactoryForTest(), this.registryKey, 100); + registry2.setRedisLockType(testRedisLockType); Lock lock1 = registry1.obtain("foo"); Lock lock2 = registry2.obtain("foo"); assertThat(lock1.tryLock()).isTrue(); @@ -361,6 +390,7 @@ public class RedisLockRegistryTests extends RedisAvailableTests { @RedisAvailable public void testExceptionOnExpire() throws Exception { RedisLockRegistry registry = new RedisLockRegistry(getConnectionFactoryForTest(), this.registryKey, 1); + registry.setRedisLockType(testRedisLockType); Lock lock1 = registry.obtain("foo"); assertThat(lock1.tryLock()).isTrue(); waitForExpire("foo"); @@ -375,8 +405,12 @@ public class RedisLockRegistryTests extends RedisAvailableTests { public void testEquals() { RedisConnectionFactory connectionFactory = getConnectionFactoryForTest(); RedisLockRegistry registry1 = new RedisLockRegistry(connectionFactory, this.registryKey); + registry1.setRedisLockType(testRedisLockType); RedisLockRegistry registry2 = new RedisLockRegistry(connectionFactory, this.registryKey); + registry2.setRedisLockType(testRedisLockType); RedisLockRegistry registry3 = new RedisLockRegistry(connectionFactory, this.registryKey2); + registry3.setRedisLockType(testRedisLockType); + Lock lock1 = registry1.obtain("foo"); Lock lock2 = registry1.obtain("foo"); assertThat(lock2).isEqualTo(lock1); @@ -407,6 +441,7 @@ public class RedisLockRegistryTests extends RedisAvailableTests { @RedisAvailable public void testThreadLocalListLeaks() { RedisLockRegistry registry = new RedisLockRegistry(getConnectionFactoryForTest(), this.registryKey, 10000); + registry.setRedisLockType(testRedisLockType); for (int i = 0; i < 10; i++) { registry.obtain("foo" + i); @@ -431,6 +466,8 @@ public class RedisLockRegistryTests extends RedisAvailableTests { public void testExpireNotChanged() throws Exception { RedisConnectionFactory connectionFactory = getConnectionFactoryForTest(); final RedisLockRegistry registry = new RedisLockRegistry(connectionFactory, this.registryKey, 10000); + registry.setRedisLockType(testRedisLockType); + Lock lock = registry.obtain("foo"); lock.lock(); @@ -457,6 +494,8 @@ public class RedisLockRegistryTests extends RedisAvailableTests { final RedisConnectionFactory connectionFactory = getConnectionFactoryForTest(); final RedisLockRegistry registry = new RedisLockRegistry(connectionFactory, this.registryKey, 10000); registry.setCacheCapacity(CAPACITY_CNT); + registry.setRedisLockType(testRedisLockType); + final ExecutorService executorService = Executors.newFixedThreadPool(THREAD_CNT); for (int i = 0; i < KEY_CNT; i++) { @@ -498,6 +537,8 @@ public class RedisLockRegistryTests extends RedisAvailableTests { final RedisConnectionFactory connectionFactory = getConnectionFactoryForTest(); final RedisLockRegistry registry = new RedisLockRegistry(connectionFactory, this.registryKey, 10000); registry.setCacheCapacity(CAPACITY_CNT); + registry.setRedisLockType(testRedisLockType); + final ExecutorService executorService = Executors.newFixedThreadPool(THREAD_CNT); final Queue remainLockCheckQueue = new LinkedBlockingQueue<>(); @@ -546,6 +587,8 @@ public class RedisLockRegistryTests extends RedisAvailableTests { final RedisConnectionFactory connectionFactory = getConnectionFactoryForTest(); final RedisLockRegistry registry = new RedisLockRegistry(connectionFactory, this.registryKey, 10000); registry.setCacheCapacity(CAPACITY_CNT); + registry.setRedisLockType(testRedisLockType); + final ExecutorService executorService = Executors.newFixedThreadPool(THREAD_CNT); final Queue remainLockCheckQueue = new LinkedBlockingQueue<>(); @@ -593,6 +636,7 @@ public class RedisLockRegistryTests extends RedisAvailableTests { final RedisConnectionFactory connectionFactory = getConnectionFactoryForTest(); final RedisLockRegistry registry = new RedisLockRegistry(connectionFactory, this.registryKey, 10000); registry.setCacheCapacity(CAPACITY_CNT); + registry.setRedisLockType(testRedisLockType); registry.obtain("foo:1"); registry.obtain("foo:2"); @@ -618,7 +662,10 @@ public class RedisLockRegistryTests extends RedisAvailableTests { public void twoRedisLockRegistryTest() throws InterruptedException { RedisConnectionFactory connectionFactory = getConnectionFactoryForTest(); RedisLockRegistry registry1 = new RedisLockRegistry(connectionFactory, registryKey, 1000000L); + registry1.setRedisLockType(testRedisLockType); RedisLockRegistry registry2 = new RedisLockRegistry(connectionFactory, registryKey, 1000000L); + registry2.setRedisLockType(testRedisLockType); + String lockKey = "test-1"; Lock obtainLock_1 = registry1.obtain(lockKey); @@ -669,6 +716,7 @@ public class RedisLockRegistryTests extends RedisAvailableTests { .mapToObj((num) -> new RedisLockRegistry( connectionFactory, registryKey, expireAfter)) .map((registry) -> { + registry.setRedisLockType(testRedisLockType); final Callable callable = () -> { Lock obtain = registry.obtain(testKey); obtain.lock(); @@ -702,8 +750,12 @@ public class RedisLockRegistryTests extends RedisAvailableTests { final CountDownLatch awaitTimeout = new CountDownLatch(THREAD_CNT); final RedisConnectionFactory connectionFactory = getConnectionFactoryForTest(); final RedisLockRegistry registry1 = new RedisLockRegistry(connectionFactory, this.registryKey); + registry1.setRedisLockType(testRedisLockType); final RedisLockRegistry registry2 = new RedisLockRegistry(connectionFactory, this.registryKey); + registry2.setRedisLockType(testRedisLockType); final RedisLockRegistry registry3 = new RedisLockRegistry(connectionFactory, this.registryKey); + registry3.setRedisLockType(testRedisLockType); + final ExecutorService executorService = Executors.newFixedThreadPool(THREAD_CNT); Lock lock1 = registry1.obtain(testKey); @@ -754,9 +806,11 @@ public class RedisLockRegistryTests extends RedisAvailableTests { willReturn(props).given(ops).execute(any(RedisCallback.class)); props.setProperty("redis_version", "3.0.0"); RedisLockRegistry registry = new RedisLockRegistry(mock(RedisConnectionFactory.class), "foo"); + registry.setRedisLockType(testRedisLockType); assertThat(TestUtils.getPropertyValue(registry, "ulinkAvailable", Boolean.class)).isFalse(); props.setProperty("redis_version", "4.0.0"); registry = new RedisLockRegistry(mock(RedisConnectionFactory.class), "foo"); + registry.setRedisLockType(testRedisLockType); assertThat(TestUtils.getPropertyValue(registry, "ulinkAvailable", Boolean.class)).isTrue(); } diff --git a/src/reference/asciidoc/redis.adoc b/src/reference/asciidoc/redis.adoc index b0a008f99b..40d5ea5c97 100644 --- a/src/reference/asciidoc/redis.adoc +++ b/src/reference/asciidoc/redis.adoc @@ -887,3 +887,12 @@ Starting with version 5.0, the `RedisLockRegistry` implements `ExpirableLockRegi String with version 5.5.6, the `RedisLockRegistry` is support automatically clean up cache for redisLocks in `RedisLockRegistry.locks` via `RedisLockRegistry.setCacheCapacity()`. See its JavaDocs for more information. + +String with version 5.5.13, the `RedisLockRegistry` exposes a `setRedisLockType(RedisLockType)` option to determine in which mode a Redis lock acquisition should happen: + +- `RedisLockType.SPIN_LOCK` - the lock is acquired by periodic loop (100ms) checking whether the lock can be acquired. +Default. +- `RedisLockType.PUB_SUB_LOCK` - The lock is acquired by redis pub-sub subscription. + +The pub-sub is preferred mode - less network chatter between client Redis server, and more performant - the lock is acquired immediately when subscription is notified about unlocking in the other process. +However, the Redis does not support pub-sub in the Master/Replica connections (for example in AWS ElastiCache environment), therefore a busy-spin mode is chosen as a default to make the registry working in any environment. \ No newline at end of file