diff --git a/spring-integration-core/src/main/java/org/springframework/integration/store/MessageGroup.java b/spring-integration-core/src/main/java/org/springframework/integration/store/MessageGroup.java index de78c69724..47deceac8a 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/store/MessageGroup.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/store/MessageGroup.java @@ -50,12 +50,6 @@ public interface MessageGroup { */ int size(); - /** - * Mark all unmarked messages in the group. A MessageGroupProcessor typically invokes this method after - * processing all unmarked messages. - */ - void markAll(); - /** * @return a single message from the group */ diff --git a/spring-integration-core/src/main/java/org/springframework/integration/store/SimpleMessageGroup.java b/spring-integration-core/src/main/java/org/springframework/integration/store/SimpleMessageGroup.java index 776d956393..0700e67625 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/store/SimpleMessageGroup.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/store/SimpleMessageGroup.java @@ -115,11 +115,15 @@ public class SimpleMessageGroup implements MessageGroup { } public Collection> getUnmarked() { - return Collections.unmodifiableCollection(unmarked); + synchronized (lock) { + return Collections.unmodifiableCollection(unmarked); + } } public Collection> getMarked() { - return Collections.unmodifiableCollection(marked); + synchronized (lock) { + return Collections.unmodifiableCollection(marked); + } } public Object getCorrelationKey() { diff --git a/spring-integration-core/src/main/java/org/springframework/integration/store/SimpleMessageStore.java b/spring-integration-core/src/main/java/org/springframework/integration/store/SimpleMessageStore.java index 4d3073a101..5831fd742b 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/store/SimpleMessageStore.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/store/SimpleMessageStore.java @@ -115,10 +115,9 @@ public class SimpleMessageStore extends AbstractMessageGroupStore implements Mes public MessageGroup markMessageGroup(MessageGroup group) { Object correlationId = group.getCorrelationKey(); - MessageGroup internal = getMessageGroupInternal(correlationId); + SimpleMessageGroup internal = getMessageGroupInternal(correlationId); internal.markAll(); - group.markAll(); - return group; + return internal; } public void removeMessageGroup(Object correlationId) {