From ebd9b4bde2eba4ce540b64052c2d53df57cdff42 Mon Sep 17 00:00:00 2001 From: Roshi Date: Thu, 10 Jun 2021 23:45:49 +0800 Subject: [PATCH] =?UTF-8?q?=E5=A2=9E=E5=8A=A0redisRegistry=20consumer?= =?UTF-8?q?=E4=BE=A7=E7=9A=84=E6=9C=8D=E5=8A=A1=E6=8E=A2=E6=B4=BB=E9=80=BB?= =?UTF-8?q?=E8=BE=91=20(#7929)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../dubbo/registry/redis/RedisRegistry.java | 27 ++++- .../registry/redis/RedisRegistryTest.java | 104 +++++++++++++++++- 2 files changed, 128 insertions(+), 3 deletions(-) diff --git a/dubbo-registry/dubbo-registry-redis/src/main/java/org/apache/dubbo/registry/redis/RedisRegistry.java b/dubbo-registry/dubbo-registry-redis/src/main/java/org/apache/dubbo/registry/redis/RedisRegistry.java index 893b09486a..cc2a8d0018 100644 --- a/dubbo-registry/dubbo-registry-redis/src/main/java/org/apache/dubbo/registry/redis/RedisRegistry.java +++ b/dubbo-registry/dubbo-registry-redis/src/main/java/org/apache/dubbo/registry/redis/RedisRegistry.java @@ -93,6 +93,11 @@ public class RedisRegistry extends FailbackRegistry { private volatile boolean admin = false; + private final Map expireCache = new ConcurrentHashMap<>(); + + // just for unit test + private volatile boolean doExpire = true; + public RedisRegistry(URL url) { super(url); String type = url.getParameter(REDIS_CLIENT_KEY, MONO_REDIS); @@ -137,6 +142,15 @@ public class RedisRegistry extends FailbackRegistry { } } } + + if (doExpire) { + for (Map.Entry expireEntry : expireCache.entrySet()) { + if (expireEntry.getValue() < System.currentTimeMillis()) { + doNotify(toCategoryPath(expireEntry.getKey())); + } + } + } + if (admin) { clean(); } @@ -290,18 +304,28 @@ public class RedisRegistry extends FailbackRegistry { continue; } List urls = new ArrayList<>(); + Set toDeleteExpireKeys = new HashSet<>(expireCache.keySet()); Map values = redisClient.hgetAll(key); if (CollectionUtils.isNotEmptyMap(values)) { for (Map.Entry entry : values.entrySet()) { URL u = URL.valueOf(entry.getKey()); + long expire = Long.parseLong(entry.getValue()); if (!u.getParameter(DYNAMIC_KEY, true) - || Long.parseLong(entry.getValue()) >= now) { + || expire >= now) { if (UrlUtils.isMatch(url, u)) { urls.add(u); + expireCache.put(u, expire); + toDeleteExpireKeys.remove(u); } } } } + + if (!toDeleteExpireKeys.isEmpty()) { + for (URL u : toDeleteExpireKeys) { + expireCache.remove(u); + } + } if (urls.isEmpty()) { urls.add(URLBuilder.from(url) .setProtocol(EMPTY_PROTOCOL) @@ -311,6 +335,7 @@ public class RedisRegistry extends FailbackRegistry { .build()); } result.addAll(urls); + if (logger.isInfoEnabled()) { logger.info("redis notify: " + key + " = " + urls); } diff --git a/dubbo-registry/dubbo-registry-redis/src/test/java/org/apache/dubbo/registry/redis/RedisRegistryTest.java b/dubbo-registry/dubbo-registry-redis/src/test/java/org/apache/dubbo/registry/redis/RedisRegistryTest.java index 2feccc4c44..2f1795d02d 100644 --- a/dubbo-registry/dubbo-registry-redis/src/test/java/org/apache/dubbo/registry/redis/RedisRegistryTest.java +++ b/dubbo-registry/dubbo-registry-redis/src/test/java/org/apache/dubbo/registry/redis/RedisRegistryTest.java @@ -22,6 +22,7 @@ import org.apache.dubbo.registry.NotifyListener; import org.apache.dubbo.registry.Registry; import org.apache.commons.lang3.SystemUtils; +import org.apache.dubbo.registry.support.AbstractRegistry; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.BeforeEach; @@ -30,6 +31,8 @@ import redis.clients.jedis.exceptions.JedisConnectionException; import redis.embedded.RedisServer; import java.io.IOException; +import java.lang.reflect.Field; +import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Random; @@ -44,7 +47,9 @@ import static redis.embedded.RedisServer.newRedisServer; public class RedisRegistryTest { private static final String SERVICE = "org.apache.dubbo.test.injvmServie"; - private static final URL SERVICE_URL = URL.valueOf("redis://redis/" + SERVICE + "?notify=false&methods=test1,test2"); + private static final URL SERVICE_URL = URL.valueOf("redis://redis/" + SERVICE + "?notify=false&methods=test1,test2&category=providers,configurators,routers"); + private static final URL PROVIDER_URL_A = URL.valueOf("redis://127.0.0.1:20880/" + SERVICE + "?notify=false&methods=test1,test2"); + private static final URL PROVIDER_URL_B = URL.valueOf("redis://127.0.0.1:20881/" + SERVICE + "?notify=false&methods=test1,test2"); private RedisServer redisServer; private RedisRegistry redisRegistry; @@ -74,7 +79,7 @@ public class RedisRegistryTest { } } Assertions.assertNull(exception); - registryUrl = URL.valueOf("redis://localhost:" + redisPort); + registryUrl = URL.valueOf("redis://localhost:" + redisPort + "?session=4000"); redisRegistry = (RedisRegistry) new RedisRegistryFactory().createRegistry(registryUrl); } @@ -105,6 +110,101 @@ public class RedisRegistryTest { }); } + @Test + public void testSubscribeExpireCache() throws Exception { + redisRegistry.register(PROVIDER_URL_A); + redisRegistry.register(PROVIDER_URL_B); + + NotifyListener listener = new NotifyListener() { + @Override + public void notify(List urls) { + } + }; + + redisRegistry.subscribe(SERVICE_URL, listener); + + Field expireCache = RedisRegistry.class.getDeclaredField("expireCache"); + expireCache.setAccessible(true); + Map cacheExpire = (Map)expireCache.get(redisRegistry); + + assertThat(cacheExpire.get(PROVIDER_URL_A) > 0, is(true)); + assertThat(cacheExpire.get(PROVIDER_URL_B) > 0, is(true)); + + redisRegistry.unregister(PROVIDER_URL_A); + + boolean success = false; + + for (int i = 0; i < 30; i++) { + cacheExpire = (Map)expireCache.get(redisRegistry); + if (cacheExpire.get(PROVIDER_URL_A) == null) { + success = true; + break; + } + Thread.sleep(500); + } + assertThat(success, is(true)); + } + + @Test + public void testSubscribeWhenProviderCrash() throws Exception { + + // unit test will fail if doExpire=false + // Field doExpireField = RedisRegistry.class.getDeclaredField("doExpire"); + // doExpireField.setAccessible(true); + // doExpireField.set(redisRegistry, false); + + redisRegistry.register(PROVIDER_URL_A); + redisRegistry.register(PROVIDER_URL_B); + assertThat(redisRegistry.getRegistered().contains(PROVIDER_URL_A), is(true)); + assertThat(redisRegistry.getRegistered().contains(PROVIDER_URL_B), is(true)); + + Set notifiedUrls = new HashSet<>(); + Object lock = new Object(); + + NotifyListener listener = new NotifyListener() { + @Override + public void notify(List urls) { + synchronized (lock) { + notifiedUrls.clear(); + notifiedUrls.addAll(urls); + } + } + }; + + redisRegistry.subscribe(SERVICE_URL, listener); + assertThat(redisRegistry.getSubscribed().size(), is(1)); + + boolean firstOk = false; + boolean secondOk = false; + + for (int i = 0; i < 30; i++) { + synchronized (lock) { + if (notifiedUrls.contains(PROVIDER_URL_A) && notifiedUrls.contains(PROVIDER_URL_B)) { + firstOk = true; + break; + } + } + Thread.sleep(500); + } + + assertThat(firstOk, is(true)); + + // kill -9 to providerB + Field registeredField = AbstractRegistry.class.getDeclaredField("registered"); + registeredField.setAccessible(true); + ((Set) registeredField.get(redisRegistry)).remove(PROVIDER_URL_B); + + for (int i = 0; i < 30; i++) { + synchronized (lock) { + if (notifiedUrls.contains(PROVIDER_URL_A) && notifiedUrls.size() == 1) { + secondOk = true; + break; + } + } + Thread.sleep(500); + } + assertThat(secondOk, is(true)); + } @Test public void testSubscribeAndUnsubscribe() {