Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
61 changes: 52 additions & 9 deletions pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -4,24 +4,67 @@
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>

<groupId>com.netflix.rxjava</groupId>
<groupId>io.reactivex</groupId>
<artifactId>rxjava-redis</artifactId>
<version>0.0.1-SNAPSHOT</version>

<dependencies>
<dependency>
<groupId>redis.clients</groupId>
<artifactId>jedis</artifactId>
<version>2.1.0</version>
<version>2.7.2</version>
</dependency>


<dependency>
<groupId>com.netflix.rxjava</groupId>
<artifactId>rxjava-core</artifactId>
<version>0.6.0</version>
<groupId>io.reactivex</groupId>
<artifactId>rxjava</artifactId>
<version>1.0.13</version>
</dependency>
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-api</artifactId>
<version>1.7.7</version>
</dependency>
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-log4j12</artifactId>
<version>1.7.7</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>log4j</groupId>
<artifactId>log4j</artifactId>
<version>1.2.17</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.testng</groupId>
<artifactId>testng</artifactId>
<version>6.8.8</version>
<scope>test</scope>
</dependency>

</dependencies>


<build>
<testResources>
<testResource>
<directory>src/test/resources</directory>
<includes>
<include>**/*.properties</include>
</includes>
</testResource>
</testResources>
<plugins>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-compiler-plugin</artifactId>
<version>3.0</version>
<configuration>
<!-- http://maven.apache.org/plugins/maven-compiler-plugin/ -->
<source>1.7</source>
<target>1.7</target>
</configuration>
</plugin>
</plugins>
</build>

</project>
154 changes: 154 additions & 0 deletions src/main/java/rx/redis/RedisPoolPubSub.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,154 @@
package rx.redis;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import redis.clients.jedis.Jedis;
import redis.clients.jedis.JedisPool;
import redis.clients.jedis.JedisPubSub;
import redis.clients.jedis.exceptions.JedisConnectionException;
import rx.Observable;
import rx.Subscriber;
import rx.functions.Action0;
import rx.functions.Func0;
import rx.subscriptions.Subscriptions;

import java.util.concurrent.Executor;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

/**
* @author Matteo Moci ( matteo (dot) moci (at) gmail (dot) com )
*/
public class RedisPoolPubSub {

private static final Logger LOGGER = LoggerFactory.getLogger(RedisPoolPubSub.class);

public static Observable<String> observe(final JedisPool jedisPool, final String channel) {

return observe(jedisPool, channel, Executors.newSingleThreadExecutor());
}

public static Observable<String> observe(final JedisPool jedisPool, final String channel,
final ExecutorService executor) {

return Observable.defer(new Func0<Observable<String>>() {
@Override
public Observable<String> call() {

return Observable.create(new RedisObservable(jedisPool, channel, executor));
}
});
}

private static class RedisObservable implements Observable.OnSubscribe<String> {

private final JedisPool jedisPool;

private final String channel;

private final Executor executor;

public RedisObservable(final JedisPool jedisPool, final String channel,
final Executor executor) {

this.jedisPool = jedisPool;

this.channel = channel;

this.executor = executor;

}

@Override
public void call(final Subscriber<? super String> subscriber) {

final JedisPubSub jedisPubSub = new OnMessageOnNext(subscriber);

//the resource is returned to the pool by UnsubscribeAction
final Jedis jedis;
try {
jedis = jedisPool.getResource();

//subscribe
executor.execute(new Runnable() {
@Override
public void run() {

jedis.subscribe(jedisPubSub, channel);

//blocked until it's called jedisPubSub.unsubscribe()
//it closes jedis in the end

jedis.close();

}
});

//unsubscribe jedisPubSub and close hedis when observer is unsubscribed
// from http://stackoverflow.com/questions/26695125/how-to-get-notified-of-a-observers-unsubscribe-action-in-a-custom-observable-in
subscriber.add(Subscriptions.create(new Action0() {
@Override
public void call() {

jedisPubSub.unsubscribe();

}
}));

} catch (final JedisConnectionException e) {
subscriber.onError(e);
}
}

private static class OnMessageOnNext extends JedisPubSub {

private Subscriber<? super String> subscriber;

public OnMessageOnNext(final Subscriber<? super String> subscriber) {

this.subscriber = subscriber;
}

@Override
public void onMessage(final String channel, final String message) {

if (!subscriber.isUnsubscribed()) {
subscriber.onNext(message);
}

}

@Override
public void onPMessage(String pattern, String channel, String message) {

}

public void onSubscribe(String channel, int subscribedChannels) {
//TODO?
// subscriber.onStart();

}

@Override
public void onUnsubscribe(String channel, int subscribedChannels) {

//TODO?
// subscriber.onCompleted();

}

@Override
public void onPUnsubscribe(String pattern, int subscribedChannels) {

}

@Override
public void onPSubscribe(String pattern, int subscribedChannels) {

}

}
}

}

