+ upgrade to latest RJC (0.6.3)

+ still left with some issues regarding the pubsub support
This commit is contained in:
Costin Leau
2011-03-30 20:05:52 +03:00
parent e118bbf98a
commit 14a1d64b39
5 changed files with 62 additions and 37 deletions

View File

@@ -16,10 +16,10 @@
<spring.range>"[3.0.0, 4.0.0)"</spring.range>
<jredis.ver>03122010</jredis.ver>
<jedis.ver>1.5.2</jedis.ver>
<rjc.ver>0.6.2</rjc.ver>
<rjc.ver>0.6.3</rjc.ver>
<jedis.range>"[1.0.0,2.0.0)"</jedis.range>
<jackson.range>"[1.6, 2.0.0)"</jackson.range>
<rjc.range>"[0.6.2, 0.6.2]"</rjc.range>
<rjc.range>"[0.6.3, 0.6.3]"</rjc.range>
</properties>
<dependencies>

View File

@@ -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);
}

View File

@@ -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));
}
}

View File

@@ -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;
}
}

View File

@@ -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);
}