DATAREDIS-612 - Polishing.

Introduce factory methods for ChannelTopic and PatternTopic.

Original Pull Request: #295
This commit is contained in:
Mark Paluch
2018-04-26 09:20:04 +02:00
committed by Christoph Strobl
parent d356ade624
commit a3b96add06
5 changed files with 36 additions and 14 deletions

View File

@@ -53,5 +53,5 @@ The message listener container itself does not require external threading resour
ReactiveRedisConnectionFactory factory = …
ReactiveRedisMessageListenerContainer container = new ReactiveRedisMessageListenerContainer(factory);
Flux<ChannelMessage<String, String>> stream = container.receive(new ChannelTopic("my-chanel"));
Flux<ChannelMessage<String, String>> stream = container.receive(ChannelTopic.of("my-chanel"));
----

View File

@@ -42,6 +42,17 @@ public class ChannelTopic implements Topic {
this.channelName = name;
}
/**
* Create a new {@link ChannelTopic} for channel subscriptions.
*
* @param name the channel name, must not be {@literal null} or empty.
* @return the {@link ChannelTopic} for {@code channelName}.
* @since 2.1
*/
public static ChannelTopic of(String name) {
return new ChannelTopic(name);
}
/**
* @return topic name.
*/

View File

@@ -42,6 +42,17 @@ public class PatternTopic implements Topic {
this.channelPattern = pattern;
}
/**
* Create a new {@link PatternTopic} for channel subscriptions based on a {@code pattern}.
*
* @param pattern the channel pattern, must not be {@literal null} or empty.
* @return the {@link PatternTopic} for {@code pattern}.
* @since 2.1
*/
static PatternTopic of(String pattern) {
return new PatternTopic(pattern);
}
/**
* @return channel pattern.
*/

View File

