INT-3382 Support Easier Aggregator Extensions

JIRA: https://jira.spring.io/browse/INT-3382

- Make private methods protected; add protected getters
- Support obtaining just metadata from stores
- Support just fetching the first message from a group
- add metaSize() to message group

INT-3382: Polishing
This commit is contained in:
Gary Russell
2014-04-25 19:31:53 +03:00
committed by Artem Bilan
parent 5e106f20dc
commit 5522b99741
8 changed files with 189 additions and 29 deletions

View File

@@ -299,6 +299,74 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageH
return messageStore;
}
protected Map<UUID, ScheduledFuture<?>> getExpireGroupScheduledFutures() {
return expireGroupScheduledFutures;
}
protected MessageGroupProcessor getOutputProcessor() {
return outputProcessor;
}
protected CorrelationStrategy getCorrelationStrategy() {
return correlationStrategy;
}
protected ReleaseStrategy getReleaseStrategy() {
return releaseStrategy;
}
protected MessageChannel getOutputChannel() {
return outputChannel;
}
protected String getOutputChannelName() {
return outputChannelName;
}
protected MessagingTemplate getMessagingTemplate() {
return messagingTemplate;
}
protected MessageChannel getDiscardChannel() {
return discardChannel;
}
protected String getDiscardChannelName() {
return discardChannelName;
}
protected boolean isSendPartialResultOnExpiry() {
return sendPartialResultOnExpiry;
}
protected boolean isSequenceAware() {
return sequenceAware;
}
protected LockRegistry getLockRegistry() {
return lockRegistry;
}
protected boolean isLockRegistrySet() {
return lockRegistrySet;
}
protected long getMinimumTimeoutForEmptyGroups() {
return minimumTimeoutForEmptyGroups;
}
protected boolean isReleasePartialSequences() {
return releasePartialSequences;
}
protected Expression getGroupTimeoutExpression() {
return groupTimeoutExpression;
}
protected EvaluationContext getEvaluationContext() {
return evaluationContext;
}
@Override
protected void handleMessageInternal(Message<?> message) throws Exception {
Object correlationKey = correlationStrategy.getCorrelationKey(message);
@@ -480,11 +548,11 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageH
return new IntegrationMessageHeaderAccessor(lastReleasedMessage).getSequenceNumber();
}
private MessageGroup store(Object correlationKey, Message<?> message) {
protected MessageGroup store(Object correlationKey, Message<?> message) {
return messageStore.addMessageToGroup(correlationKey, message);
}
private void expireGroup(Object correlationKey, MessageGroup group) {
protected void expireGroup(Object correlationKey, MessageGroup group) {
if (logger.isInfoEnabled()) {
logger.info("Expiring MessageGroup with correlationKey[" + correlationKey + "]");
}
@@ -506,7 +574,7 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageH
}
}
private void completeGroup(Object correlationKey, MessageGroup group) {
protected void completeGroup(Object correlationKey, MessageGroup group) {
Message<?> first = null;
if (group != null) {
first = group.getOne();
@@ -515,7 +583,7 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageH
}
@SuppressWarnings("unchecked")
private Collection<Message<?>> completeGroup(Message<?> message, Object correlationKey, MessageGroup group) {
protected Collection<Message<?>> completeGroup(Message<?> message, Object correlationKey, MessageGroup group) {
if (logger.isDebugEnabled()) {
logger.debug("Completing group with correlationKey [" + correlationKey + "]");
}
@@ -530,13 +598,13 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageH
return partialSequence;
}
private void verifyResultCollectionConsistsOfMessages(Collection<?> elements){
protected void verifyResultCollectionConsistsOfMessages(Collection<?> elements){
Class<?> commonElementType = CollectionUtils.findCommonElementType(elements);
Assert.isAssignable(Message.class, commonElementType, "The expected collection of Messages contains non-Message element: " + commonElementType);
}
@SuppressWarnings("rawtypes")
private void sendReplies(Object processorResult, Message message) {
protected void sendReplies(Object processorResult, Message message) {
Object replyChannelHeader = null;
if (message != null) {
replyChannelHeader = message.getHeaders().getReplyChannel();
@@ -555,7 +623,7 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageH
}
}
private void sendReplyMessage(Object reply, Object replyChannel) {
protected void sendReplyMessage(Object reply, Object replyChannel) {
if (replyChannel instanceof MessageChannel) {
if (reply instanceof Message<?>) {
this.messagingTemplate.send((MessageChannel) replyChannel, (Message<?>) reply);
@@ -573,7 +641,7 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageH
}
}
private boolean shouldSendMultipleReplies(Iterable<?> iter) {
protected boolean shouldSendMultipleReplies(Iterable<?> iter) {
for (Object next : iter) {
if (next instanceof Message<?>) {
return true;
@@ -582,7 +650,7 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageH
return false;
}
private Long obtainGroupTimeout(MessageGroup group) {
protected Long obtainGroupTimeout(MessageGroup group) {
return this.groupTimeoutExpression != null
? this.groupTimeoutExpression.getValue(this.evaluationContext, group, Long.class) : null;
}
@@ -594,7 +662,7 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageH
}
}
private static class SequenceAwareMessageGroup extends SimpleMessageGroup {
protected static class SequenceAwareMessageGroup extends SimpleMessageGroup {
public SequenceAwareMessageGroup(MessageGroup messageGroup) {
super(messageGroup);

View File

@@ -26,6 +26,7 @@ import org.springframework.integration.support.MessageBuilderFactory;
import org.springframework.integration.support.utils.IntegrationUtils;
import org.springframework.jmx.export.annotation.ManagedAttribute;
import org.springframework.jmx.export.annotation.ManagedResource;
import org.springframework.messaging.Message;
/**
* @author Dave Syer
@@ -135,6 +136,16 @@ public abstract class AbstractMessageGroupStore implements MessageGroupStore, It
return count;
}
@Override
public MessageGroupMetadata getGroupMetadata(Object groupId) {
throw new UnsupportedOperationException("Not yet implemented for this store");
}
@Override
public Message<?> getOneMessageFromGroup(Object groupId) {
throw new UnsupportedOperationException("Not yet implemented for this store");
}
private void expire(MessageGroup group) {
RuntimeException exception = null;

View File

@@ -28,6 +28,7 @@ import org.springframework.messaging.Message;
*
* @author Dave Syer
* @author Oleg Zhurakousky
* @author Gary Russell
*/
public interface MessageGroup {

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2002-2011 the original author or authors.
* Copyright 2002-2014 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
@@ -27,43 +27,65 @@ import org.springframework.util.Assert;
/**
* Immutable Value Object holding metadata about a MessageGroup.
*
*
* @author Oleg Zhurakousky
* @author Gary Russell
* @since 2.1
*/
public class MessageGroupMetadata implements Serializable{
public class MessageGroupMetadata implements Serializable {
private static final long serialVersionUID = 1L;
private final Object groupId;
private final List<UUID> messageIds = new LinkedList<UUID>();
private final boolean complete;
private final long timestamp;
private volatile long lastModified;
private final int lastReleasedMessageSequenceNumber;
private final boolean hasMessages;
private final int size;
private final UUID first;
public MessageGroupMetadata(MessageGroup messageGroup) {
this(messageGroup, true, null);
}
public MessageGroupMetadata(MessageGroup messageGroup, boolean hasMessages, UUID first) {
Assert.notNull(messageGroup, "'messageGroup' must not be null");
this.groupId = messageGroup.getGroupId();
for (Message<?> message : messageGroup.getMessages()) {
this.messageIds.add(message.getHeaders().getId());
if (hasMessages) {
for (Message<?> message : messageGroup.getMessages()) {
this.messageIds.add(message.getHeaders().getId());
}
this.size = this.messageIds.size();
}
else {
this.size = messageGroup.size();
}
this.complete = messageGroup.isComplete();
this.timestamp = messageGroup.getTimestamp();
this.lastReleasedMessageSequenceNumber = messageGroup.getLastReleasedMessageSequenceNumber();
this.lastModified = messageGroup.getLastModified();
this.hasMessages = hasMessages;
this.first = first;
}
public void remove(UUID messageId){
if (!this.hasMessages) {
throw new IllegalStateException("Messages are not available, fetch the entire group");
}
this.messageIds.remove(messageId);
}
public void setLastModified(long lastModified) {
this.lastModified = lastModified;
}
@@ -71,27 +93,35 @@ public class MessageGroupMetadata implements Serializable{
public Object getGroupId() {
return this.groupId;
}
public Iterator<UUID> messageIdIterator(){
if (!this.hasMessages) {
throw new IllegalStateException("Messages are not available, fetch the entire group");
}
return this.messageIds.iterator();
}
public int size(){
return this.messageIds.size();
return this.size;
}
public UUID firstId(){
if (this.first != null) {
return this.first;
}
if (!this.hasMessages) {
throw new IllegalStateException("Messages are not available, fetch the entire group");
}
if (this.messageIds.size() > 0){
return this.messageIds.iterator().next();
}
return null;
}
public boolean isComplete() {
return this.complete;
}
public long getLastModified() {
return lastModified;
}

View File

@@ -101,6 +101,23 @@ public interface MessageGroupStore extends BasicMessageGroupStore {
*/
void completeGroup(Object groupId);
/**
* Obtain the group metadata without fetching any messages; must supply all other
* group properties; may include the id of the first message.
* @param groupId The group id.
* @return The metadata.
* @since 4.0
*/
MessageGroupMetadata getGroupMetadata(Object groupId);
/**
* Return the one {@link org.springframework.messaging.Message} from {@link org.springframework.integration.store.MessageGroup}.
* @param groupId The group identifier.
* @return the {@link org.springframework.messaging.Message}.
* @since 4.0
*/
Message<?> getOneMessageFromGroup(Object groupId);
/**
* Invoked when a MessageGroupStore expires a group.
*/

View File

@@ -68,6 +68,7 @@ public class SimpleMessageGroup implements MessageGroup {
this(messageGroup.getMessages(), messageGroup.getGroupId(), messageGroup.getTimestamp(), messageGroup.isComplete());
}
@Override
public long getTimestamp() {
return timestamp;
}
@@ -76,10 +77,12 @@ public class SimpleMessageGroup implements MessageGroup {
this.lastModified = lastModified;
}
@Override
public long getLastModified() {
return lastModified;
}
@Override
public boolean canAdd(Message<?> message) {
return true;
}
@@ -92,6 +95,7 @@ public class SimpleMessageGroup implements MessageGroup {
messages.remove(message);
}
@Override
public int getLastReleasedMessageSequenceNumber() {
return lastReleasedMessageSequence;
}
@@ -100,6 +104,7 @@ public class SimpleMessageGroup implements MessageGroup {
return this.messages.offer(message);
}
@Override
public Collection<Message<?>> getMessages() {
return Collections.unmodifiableCollection(messages);
}
@@ -108,18 +113,22 @@ public class SimpleMessageGroup implements MessageGroup {
this.lastReleasedMessageSequence = sequenceNumber;
}
@Override
public Object getGroupId() {
return groupId;
}
@Override
public boolean isComplete() {
return this.complete;
}
@Override
public void complete() {
this.complete = true;
}
@Override
public int getSequenceSize() {
if (size() == 0) {
return 0;
@@ -127,10 +136,12 @@ public class SimpleMessageGroup implements MessageGroup {
return new IntegrationMessageHeaderAccessor(getOne()).getSequenceSize();
}
@Override
public int size() {
return this.messages.size();
}
@Override
public Message<?> getOne() {
Message<?> one = messages.peek();
return one;

View File

@@ -298,4 +298,15 @@ public class SimpleMessageStore extends AbstractMessageGroupStore
public int messageGroupSize(Object groupId) {
return this.getMessageGroup(groupId).size();
}
@Override
public MessageGroupMetadata getGroupMetadata(Object groupId) {
return new MessageGroupMetadata(this.getMessageGroup(groupId));
}
@Override
public Message<?> getOneMessageFromGroup(Object groupId) {
return this.getMessageGroup(groupId).getOne();
}
}

View File

@@ -23,9 +23,6 @@ import java.util.LinkedHashSet;
import java.util.List;
import java.util.UUID;
import com.mongodb.DB;
import com.mongodb.MongoException;
import org.springframework.dao.DataAccessException;
import org.springframework.data.domain.Sort;
import org.springframework.data.mongodb.MongoDbFactory;
@@ -36,6 +33,7 @@ import org.springframework.data.mongodb.core.query.Criteria;
import org.springframework.data.mongodb.core.query.Query;
import org.springframework.data.mongodb.core.query.Update;
import org.springframework.integration.store.MessageGroup;
import org.springframework.integration.store.MessageGroupMetadata;
import org.springframework.integration.store.MessageGroupStore;
import org.springframework.integration.store.MessageStore;
import org.springframework.integration.store.SimpleMessageGroup;
@@ -43,6 +41,9 @@ import org.springframework.jmx.export.annotation.ManagedAttribute;
import org.springframework.messaging.Message;
import org.springframework.util.Assert;
import com.mongodb.DB;
import com.mongodb.MongoException;
/**
* An alternate MongoDB {@link MessageStore} and {@link MessageGroupStore} which allows the user to
* configure the instance of {@link MongoTemplate}. The mechanism of storing the messages/group of messages
@@ -322,6 +323,16 @@ public class ConfigurableMongoDbMessageStore extends AbstractConfigurableMongoDb
.size();
}
@Override
public MessageGroupMetadata getGroupMetadata(Object groupId) {
throw new UnsupportedOperationException("Not yet implemented for this store");
}
@Override
public Message<?> getOneMessageFromGroup(Object groupId) {
throw new UnsupportedOperationException("Not yet implemented for this store");
}
private void expire(MessageGroup group) {
RuntimeException exception = null;