diff --git a/spring-data-redis/pom.xml b/spring-data-redis/pom.xml index c0fe69d97..0c33106f1 100644 --- a/spring-data-redis/pom.xml +++ b/spring-data-redis/pom.xml @@ -16,10 +16,10 @@ "[3.0.0, 4.0.0)" 03122010 1.5.2 - 0.6.2 + 0.6.3 "[1.0.0,2.0.0)" "[1.6, 2.0.0)" - "[0.6.2, 0.6.2]" + "[0.6.3, 0.6.3]" diff --git a/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/rjc/RjcConnection.java b/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/rjc/RjcConnection.java index 77d3a576b..dfbaf20ec 100644 --- a/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/rjc/RjcConnection.java +++ b/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/rjc/RjcConnection.java @@ -54,12 +54,11 @@ public class RjcConnection implements RedisConnection { private volatile RjcSubscription subscription; private volatile RedisNodeSubscriber subscriber; - private final Object pubSubMonitor = new Object(); - public RjcConnection(org.idevlab.rjc.ds.RedisConnection connection, int dbIndex) { SingleDataSource connectionDataSource = new SingleDataSource(connection); session = new SessionFactoryImpl(connectionDataSource).create(); - subscriber = new RedisNodeSubscriber(connectionDataSource); + subscriber = new RedisNodeSubscriber(); + subscriber.setDataSource(connectionDataSource); client = new Client(connection); this.dbIndex = dbIndex; @@ -82,6 +81,11 @@ public class RjcConnection implements RedisConnection { isClosed = true; try { subscriber.close(); + } catch (Exception ex) { + // ignore + } + + try { session.close(); } catch (Exception ex) { throw convertRjcAccessException(ex); @@ -2024,12 +2028,9 @@ public class RjcConnection implements RedisConnection { throw new UnsupportedOperationException(); } - subscription = new RjcSubscription(listener, subscriber, pubSubMonitor); + subscription = new RjcSubscription(listener, subscriber, client); subscription.pSubscribe(patterns); - synchronized (pubSubMonitor) { - pubSubMonitor.wait(); - } } catch (Exception ex) { throw convertRjcAccessException(ex); } @@ -2050,13 +2051,9 @@ public class RjcConnection implements RedisConnection { throw new UnsupportedOperationException(); } - subscription = new RjcSubscription(listener, subscriber, pubSubMonitor); + subscription = new RjcSubscription(listener, subscriber, client); subscription.subscribe(channels); - synchronized (pubSubMonitor) { - pubSubMonitor.wait(); - } - } catch (Exception ex) { throw convertRjcAccessException(ex); } 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 0cd0d48cf..64d8bc959 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 @@ -15,6 +15,7 @@ */ package org.springframework.data.keyvalue.redis.connection.rjc; +import org.idevlab.rjc.Client; import org.idevlab.rjc.message.RedisNodeSubscriber; import org.springframework.data.keyvalue.redis.connection.MessageListener; import org.springframework.data.keyvalue.redis.connection.util.AbstractSubscription; @@ -27,52 +28,62 @@ import org.springframework.data.keyvalue.redis.connection.util.AbstractSubscript class RjcSubscription extends AbstractSubscription { private final RedisNodeSubscriber subscriber; - private final RjcMessageListener listenerAdapter; - private final Object pubSubMonitor; + private final Client client; + // rjc does not support subscription while listening + // so we have to handle this ourselves through the client + private volatile boolean subscribed = false; - RjcSubscription(MessageListener listener, RedisNodeSubscriber subscriber, Object pubSubMonitor) { + RjcSubscription(MessageListener listener, RedisNodeSubscriber subscriber, Client client) { super(listener); this.subscriber = subscriber; - this.listenerAdapter = new RjcMessageListener(listener); - this.pubSubMonitor = pubSubMonitor; + subscriber.setMessageListener(new RjcMessageListener(listener)); + subscriber.setPMessageListener(new RjcMessageListener(listener)); + this.client = client; } @Override protected void doClose() { - try { - subscriber.close(); - } finally { - synchronized (pubSubMonitor) { - pubSubMonitor.notifyAll(); - } - } + subscribed = false; + client.unsubscribe(); + client.punsubscribe(); + client.rollbackTimeout(); } @Override protected void doPsubscribe(byte[]... patterns) { - for (String str : RjcUtils.decodeMultiple(patterns)) { - subscriber.psubscribe(str, listenerAdapter); + String[] pats = RjcUtils.decodeMultiple(patterns); + + if (subscribed) { + client.psubscribe(pats); + } + else { + subscriber.setPatterns(RjcUtils.addArray(subscriber.getPatterns(), pats)); + subscribed = true; + subscriber.subscribe(); } } @Override protected void doPUnsubscribe(boolean all, byte[]... patterns) { - for (String str : RjcUtils.decodeMultiple(patterns)) { - subscriber.punsubscribe(str); - } + client.punsubscribe(RjcUtils.decodeMultiple(patterns)); } @Override protected void doSubscribe(byte[]... channels) { - for (String str : RjcUtils.decodeMultiple(channels)) { - subscriber.subscribe(str, listenerAdapter); + String[] chs = RjcUtils.decodeMultiple(channels); + + if (subscribed) { + client.subscribe(chs); + } + else { + subscriber.setPatterns(RjcUtils.addArray(subscriber.getChannels(), chs)); + subscribed = true; + subscriber.subscribe(); } } @Override protected void doUnsubscribe(boolean all, byte[]... channels) { - for (String str : RjcUtils.decodeMultiple(channels)) { - subscriber.unsubscribe(str); - } + client.punsubscribe(RjcUtils.decodeMultiple(channels)); } } \ No newline at end of file diff --git a/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/rjc/RjcUtils.java b/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/rjc/RjcUtils.java index 817d47589..e373fe469 100644 --- a/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/rjc/RjcUtils.java +++ b/spring-data-redis/src/main/java/org/springframework/data/keyvalue/redis/connection/rjc/RjcUtils.java @@ -16,6 +16,7 @@ package org.springframework.data.keyvalue.redis.connection.rjc; import java.io.StringReader; +import java.util.Arrays; import java.util.Collection; import java.util.LinkedHashMap; import java.util.LinkedHashSet; @@ -41,6 +42,7 @@ import org.springframework.data.keyvalue.redis.connection.RedisZSetCommands.Tupl import org.springframework.data.keyvalue.redis.connection.SortParameters.Order; import org.springframework.data.keyvalue.redis.connection.SortParameters.Range; import org.springframework.data.keyvalue.redis.connection.util.DecodeUtils; +import org.springframework.util.ObjectUtils; /** @@ -220,4 +222,18 @@ public abstract class RjcUtils { static Double convert(String zscore) { return (zscore == null ? null : Double.valueOf(zscore)); } + + + static String[] addArray(String[] one, String[] two) { + if (ObjectUtils.isEmpty(one)) { + return two; + } + if (ObjectUtils.isEmpty(two)) { + return one; + } + + String[] result = Arrays.copyOf(one, one.length + two.length); + System.arraycopy(two, 0, result, one.length, two.length); + return result; + } } \ No newline at end of file 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 116908460..875d65b79 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 @@ -235,7 +235,8 @@ public abstract class AbstractConnectionIntegrationTests { connection.publish(channel, "two".getBytes()); connection.publish(channel, "I see you".getBytes()); System.out.println("Done publishing..."); - Thread.sleep(3000); + Thread.sleep(5000); + System.out.println("Done waiting ..."); } finally { flag.set(false); }