diff --git a/src/main/java/org/springframework/data/redis/listener/adapter/MessageListenerAdapter.java b/src/main/java/org/springframework/data/redis/listener/adapter/MessageListenerAdapter.java index 041a1eff1..a6d8dd423 100644 --- a/src/main/java/org/springframework/data/redis/listener/adapter/MessageListenerAdapter.java +++ b/src/main/java/org/springframework/data/redis/listener/adapter/MessageListenerAdapter.java @@ -372,7 +372,7 @@ public class MessageListenerAdapter implements InitializingBean, MessageListener if (serializer != null) { return serializer.deserialize(message.getBody()); } - return message; + return message.getBody(); } /** diff --git a/src/test/java/org/springframework/data/redis/listener/PubSubTestParams.java b/src/test/java/org/springframework/data/redis/listener/PubSubTestParams.java index 1de87e1b9..32aad4d91 100644 --- a/src/test/java/org/springframework/data/redis/listener/PubSubTestParams.java +++ b/src/test/java/org/springframework/data/redis/listener/PubSubTestParams.java @@ -21,6 +21,7 @@ import java.util.Collection; import org.springframework.data.redis.ObjectFactory; import org.springframework.data.redis.Person; import org.springframework.data.redis.PersonObjectFactory; +import org.springframework.data.redis.RawObjectFactory; import org.springframework.data.redis.SettingsUtils; import org.springframework.data.redis.StringObjectFactory; import org.springframework.data.redis.connection.jedis.JedisConnectionFactory; @@ -39,6 +40,7 @@ public class PubSubTestParams { // create Jedis Factory ObjectFactory stringFactory = new StringObjectFactory(); ObjectFactory personFactory = new PersonObjectFactory(); + ObjectFactory rawFactory = new RawObjectFactory(); JedisConnectionFactory jedisConnFactory = new JedisConnectionFactory(); jedisConnFactory.setUsePool(true); @@ -52,6 +54,10 @@ public class PubSubTestParams { RedisTemplate personTemplate = new RedisTemplate(); personTemplate.setConnectionFactory(jedisConnFactory); personTemplate.afterPropertiesSet(); + RedisTemplate rawTemplate = new RedisTemplate(); + rawTemplate.setEnableDefaultSerializer(false); + rawTemplate.setConnectionFactory(jedisConnFactory); + rawTemplate.afterPropertiesSet(); // add Lettuce LettuceConnectionFactory lettuceConnFactory = new LettuceConnectionFactory(); @@ -63,6 +69,10 @@ public class PubSubTestParams { RedisTemplate personTemplateLtc = new RedisTemplate(); personTemplateLtc.setConnectionFactory(lettuceConnFactory); personTemplateLtc.afterPropertiesSet(); + RedisTemplate rawTemplateLtc = new RedisTemplate(); + rawTemplateLtc.setEnableDefaultSerializer(false); + rawTemplateLtc.setConnectionFactory(lettuceConnFactory); + rawTemplateLtc.afterPropertiesSet(); // SRP SrpConnectionFactory srpConnFactory = new SrpConnectionFactory(); @@ -74,12 +84,17 @@ public class PubSubTestParams { RedisTemplate personTemplateSrp = new RedisTemplate(); personTemplateSrp.setConnectionFactory(srpConnFactory); personTemplateSrp.afterPropertiesSet(); + RedisTemplate rawTemplateSrp = new RedisTemplate(); + rawTemplateSrp.setEnableDefaultSerializer(false); + rawTemplateSrp.setConnectionFactory(srpConnFactory); + rawTemplateSrp.afterPropertiesSet(); // JRedis does not support pub/sub return Arrays.asList(new Object[][] { { stringFactory, stringTemplate }, { personFactory, personTemplate }, - { stringFactory, stringTemplateLtc }, { personFactory, personTemplateLtc }, - { stringFactory, stringTemplateSrp }, { personFactory, personTemplateSrp } + {rawFactory, rawTemplate}, { stringFactory, stringTemplateLtc }, { personFactory, personTemplateLtc }, + {rawFactory, rawTemplateLtc}, { stringFactory, stringTemplateSrp }, { personFactory, personTemplateSrp }, + {rawFactory, rawTemplateSrp} }); } } diff --git a/src/test/java/org/springframework/data/redis/listener/PubSubTests.java b/src/test/java/org/springframework/data/redis/listener/PubSubTests.java index 74e29bf69..71ee2dad7 100644 --- a/src/test/java/org/springframework/data/redis/listener/PubSubTests.java +++ b/src/test/java/org/springframework/data/redis/listener/PubSubTests.java @@ -16,8 +16,9 @@ package org.springframework.data.redis.listener; import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertTrue; import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertThat; +import static org.junit.matchers.JUnitMatchers.hasItems; import java.util.Arrays; import java.util.Collection; @@ -57,11 +58,11 @@ public class PubSubTests { @SuppressWarnings("rawtypes") protected RedisTemplate template; - private final BlockingDeque bag = new LinkedBlockingDeque(99); + private final BlockingDeque bag = new LinkedBlockingDeque(99); private final Object handler = new Object() { @SuppressWarnings("unused") - public void handleMessage(String message) { + public void handleMessage(Object message) { bag.add(message); } }; @@ -117,27 +118,27 @@ public class PubSubTests { return factory.instance(); } + @SuppressWarnings("unchecked") @Test public void testContainerSubscribe() throws Exception { - String payload1 = "do"; - String payload2 = "re mi"; + T payload1 = getT(); + T payload2 = getT(); template.convertAndSend(CHANNEL, payload1); template.convertAndSend(CHANNEL, payload2); - Set set = new LinkedHashSet(); - set.add(bag.poll(1, TimeUnit.SECONDS)); - set.add(bag.poll(1, TimeUnit.SECONDS)); + Set set = new LinkedHashSet(); + set.add((T) bag.poll(1, TimeUnit.SECONDS)); + set.add((T) bag.poll(1, TimeUnit.SECONDS)); - assertTrue(set.contains(payload1)); - assertTrue(set.contains(payload2)); + assertThat(set, hasItems(payload1, payload2)); } @Test public void testMessageBatch() throws Exception { int COUNT = 10; for (int i = 0; i < COUNT; i++) { - template.convertAndSend(CHANNEL, "message=" + i); + template.convertAndSend(CHANNEL, getT()); } Thread.sleep(1000); @@ -146,8 +147,8 @@ public class PubSubTests { @Test public void testContainerUnsubscribe() throws Exception { - String payload1 = "do"; - String payload2 = "re mi"; + T payload1 = getT(); + T payload2 = getT(); container.removeMessageListener(adapter, new ChannelTopic(CHANNEL)); template.convertAndSend(CHANNEL, payload1);