From 188a90e498b82f363c45f411eb93f84ff7a4f727 Mon Sep 17 00:00:00 2001 From: Costin Leau Date: Thu, 17 Mar 2011 18:18:54 +0200 Subject: [PATCH] DATAKV-49 DATAKV-46 + finish up the pubsub improvements and fixed last remaining bugs + added blocking behaviour to RJC pubsub --- .../redis/connection/rjc/RjcSubscription.java | 8 +++++++- .../redis/connection/util/AbstractSubscription.java | 12 ++++++++++-- .../keyvalue/redis/ConnectionFactoryTracker.java | 2 +- .../AbstractConnectionIntegrationTests.java | 3 ++- .../rjc/RjcConnectionIntegrationTests.java | 2 +- 5 files changed, 21 insertions(+), 6 deletions(-) diff --git a/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/rjc/RjcSubscription.java b/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/rjc/RjcSubscription.java index fa0aa5a1d..0cd0d48cf 100644 --- a/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/rjc/RjcSubscription.java +++ b/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/rjc/RjcSubscription.java @@ -39,7 +39,13 @@ class RjcSubscription extends AbstractSubscription { @Override protected void doClose() { - subscriber.close(); + try { + subscriber.close(); + } finally { + synchronized (pubSubMonitor) { + pubSubMonitor.notifyAll(); + } + } } @Override diff --git a/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/util/AbstractSubscription.java b/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/util/AbstractSubscription.java index 76dfbfec1..6d20a8bf5 100644 --- a/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/util/AbstractSubscription.java +++ b/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/util/AbstractSubscription.java @@ -56,10 +56,10 @@ public abstract class AbstractSubscription implements Subscription { this.listener = listener; synchronized (this.channels) { - remove(this.channels, channels); + add(this.channels, channels); } synchronized (this.patterns) { - remove(this.patterns, patterns); + add(this.patterns, patterns); } } @@ -168,6 +168,10 @@ public abstract class AbstractSubscription implements Subscription { this.patterns.clear(); } } + else { + // nothing to unsubscribe from + return; + } } else { synchronized (this.patterns) { @@ -194,6 +198,10 @@ public abstract class AbstractSubscription implements Subscription { this.channels.clear(); } } + else { + // nothing to unsubscribe from + return; + } } else { synchronized (this.channels) { diff --git a/spring-data-redis/src/test/java/org/springframework/data/keyvalue/redis/ConnectionFactoryTracker.java b/spring-data-redis/src/test/java/org/springframework/data/keyvalue/redis/ConnectionFactoryTracker.java index 9ef6c5e59..634d3448a 100644 --- a/spring-data-redis/src/test/java/org/springframework/data/keyvalue/redis/ConnectionFactoryTracker.java +++ b/spring-data-redis/src/test/java/org/springframework/data/keyvalue/redis/ConnectionFactoryTracker.java @@ -40,7 +40,7 @@ public abstract class ConnectionFactoryTracker { for (RedisConnectionFactory connectionFactory : connFactories) { try { ((DisposableBean) connectionFactory).destroy(); - System.out.println("Succesfully cleaned up factory " + connectionFactory); + //System.out.println("Succesfully cleaned up factory " + connectionFactory); } catch (Exception ex) { System.err.println("Cannot clean factory " + connectionFactory + ex); } diff --git a/spring-data-redis/src/test/java/org/springframework/data/keyvalue/redis/connection/AbstractConnectionIntegrationTests.java b/spring-data-redis/src/test/java/org/springframework/data/keyvalue/redis/connection/AbstractConnectionIntegrationTests.java index cb8c408ce..116908460 100644 --- a/spring-data-redis/src/test/java/org/springframework/data/keyvalue/redis/connection/AbstractConnectionIntegrationTests.java +++ b/spring-data-redis/src/test/java/org/springframework/data/keyvalue/redis/connection/AbstractConnectionIntegrationTests.java @@ -218,7 +218,7 @@ public abstract class AbstractConnectionIntegrationTests { System.out.println("Subscribed"); while (flag.get()) { try { - Thread.currentThread().wait(2000); + Thread.currentThread().sleep(2000); } catch (Exception ex) { return; } @@ -239,6 +239,7 @@ public abstract class AbstractConnectionIntegrationTests { } finally { flag.set(false); } + System.out.println(queue); assertEquals(3, queue.size()); } diff --git a/spring-data-redis/src/test/java/org/springframework/data/keyvalue/redis/connection/rjc/RjcConnectionIntegrationTests.java b/spring-data-redis/src/test/java/org/springframework/data/keyvalue/redis/connection/rjc/RjcConnectionIntegrationTests.java index 8bbe97356..4fa5b3ed8 100644 --- a/spring-data-redis/src/test/java/org/springframework/data/keyvalue/redis/connection/rjc/RjcConnectionIntegrationTests.java +++ b/spring-data-redis/src/test/java/org/springframework/data/keyvalue/redis/connection/rjc/RjcConnectionIntegrationTests.java @@ -33,7 +33,7 @@ public class RjcConnectionIntegrationTests extends AbstractConnectionIntegration factory.setPort(SettingsUtils.getPort()); factory.setHostName(SettingsUtils.getHost()); - factory.setUsePool(true); + factory.setUsePool(false); factory.afterPropertiesSet(); }