60 changes: 44 additions & 16 deletions src/main/java/rx/redis/RedisPubSub.java
Original file line number Diff line number Diff line change
@@ -1,68 +1,90 @@
package rx.redis;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import redis.clients.jedis.Jedis;
import redis.clients.jedis.JedisPubSub;
import rx.Observable;
import rx.Observer;
import rx.Subscription;
import rx.util.functions.Func0;
import rx.util.functions.Func1;
import rx.Subscriber;
import rx.functions.Action0;
import rx.functions.Func0;
import rx.subscriptions.Subscriptions;

import java.util.concurrent.Executor;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

public final class RedisPubSub {

private static final Logger LOGGER = LoggerFactory.getLogger(RedisPubSub.class);

public static Observable<String> observe(final Jedis jedis, final String channel) {

return observe(jedis, channel, Executors.newSingleThreadExecutor());
}

public static Observable<String> observe(final Jedis jedis, final String channel, final ExecutorService executor) {
public static Observable<String> observe(final Jedis jedis, final String channel,
final ExecutorService executor) {

return Observable.defer(new Func0<Observable<String>>() {
@Override
public Observable<String> call() {

return Observable.create(new RedisObservable(jedis, channel, executor));
}
});
}

private static class RedisObservable implements Func1<Observer<String>, Subscription> {
private static class RedisObservable implements Observable.OnSubscribe<String> {

private final Jedis jedis;

private final String channel;

private final Executor executor;

public RedisObservable(Jedis jedis, String channel, Executor executor) {
public RedisObservable(final Jedis jedis, final String channel, final Executor executor) {

this.jedis = jedis;

this.channel = channel;

this.executor = executor;
}

@Override
public Subscription call(final Observer<String> observer) {
public void call(final Subscriber<? super String> subscriber) {

final JedisPubSub pubSub = new JedisPubSub() {
@Override
public void onMessage(String channel, String message) {
observer.onNext(message);

subscriber.onNext(message);
}

@Override
public void onPMessage(String pattern, String channel,
String message) {
public void onPMessage(String pattern, String channel, String message) {

}

public void onSubscribe(String channel, int subscribedChannels) {

//TODO
// subscriber.onStart();

}

@Override
public void onUnsubscribe(String channel, int subscribedChannels) {
}

//TODO
// subscriber.onCompleted();

}

@Override
public void onPUnsubscribe(String pattern, int subscribedChannels) {

}

@Override
Expand All @@ -74,16 +96,22 @@ public void onPSubscribe(String pattern, int subscribedChannels) {
executor.execute(new Runnable() {
@Override
public void run() {

jedis.subscribe(pubSub, channel);

}
});

return new Subscription() {
// from http://stackoverflow.com/questions/26695125/how-to-get-notified-of-a-observers-unsubscribe-action-in-a-custom-observable-in
subscriber.add(Subscriptions.create(new Action0() {
@Override
public void unsubscribe() {
pubSub.unsubscribe();
public void call() {

pubSub.unsubscribe(channel);

}
};
}));

}
}

Expand Down
Loading