diff --git a/src/main/java/org/springframework/data/redis/connection/ReactivePubSubCommands.java b/src/main/java/org/springframework/data/redis/connection/ReactivePubSubCommands.java index 38ce433d9..ec847efbc 100644 --- a/src/main/java/org/springframework/data/redis/connection/ReactivePubSubCommands.java +++ b/src/main/java/org/springframework/data/redis/connection/ReactivePubSubCommands.java @@ -34,11 +34,13 @@ public interface ReactivePubSubCommands { /** * Creates a subscription for this connection. Connections can have multiple {@link ReactiveSubscription}s. + *

+ * Use {@link #createSubscription(SubscriptionListener)} to get notified when the subscription completes. * * @return the subscription. */ default Mono createSubscription() { - return createSubscription(SubscriptionListener.EMPTY); + return createSubscription(SubscriptionListener.NO_OP_SUBSCRIPTION_LISTENER); } /** diff --git a/src/main/java/org/springframework/data/redis/connection/SubscriptionListener.java b/src/main/java/org/springframework/data/redis/connection/SubscriptionListener.java index 1f5f616a9..fe93c9dee 100644 --- a/src/main/java/org/springframework/data/redis/connection/SubscriptionListener.java +++ b/src/main/java/org/springframework/data/redis/connection/SubscriptionListener.java @@ -29,7 +29,7 @@ public interface SubscriptionListener { /** * Empty {@link SubscriptionListener}. */ - SubscriptionListener EMPTY = new SubscriptionListener() {}; + SubscriptionListener NO_OP_SUBSCRIPTION_LISTENER = new SubscriptionListener() {}; /** * Notification when Redis has confirmed a channel subscription. diff --git a/src/main/java/org/springframework/data/redis/connection/jedis/JedisMessageListener.java b/src/main/java/org/springframework/data/redis/connection/jedis/JedisMessageListener.java index 114d96749..c9b6e7223 100644 --- a/src/main/java/org/springframework/data/redis/connection/jedis/JedisMessageListener.java +++ b/src/main/java/org/springframework/data/redis/connection/jedis/JedisMessageListener.java @@ -34,10 +34,12 @@ class JedisMessageListener extends BinaryJedisPubSub { private final SubscriptionListener subscriptionListener; JedisMessageListener(MessageListener listener) { + Assert.notNull(listener, "MessageListener is required"); + this.listener = listener; this.subscriptionListener = listener instanceof SubscriptionListener ? (SubscriptionListener) listener - : SubscriptionListener.EMPTY; + : SubscriptionListener.NO_OP_SUBSCRIPTION_LISTENER; } public void onMessage(byte[] channel, byte[] message) { diff --git a/src/main/java/org/springframework/data/redis/connection/jedis/JedisSubscription.java b/src/main/java/org/springframework/data/redis/connection/jedis/JedisSubscription.java index 5358b9f95..915c4dba4 100644 --- a/src/main/java/org/springframework/data/redis/connection/jedis/JedisSubscription.java +++ b/src/main/java/org/springframework/data/redis/connection/jedis/JedisSubscription.java @@ -30,7 +30,7 @@ class JedisSubscription extends AbstractSubscription { private final BinaryJedisPubSub jedisPubSub; - JedisSubscription(MessageListener listener, JedisMessageListener jedisPubSub, @Nullable byte[][] channels, + JedisSubscription(MessageListener listener, BinaryJedisPubSub jedisPubSub, @Nullable byte[][] channels, @Nullable byte[][] patterns) { super(listener, channels, patterns); this.jedisPubSub = jedisPubSub; diff --git a/src/main/java/org/springframework/data/redis/connection/lettuce/LettuceByteBufferPubSubListenerWrapper.java b/src/main/java/org/springframework/data/redis/connection/lettuce/LettuceByteBufferPubSubListenerWrapper.java index 03814c046..4cec5ca24 100644 --- a/src/main/java/org/springframework/data/redis/connection/lettuce/LettuceByteBufferPubSubListenerWrapper.java +++ b/src/main/java/org/springframework/data/redis/connection/lettuce/LettuceByteBufferPubSubListenerWrapper.java @@ -19,6 +19,7 @@ import io.lettuce.core.pubsub.RedisPubSubListener; import java.nio.ByteBuffer; +import org.springframework.data.redis.util.ByteUtils; import org.springframework.lang.Nullable; import org.springframework.util.Assert; @@ -99,14 +100,6 @@ class LettuceByteBufferPubSubListenerWrapper implements RedisPubSubListener Flux> receive(Iterable topics, SerializationPair channelSerializer, SerializationPair messageSerializer) { - return receive(topics, channelSerializer, messageSerializer, SubscriptionListener.EMPTY); + return receive(topics, channelSerializer, messageSerializer, SubscriptionListener.NO_OP_SUBSCRIPTION_LISTENER); } /** diff --git a/src/main/java/org/springframework/data/redis/util/ByteUtils.java b/src/main/java/org/springframework/data/redis/util/ByteUtils.java index fb71c3317..200eba577 100644 --- a/src/main/java/org/springframework/data/redis/util/ByteUtils.java +++ b/src/main/java/org/springframework/data/redis/util/ByteUtils.java @@ -139,6 +139,10 @@ public final class ByteUtils { Assert.notNull(byteBuffer, "ByteBuffer must not be null!"); + if (byteBuffer.hasArray()) { + return byteBuffer.array(); + } + ByteBuffer duplicate = byteBuffer.duplicate(); byte[] bytes = new byte[duplicate.remaining()]; duplicate.get(bytes); diff --git a/src/test/java/org/springframework/data/redis/connection/jedis/JedisSubscriptionUnitTests.java b/src/test/java/org/springframework/data/redis/connection/jedis/JedisSubscriptionUnitTests.java index beec82b75..27e9b71c4 100644 --- a/src/test/java/org/springframework/data/redis/connection/jedis/JedisSubscriptionUnitTests.java +++ b/src/test/java/org/springframework/data/redis/connection/jedis/JedisSubscriptionUnitTests.java @@ -39,7 +39,7 @@ import org.springframework.data.redis.connection.RedisInvalidSubscriptionExcepti @ExtendWith(MockitoExtension.class) class JedisSubscriptionUnitTests { - @Mock JedisMessageListener jedisPubSub; + @Mock BinaryJedisPubSub jedisPubSub; @Mock MessageListener listener; diff --git a/src/test/java/org/springframework/data/redis/core/ReactiveRedisTemplateIntegrationTests.java b/src/test/java/org/springframework/data/redis/core/ReactiveRedisTemplateIntegrationTests.java index d38205bc1..7b127b4aa 100644 --- a/src/test/java/org/springframework/data/redis/core/ReactiveRedisTemplateIntegrationTests.java +++ b/src/test/java/org/springframework/data/redis/core/ReactiveRedisTemplateIntegrationTests.java @@ -465,6 +465,7 @@ public class ReactiveRedisTemplateIntegrationTests { redisTemplate.listenToChannelLater(channel) // .doOnNext(it -> redisTemplate.convertAndSend(channel, message).subscribe()).flatMapMany(Function.identity()) // + .cast(Message.class) // why? java16 why? .as(StepVerifier::create) // .assertNext(received -> { @@ -515,6 +516,7 @@ public class ReactiveRedisTemplateIntegrationTests { stream.doOnNext(it -> redisTemplate.convertAndSend(channel, message).subscribe()) // .flatMapMany(Function.identity()) // + .cast(Message.class) // why? java16 why? .as(StepVerifier::create) // .assertNext(received -> {