diff --git a/spring-integration-redis/src/main/java/org/springframework/integration/redis/store/RedisMessageGroup.java b/spring-integration-redis/src/main/java/org/springframework/integration/redis/store/RedisMessageGroup.java index c3a491a007..8f05d01026 100644 --- a/spring-integration-redis/src/main/java/org/springframework/integration/redis/store/RedisMessageGroup.java +++ b/spring-integration-redis/src/main/java/org/springframework/integration/redis/store/RedisMessageGroup.java @@ -15,12 +15,14 @@ package org.springframework.integration.redis.store; import java.util.Collection; import java.util.LinkedList; import java.util.List; +import java.util.UUID; import org.springframework.data.redis.core.BoundSetOperations; import org.springframework.data.redis.core.BoundValueOperations; import org.springframework.data.redis.core.RedisTemplate; import org.springframework.integration.Message; import org.springframework.integration.store.MessageGroup; +import org.springframework.integration.store.MessageStore; import org.springframework.util.Assert; /** @@ -33,10 +35,12 @@ class RedisMessageGroup implements MessageGroup { private final Object groupId; private final RedisTemplate redisTemplate; + private final MessageStore messageStore; - public RedisMessageGroup(RedisTemplate redisTemplate, Object groupId){ + public RedisMessageGroup(MessageStore messageStore, RedisTemplate redisTemplate, Object groupId){ this.groupId = groupId; this.redisTemplate = redisTemplate; + this.messageStore = messageStore; } public boolean canAdd(Message message) { @@ -44,15 +48,14 @@ class RedisMessageGroup implements MessageGroup { return false; } /** - * + * Will add message to this MessageGroup, b * @param message */ public void add(Message message) { BoundSetOperations mGroupOps = this.redisTemplate.boundSetOps(UNMARKED_PREFIX + this.groupId.toString()); String key = message.getHeaders().getId().toString(); mGroupOps.add(key); - BoundValueOperations bvOps = this.redisTemplate.boundValueOps(key); - bvOps.set(message); + this.messageStore.addMessage(message); } /** * @@ -61,13 +64,15 @@ class RedisMessageGroup implements MessageGroup { // delete unmarked BoundSetOperations mGroupOps = this.redisTemplate.boundSetOps(UNMARKED_PREFIX + this.groupId.toString()); for (Object messageId : mGroupOps.members()) { - this.redisTemplate.delete(messageId.toString()); + this.messageStore.removeMessage(UUID.fromString(messageId.toString())); + //this.redisTemplate.delete(messageId.toString()); } this.redisTemplate.delete(UNMARKED_PREFIX + this.groupId.toString()); // delete marked mGroupOps = this.redisTemplate.boundSetOps(MARKED_PREFIX + this.groupId.toString()); for (Object messageId : mGroupOps.members()) { - this.redisTemplate.delete(messageId.toString()); + //this.redisTemplate.delete(messageId.toString()); + this.messageStore.removeMessage(UUID.fromString(messageId.toString())); } this.redisTemplate.delete(MARKED_PREFIX + this.groupId.toString()); } @@ -150,9 +155,7 @@ class RedisMessageGroup implements MessageGroup { private Collection> getMessages(BoundSetOperations mGroupOps){ List> messages = new LinkedList>(); for (Object messageId : mGroupOps.members()) { - BoundValueOperations bvOps = this.redisTemplate.boundValueOps(messageId.toString()); - Object message = bvOps.get(); - Assert.isInstanceOf(Message.class, message, "Retrieved object is not an instance of Message"); + Message message = this.messageStore.getMessage(UUID.fromString(messageId.toString())); messages.add((Message) message); } return messages; diff --git a/spring-integration-redis/src/main/java/org/springframework/integration/redis/store/RedisMessageStore.java b/spring-integration-redis/src/main/java/org/springframework/integration/redis/store/RedisMessageStore.java index 54151f5d22..ca84b025bf 100644 --- a/spring-integration-redis/src/main/java/org/springframework/integration/redis/store/RedisMessageStore.java +++ b/spring-integration-redis/src/main/java/org/springframework/integration/redis/store/RedisMessageStore.java @@ -149,7 +149,7 @@ public class RedisMessageStore extends AbstractMessageGroupStore implements Mess if (!mGroupsOps.isMember(groupId)){ mGroupsOps.add(groupId); } - return new RedisMessageGroup(this.redisTemplate, groupId); + return new RedisMessageGroup(this, this.redisTemplate, groupId); } /** * @@ -195,7 +195,7 @@ public class RedisMessageStore extends AbstractMessageGroupStore implements Mess BoundSetOperations mGroupsOps = this.redisTemplate.boundSetOps(MESSAGE_GROUPS_KEY); List messageGroups = new ArrayList(); for (Object msgGroupId : mGroupsOps.members()) { - messageGroups.add(new RedisMessageGroup(this.redisTemplate, msgGroupId)); + messageGroups.add(new RedisMessageGroup(this, this.redisTemplate, msgGroupId)); } return messageGroups.iterator(); }