From e51513230afd0ddad403d8f0ab093c3d7ee9eb59 Mon Sep 17 00:00:00 2001 From: unseok kim Date: Mon, 8 Nov 2021 09:13:37 -0500 Subject: [PATCH] GH-3655: Add automatically delete for Redis Locks Fixes https://github.com/spring-projects/spring-integration/issues/3655 * support automatically clean up cache * RedisLockRegistry.capacity desc --- .../redis/util/RedisLockRegistry.java | 45 ++++- .../redis/util/RedisLockRegistryTests.java | 180 +++++++++++++++++- src/reference/asciidoc/redis.adoc | 3 + 3 files changed, 217 insertions(+), 11 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 76ec16ccbb..d398196add 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 @@ -1,5 +1,5 @@ /* - * Copyright 2014-2020 the original author or authors. + * Copyright 2014-2021 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -19,9 +19,10 @@ package org.springframework.integration.redis.util; import java.text.SimpleDateFormat; import java.util.Collections; import java.util.Date; +import java.util.LinkedHashMap; import java.util.Map; +import java.util.Map.Entry; import java.util.UUID; -import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.Executor; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; @@ -68,6 +69,7 @@ import org.springframework.util.ReflectionUtils; * @author Konstantin Yakimov * @author Artem Bilan * @author Vedran Pavic + * @author Unseok Kim * * @since 4.0 * @@ -78,6 +80,8 @@ public final class RedisLockRegistry implements ExpirableLockRegistry, Disposabl private static final long DEFAULT_EXPIRE_AFTER = 60000L; + private static final int DEFAULT_CAPACITY = 1_000_000; + private static final String OBTAIN_LOCK_SCRIPT = "local lockClientId = redis.call('GET', KEYS[1])\n" + "if lockClientId == ARGV[1] then\n" + @@ -90,7 +94,15 @@ public final class RedisLockRegistry implements ExpirableLockRegistry, Disposabl "return false"; - private final Map locks = new ConcurrentHashMap<>(); + private final Map locks = + new LinkedHashMap(16, 0.75F, true) { + + @Override + protected boolean removeEldestEntry(Entry eldest) { + return size() > RedisLockRegistry.this.capacity; + } + + }; private final String clientId = UUID.randomUUID().toString(); @@ -102,6 +114,8 @@ public final class RedisLockRegistry implements ExpirableLockRegistry, Disposabl private final long expireAfter; + private int capacity = DEFAULT_CAPACITY; + /** * An {@link ExecutorService} to call {@link StringRedisTemplate#delete} in * the separate thread when the current one is interrupted. @@ -152,21 +166,34 @@ public final class RedisLockRegistry implements ExpirableLockRegistry, Disposabl this.executorExplicitlySet = true; } + /** + * Set the capacity of cached locks. + * @param capacity The capacity of cached lock, (default 1_000_000). + * @since 5.5.6 + */ + public void setCapacity(int capacity) { + this.capacity = capacity; + } + @Override public Lock obtain(Object lockKey) { Assert.isInstanceOf(String.class, lockKey); String path = (String) lockKey; - return this.locks.computeIfAbsent(path, RedisLock::new); + synchronized (this.locks) { + return this.locks.computeIfAbsent(path, RedisLock::new); + } } @Override public void expireUnusedOlderThan(long age) { long now = System.currentTimeMillis(); - this.locks.entrySet() - .removeIf((entry) -> { - RedisLock lock = entry.getValue(); - return now - lock.getLockedAt() > age && !lock.isAcquiredInThisProcess(); - }); + synchronized (this.locks) { + this.locks.entrySet() + .removeIf((entry) -> { + RedisLock lock = entry.getValue(); + return now - lock.getLockedAt() > age && !lock.isAcquiredInThisProcess(); + }); + } } @Override 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 947fe5f00f..6ca5cf0e53 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 @@ -1,5 +1,5 @@ /* - * Copyright 2014-2019 the original author or authors. + * Copyright 2014-2021 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -24,10 +24,13 @@ import static org.mockito.Mockito.mock; import java.util.Map; import java.util.Properties; +import java.util.Queue; import java.util.UUID; import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.Future; +import java.util.concurrent.LinkedBlockingQueue; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.locks.Lock; @@ -51,6 +54,7 @@ import org.springframework.integration.test.util.TestUtils; * @author Konstantin Yakimov * @author Artem Bilan * @author Vedran Pavic + * @author Unseok Kim * * @since 4.0 * @@ -435,9 +439,176 @@ public class RedisLockRegistryTests extends RedisAvailableTests { lock.unlock(); } + @Test + @RedisAvailable + public void concurrentObtainCapacityTest() throws InterruptedException { + final int KEY_CNT = 500; + final int CAPACITY_CNT = 179; + final int THREAD_CNT = 4; + + final CountDownLatch countDownLatch = new CountDownLatch(THREAD_CNT); + final RedisConnectionFactory connectionFactory = getConnectionFactoryForTest(); + final RedisLockRegistry registry = new RedisLockRegistry(connectionFactory, this.registryKey, 10000); + registry.setCapacity(CAPACITY_CNT); + final ExecutorService executorService = Executors.newFixedThreadPool(THREAD_CNT); + + for (int i = 0; i < KEY_CNT; i++) { + int finalI = i; + executorService.submit(() -> { + countDownLatch.countDown(); + try { + countDownLatch.await(); + } + catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + String keyId = "foo:" + finalI; + Lock obtain = registry.obtain(keyId); + obtain.lock(); + obtain.unlock(); + }); + } + executorService.shutdown(); + executorService.awaitTermination(5, TimeUnit.SECONDS); + + //capacity limit test + assertThat(TestUtils.getPropertyValue(registry, "locks", Map.class).size()).isEqualTo(CAPACITY_CNT); + + + registry.expireUnusedOlderThan(-1000); + assertThat(TestUtils.getPropertyValue(registry, "locks", Map.class).size()).isEqualTo(0); + } + + @Test + @RedisAvailable + public void concurrentObtainRemoveOrderTest() throws InterruptedException { + final int THREAD_CNT = 2; + final int DUMMY_LOCK_CNT = 3; + + final int CAPACITY_CNT = THREAD_CNT; + + final CountDownLatch countDownLatch = new CountDownLatch(THREAD_CNT); + final RedisConnectionFactory connectionFactory = getConnectionFactoryForTest(); + final RedisLockRegistry registry = new RedisLockRegistry(connectionFactory, this.registryKey, 10000); + registry.setCapacity(CAPACITY_CNT); + final ExecutorService executorService = Executors.newFixedThreadPool(THREAD_CNT); + final Queue remainLockCheckQueue = new LinkedBlockingQueue<>(); + + //Removed due to capcity limit + for (int i = 0; i < DUMMY_LOCK_CNT; i++) { + Lock obtainLock0 = registry.obtain("foo:" + i); + obtainLock0.lock(); + obtainLock0.unlock(); + } + + for (int i = DUMMY_LOCK_CNT; i < THREAD_CNT + DUMMY_LOCK_CNT; i++) { + int finalI = i; + executorService.submit(() -> { + countDownLatch.countDown(); + try { + countDownLatch.await(); + } + catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + String keyId = "foo:" + finalI; + remainLockCheckQueue.offer(keyId); + Lock obtain = registry.obtain(keyId); + obtain.lock(); + obtain.unlock(); + }); + } + + executorService.shutdown(); + executorService.awaitTermination(5, TimeUnit.SECONDS); + + assertThat(getRedisLockRegistryLocks(registry)).containsKeys( + remainLockCheckQueue.toArray(new String[remainLockCheckQueue.size()])); + } + + @Test + @RedisAvailable + public void concurrentObtainAccessRemoveOrderTest() throws InterruptedException { + final int THREAD_CNT = 2; + final int DUMMY_LOCK_CNT = 3; + + final int CAPACITY_CNT = THREAD_CNT + 1; + final String REMAIN_DUMMY_LOCK_KEY = "foo:1"; + + final CountDownLatch countDownLatch = new CountDownLatch(THREAD_CNT); + final RedisConnectionFactory connectionFactory = getConnectionFactoryForTest(); + final RedisLockRegistry registry = new RedisLockRegistry(connectionFactory, this.registryKey, 10000); + registry.setCapacity(CAPACITY_CNT); + final ExecutorService executorService = Executors.newFixedThreadPool(THREAD_CNT); + final Queue remainLockCheckQueue = new LinkedBlockingQueue<>(); + + //Removed due to capcity limit + for (int i = 0; i < DUMMY_LOCK_CNT; i++) { + Lock obtainLock0 = registry.obtain("foo:" + i); + obtainLock0.lock(); + obtainLock0.unlock(); + } + + Lock obtainLock0 = registry.obtain(REMAIN_DUMMY_LOCK_KEY); + obtainLock0.lock(); + obtainLock0.unlock(); + remainLockCheckQueue.offer(REMAIN_DUMMY_LOCK_KEY); + + for (int i = DUMMY_LOCK_CNT; i < THREAD_CNT + DUMMY_LOCK_CNT; i++) { + int finalI = i; + executorService.submit(() -> { + countDownLatch.countDown(); + try { + countDownLatch.await(); + } + catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + String keyId = "foo:" + finalI; + remainLockCheckQueue.offer(keyId); + Lock obtain = registry.obtain(keyId); + obtain.lock(); + obtain.unlock(); + }); + } + + executorService.shutdown(); + executorService.awaitTermination(5, TimeUnit.SECONDS); + + assertThat(getRedisLockRegistryLocks(registry)).containsKeys( + remainLockCheckQueue.toArray(new String[remainLockCheckQueue.size()])); + } + + @Test + @RedisAvailable + public void setCapacityTest() { + final int CAPACITY_CNT = 4; + final RedisConnectionFactory connectionFactory = getConnectionFactoryForTest(); + final RedisLockRegistry registry = new RedisLockRegistry(connectionFactory, this.registryKey, 10000); + registry.setCapacity(CAPACITY_CNT); + + registry.obtain("foo:1"); + registry.obtain("foo:2"); + registry.obtain("foo:3"); + + //capacity 4->3 + registry.setCapacity(CAPACITY_CNT - 1); + + registry.obtain("foo:4"); + + assertThat(TestUtils.getPropertyValue(registry, "locks", Map.class).size()).isEqualTo(3); + assertThat(getRedisLockRegistryLocks(registry)).containsKeys("foo:2", "foo:3", "foo:4"); + + //capacity 3->4 + registry.setCapacity(CAPACITY_CNT); + registry.obtain("foo:5"); + assertThat(TestUtils.getPropertyValue(registry, "locks", Map.class).size()).isEqualTo(4); + assertThat(getRedisLockRegistryLocks(registry)).containsKeys("foo:3", "foo:4", "foo:5"); + } + @SuppressWarnings({ "unchecked", "rawtypes" }) @Test - public void ntestUlink() { + public void testUlink() { RedisOperations ops = mock(RedisOperations.class); Properties props = new Properties(); willReturn(props).given(ops).execute(any(RedisCallback.class)); @@ -464,4 +635,9 @@ public class RedisLockRegistryTests extends RedisAvailableTests { assertThat(n < 100).as(key + " key did not expire").isTrue(); } + @SuppressWarnings("unchecked") + private Map getRedisLockRegistryLocks(RedisLockRegistry registry) { + return TestUtils.getPropertyValue(registry, "locks", Map.class); + } + } diff --git a/src/reference/asciidoc/redis.adoc b/src/reference/asciidoc/redis.adoc index 8b298ac46b..6b7bf80c4c 100644 --- a/src/reference/asciidoc/redis.adoc +++ b/src/reference/asciidoc/redis.adoc @@ -883,3 +883,6 @@ However, the resources protected by such a lock may have been compromised, so su You should set the expiry at a large enough value to prevent this condition, but set it low enough that the lock can be recovered after a server failure in a reasonable amount of time. Starting with version 5.0, the `RedisLockRegistry` implements `ExpirableLockRegistry`, which removes locks last acquired more than `age` ago and that are not currently locked. + +String with version 5.5.6, the `RedisLockRegistry` is support automatically clean up cache for redisLocks in `RedisLockRegistry.locks` via `RedisLockRegistry.setCapacity()`. +See its JavaDocs for more information. \ No newline at end of file