@@ -96,7 +96,7 @@ public class ReactiveRedisMessageListenerContainerIntegrationTests {
ReactiveRedisMessageListenerContainer container = new ReactiveRedisMessageListenerContainer(connectionFactory);
StepVerifier.create(container.receive(new ChannelTopic(CHANNEL1))) //
StepVerifier.create(container.receive(ChannelTopic.of(CHANNEL1))) //
.then(awaitSubscription(container::getActiveSubscriptions))
.then(() -> connection.publish(CHANNEL1.getBytes(), MESSAGE.getBytes())) //
.assertNext(c -> {
@@ -114,7 +114,7 @@ public class ReactiveRedisMessageListenerContainerIntegrationTests {
ReactiveRedisMessageListenerContainer container = new ReactiveRedisMessageListenerContainer(connectionFactory);
StepVerifier.create(container.receive(new PatternTopic(PATTERN1))) //
StepVerifier.create(container.receive(PatternTopic.of(PATTERN1))) //
.then(awaitSubscription(container::getActiveSubscriptions))
.then(() -> connection.publish(CHANNEL1.getBytes(), MESSAGE.getBytes())) //
.assertNext(c -> {
@@ -136,7 +136,7 @@ public class ReactiveRedisMessageListenerContainerIntegrationTests {
RedisSerializationContext.string());
BlockingQueue<PatternMessage<String, String, String>> messages = new LinkedBlockingDeque<>();
Disposable subscription = container.receive(new PatternTopic(PATTERN1)).doOnNext(messages::add).subscribe();
Disposable subscription = container.receive(PatternTopic.of(PATTERN1)).doOnNext(messages::add).subscribe();
StepVerifier.create(template.convertAndSend(CHANNEL1, MESSAGE), 0) //
.then(awaitSubscription(container::getActiveSubscriptions)) //

View File

@@ -74,7 +74,7 @@ public class ReactiveRedisMessageListenerContainerUnitTests {
container = createContainer();
StepVerifier.create(container.receive(new PatternTopic("foo*"))).thenAwait().thenCancel().verify();
StepVerifier.create(container.receive(PatternTopic.of("foo*"))).thenAwait().thenCancel().verify();
verify(subscriptionMock).pSubscribe(getByteBuffer("foo*"));
}
@@ -85,7 +85,7 @@ public class ReactiveRedisMessageListenerContainerUnitTests {
when(subscriptionMock.receive()).thenReturn(Flux.never());
container = createContainer();
StepVerifier.create(container.receive(new PatternTopic("foo*"), new PatternTopic("bar*"))).thenRequest(1)
StepVerifier.create(container.receive(PatternTopic.of("foo*"), PatternTopic.of("bar*"))).thenRequest(1)
.thenAwait().thenCancel().verify();
verify(subscriptionMock).pSubscribe(getByteBuffer("foo*"), getByteBuffer("bar*"));
@@ -97,7 +97,7 @@ public class ReactiveRedisMessageListenerContainerUnitTests {
when(subscriptionMock.receive()).thenReturn(Flux.never());
container = createContainer();
StepVerifier.create(container.receive(new ChannelTopic("foo"))).thenAwait().thenCancel().verify();
StepVerifier.create(container.receive(ChannelTopic.of("foo"))).thenAwait().thenCancel().verify();
verify(subscriptionMock).subscribe(getByteBuffer("foo"));
}
@@ -108,7 +108,7 @@ public class ReactiveRedisMessageListenerContainerUnitTests {
when(subscriptionMock.receive()).thenReturn(Flux.never());
container = createContainer();
StepVerifier.create(container.receive(new ChannelTopic("foo"), new ChannelTopic("bar"))).thenAwait().thenCancel()
StepVerifier.create(container.receive(ChannelTopic.of("foo"), ChannelTopic.of("bar"))).thenAwait().thenCancel()
.verify();
verify(subscriptionMock).subscribe(getByteBuffer("foo"), getByteBuffer("bar"));
@@ -122,7 +122,7 @@ public class ReactiveRedisMessageListenerContainerUnitTests {
when(subscriptionMock.receive()).thenReturn(processor);
container = createContainer();
Flux<ChannelMessage<String, String>> messageStream = container.receive(new ChannelTopic("foo"));
Flux<ChannelMessage<String, String>> messageStream = container.receive(ChannelTopic.of("foo"));
StepVerifier.create(messageStream).then(() -> {
processor.onNext(createChannelMessage("foo", "message"));
@@ -141,7 +141,7 @@ public class ReactiveRedisMessageListenerContainerUnitTests {
when(subscriptionMock.receive()).thenReturn(processor);
container = createContainer();
Flux<PatternMessage<String, String, String>> messageStream = container.receive(new PatternTopic("foo*"));
Flux<PatternMessage<String, String, String>> messageStream = container.receive(PatternTopic.of("foo*"));
StepVerifier.create(messageStream).then(() -> {
processor.onNext(createPatternMessage("foo*", "foo", "message"));
@@ -164,7 +164,7 @@ public class ReactiveRedisMessageListenerContainerUnitTests {
when(subscriptionMock.receive()).thenReturn(DirectProcessor.create());
container = createContainer();
Flux<ChannelMessage<String, String>> messageStream = container.receive(new ChannelTopic("foo*"));
Flux<ChannelMessage<String, String>> messageStream = container.receive(ChannelTopic.of("foo*"));
Disposable subscription = messageStream.subscribe();
@@ -206,7 +206,7 @@ public class ReactiveRedisMessageListenerContainerUnitTests {
when(subscriptionMock.receive()).thenReturn(DirectProcessor.create());
container = createContainer();
Flux<PatternMessage<String, String, String>> messageStream = container.receive(new PatternTopic("foo*"));
Flux<PatternMessage<String, String, String>> messageStream = container.receive(PatternTopic.of("foo*"));
StepVerifier.create(messageStream).then(() -> {
@@ -230,7 +230,7 @@ public class ReactiveRedisMessageListenerContainerUnitTests {
}));
container = createContainer();
Flux<PatternMessage<String, String, String>> messageStream = container.receive(new PatternTopic("foo*"));
Flux<PatternMessage<String, String, String>> messageStream = container.receive(PatternTopic.of("foo*"));
StepVerifier.create(messageStream).then(() -> {
container.destroy();
@@ -245,7 +245,7 @@ public class ReactiveRedisMessageListenerContainerUnitTests {
when(subscriptionMock.receive()).thenReturn(processor);
container = createContainer();
Flux<PatternMessage<String, String, String>> messageStream = container.receive(new PatternTopic("foo*"));
Flux<PatternMessage<String, String, String>> messageStream = container.receive(PatternTopic.of("foo*"));
StepVerifier.create(messageStream).then(() -> {
assertThat(processor.hasDownstreams()).isTrue();