diff --git a/build.gradle b/build.gradle index b770d7c30f..cc9683e592 100644 --- a/build.gradle +++ b/build.gradle @@ -123,9 +123,9 @@ subprojects { subproject -> smack3Version = '3.2.1' smackVersion = '4.0.6' springAmqpVersion = project.hasProperty('springAmqpVersion') ? project.springAmqpVersion : '1.5.0.BUILD-SNAPSHOT' - springDataMongoVersion = '1.7.0.M1' - springDataRedisVersion = '1.5.0.M1' - springGemfireVersion = '1.6.0.M1' + springDataMongoVersion = '1.7.0.RELEASE' + springDataRedisVersion = '1.5.0.RELEASE' + springGemfireVersion = '1.6.0.RELEASE' springSecurityVersion = '3.2.5.RELEASE' springSocialTwitterVersion = '1.1.0.RELEASE' springRetryVersion = '1.1.2.RELEASE' diff --git a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/AbstractConfigurableMongoDbMessageStore.java b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/AbstractConfigurableMongoDbMessageStore.java index 01c5f926e0..8f202880dd 100644 --- a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/AbstractConfigurableMongoDbMessageStore.java +++ b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/AbstractConfigurableMongoDbMessageStore.java @@ -36,10 +36,8 @@ import org.springframework.core.convert.converter.Converter; import org.springframework.core.convert.converter.GenericConverter; import org.springframework.core.serializer.support.DeserializingConverter; import org.springframework.core.serializer.support.SerializingConverter; -import org.springframework.dao.DataAccessException; import org.springframework.data.domain.Sort; import org.springframework.data.mongodb.MongoDbFactory; -import org.springframework.data.mongodb.core.DbCallback; import org.springframework.data.mongodb.core.FindAndModifyOptions; import org.springframework.data.mongodb.core.IndexOperations; import org.springframework.data.mongodb.core.MongoTemplate; @@ -59,9 +57,6 @@ import org.springframework.messaging.Message; import org.springframework.messaging.MessageHeaders; import org.springframework.util.Assert; -import com.mongodb.DB; -import com.mongodb.MongoException; - /** * The abstract MongoDB {@link BasicMessageGroupStore} implementation to provide configuration for common options * for implementations of this class. @@ -112,7 +107,7 @@ public abstract class AbstractConfigurableMongoDbMessageStore implements BasicMe this(mongoDbFactory, null, collectionName); } - public AbstractConfigurableMongoDbMessageStore(MongoDbFactory mongoDbFactory, + public AbstractConfigurableMongoDbMessageStore(MongoDbFactory mongoDbFactory, MappingMongoConverter mappingMongoConverter, String collectionName) { Assert.notNull("'mongoDbFactory' must not be null"); Assert.hasText("'collectionName' must not be empty"); @@ -143,9 +138,9 @@ public abstract class AbstractConfigurableMongoDbMessageStore implements BasicMe this.mongoTemplate.setApplicationContext(this.applicationContext); } } - + this.messageBuilderFactory = IntegrationUtils.getMessageBuilderFactory(this.applicationContext); - + IndexOperations indexOperations = this.mongoTemplate.indexOps(this.collectionName); indexOperations.ensureIndex(new Index(MessageDocumentFields.MESSAGE_ID, Sort.Direction.ASC)); @@ -191,38 +186,32 @@ public abstract class AbstractConfigurableMongoDbMessageStore implements BasicMe } protected void addMessageDocument(final MessageDocument document) { - this.mongoTemplate.executeInSession(new DbCallback() { - @Override - public Void doInDB(DB db) throws MongoException, DataAccessException { - Message message = document.getMessage(); - if (message.getHeaders().containsKey(SAVED_KEY)) { - Message saved = getMessage(message.getHeaders().getId()); - if (saved != null) { - if (saved.equals(message)) { - return null; - } // We need to save it under its own id - } - } - - final long createdDate = document.getCreatedTime() == 0 - ? System.currentTimeMillis() - : document.getCreatedTime(); - - Message result = messageBuilderFactory.fromMessage(message).setHeader(SAVED_KEY, Boolean.TRUE) - .setHeader(CREATED_DATE_KEY, createdDate).build(); - - @SuppressWarnings("unchecked") - Map innerMap = (Map) new DirectFieldAccessor(result.getHeaders()) - .getPropertyValue("headers"); - // using reflection to set ID since it is immutable through MessageHeaders - innerMap.put(MessageHeaders.ID, message.getHeaders().get(MessageHeaders.ID)); - innerMap.put(MessageHeaders.TIMESTAMP, message.getHeaders().get(MessageHeaders.TIMESTAMP)); - - document.setCreatedTime(createdDate); - mongoTemplate.insert(document, collectionName); - return null; + Message message = document.getMessage(); + if (message.getHeaders().containsKey(SAVED_KEY)) { + Message saved = getMessage(message.getHeaders().getId()); + if (saved != null) { + if (saved.equals(message)) { + return; + } // We need to save it under its own id } - }); + } + + final long createdDate = document.getCreatedTime() == 0 + ? System.currentTimeMillis() + : document.getCreatedTime(); + + Message result = messageBuilderFactory.fromMessage(message).setHeader(SAVED_KEY, Boolean.TRUE) + .setHeader(CREATED_DATE_KEY, createdDate).build(); + + @SuppressWarnings("unchecked") + Map innerMap = (Map) new DirectFieldAccessor(result.getHeaders()) + .getPropertyValue("headers"); + // using reflection to set ID since it is immutable through MessageHeaders + innerMap.put(MessageHeaders.ID, message.getHeaders().get(MessageHeaders.ID)); + innerMap.put(MessageHeaders.TIMESTAMP, message.getHeaders().get(MessageHeaders.TIMESTAMP)); + + document.setCreatedTime(createdDate); + mongoTemplate.insert(document, collectionName); } protected static Query groupIdQuery(Object groupId) { @@ -230,9 +219,9 @@ public abstract class AbstractConfigurableMongoDbMessageStore implements BasicMe } /** - * A {@link org.springframework.core.convert.converter.GenericConverter} implementation to convert {@link org.springframework.messaging.Message} to - * serialized {@link byte[]} to store {@link org.springframework.messaging.Message} to the MongoDB. - * And vice versa - to convert {@link byte[]} from the MongoDB to the {@link org.springframework.messaging.Message}. + * A {@link GenericConverter} implementation to convert {@link Message} to + * serialized {@link byte[]} to store {@link Message} to the MongoDB. + * And vice versa - to convert {@link byte[]} from the MongoDB to the {@link Message}. */ private static class MongoDbMessageBytesConverter implements GenericConverter { diff --git a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/ConfigurableMongoDbMessageStore.java b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/ConfigurableMongoDbMessageStore.java index 8c0e750352..eb46420a50 100644 --- a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/ConfigurableMongoDbMessageStore.java +++ b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/ConfigurableMongoDbMessageStore.java @@ -1,5 +1,5 @@ /* - * Copyright 2013-2014 the original author or authors. + * Copyright 2013-2015 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. @@ -23,10 +23,8 @@ import java.util.LinkedHashSet; import java.util.List; import java.util.UUID; -import org.springframework.dao.DataAccessException; import org.springframework.data.domain.Sort; import org.springframework.data.mongodb.MongoDbFactory; -import org.springframework.data.mongodb.core.DbCallback; import org.springframework.data.mongodb.core.MongoTemplate; import org.springframework.data.mongodb.core.convert.MappingMongoConverter; import org.springframework.data.mongodb.core.query.Criteria; @@ -41,9 +39,6 @@ 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 @@ -85,14 +80,14 @@ public class ConfigurableMongoDbMessageStore extends AbstractConfigurableMongoDb this(mongoDbFactory, null, collectionName); } - public ConfigurableMongoDbMessageStore(MongoDbFactory mongoDbFactory, MappingMongoConverter mappingMongoConverter, String collectionName) { + public ConfigurableMongoDbMessageStore(MongoDbFactory mongoDbFactory, MappingMongoConverter mappingMongoConverter, + String collectionName) { super(mongoDbFactory, mappingMongoConverter, collectionName); } /** * Convenient injection point for expiry callbacks in the message store. Each of the callbacks provided will simply * be registered with the store using {@link #registerMessageGroupExpiryCallback(MessageGroupCallback)}. - * * @param expiryCallbacks the expiry callbacks to add */ public void setExpiryCallbacks(Collection expiryCallbacks) { @@ -110,7 +105,6 @@ public class ConfigurableMongoDbMessageStore extends AbstractConfigurableMongoDb * the {@link MessageGroup} was created. If you want the timeout to be based on the time * the {@link MessageGroup} was idling (e.g., inactive from the last update) invoke this method with 'true'. * Default is 'false'. - * * @param timeoutOnIdle The boolean. */ public void setTimeoutOnIdle(boolean timeoutOnIdle) { @@ -177,36 +171,31 @@ public class ConfigurableMongoDbMessageStore extends AbstractConfigurableMongoDb Assert.notNull(groupId, "'groupId' must not be null"); Assert.notNull(message, "'message' must not be null"); - return this.mongoTemplate.executeInSession(new DbCallback() { + Query query = groupOrderQuery(groupId); + MessageDocument messageDocument = mongoTemplate.findOne(query, MessageDocument.class, collectionName); - @Override - public MessageGroup doInDB(DB db) throws MongoException, DataAccessException { - Query query = groupOrderQuery(groupId); - MessageDocument messageDocument = mongoTemplate.findOne(query, MessageDocument.class, collectionName); + long createdTime = 0; + int lastReleasedSequence = 0; + boolean complete = false; - long createdTime = 0; - int lastReleasedSequence = 0; - boolean complete = false; + if (messageDocument != null) { + createdTime = messageDocument.getCreatedTime(); + lastReleasedSequence = messageDocument.getLastReleasedSequence(); + complete = messageDocument.isComplete(); + } - if (messageDocument != null) { - createdTime = messageDocument.getCreatedTime(); - lastReleasedSequence = messageDocument.getLastReleasedSequence(); - complete = messageDocument.isComplete(); - } + MessageDocument document = new MessageDocument(message); + document.setGroupId(groupId); + document.setComplete(complete); + document.setLastReleasedSequence(lastReleasedSequence); + document.setCreatedTime(createdTime == 0 ? System.currentTimeMillis() : createdTime); + document.setLastModifiedTime(System.currentTimeMillis()); + document.setSequence(getNextId()); - MessageDocument document = new MessageDocument(message); - document.setGroupId(groupId); - document.setComplete(complete); - document.setLastReleasedSequence(lastReleasedSequence); - document.setCreatedTime(createdTime == 0 ? System.currentTimeMillis() : createdTime); - document.setLastModifiedTime(System.currentTimeMillis()); - document.setSequence(getNextId()); + addMessageDocument(document); - addMessageDocument(document); + return getMessageGroup(groupId); - return getMessageGroup(groupId); - } - }); } @Override @@ -214,38 +203,27 @@ public class ConfigurableMongoDbMessageStore extends AbstractConfigurableMongoDb Assert.notNull(groupId, "'groupId' must not be null"); Assert.notNull(messageToRemove, "'messageToRemove' must not be null"); - return this.mongoTemplate.executeInSession(new DbCallback() { + Query query = groupIdQuery(groupId) + .addCriteria(Criteria.where(MessageDocumentFields.MESSAGE_ID).is(messageToRemove.getHeaders().getId())); + mongoTemplate.remove(query, collectionName); + updateGroup(groupId, lastModifiedUpdate()); + return getMessageGroup(groupId); - @Override - public MessageGroup doInDB(DB db) throws MongoException, DataAccessException { - Query query = groupIdQuery(groupId) - .addCriteria(Criteria.where(MessageDocumentFields.MESSAGE_ID).is(messageToRemove.getHeaders().getId())); - mongoTemplate.remove(query, collectionName); - updateGroup(groupId, lastModifiedUpdate()); - return getMessageGroup(groupId); - } - }); } @Override public Message pollMessageFromGroup(final Object groupId) { Assert.notNull(groupId, "'groupId' must not be null"); - return this.mongoTemplate.executeInSession(new DbCallback>() { - - @Override - public Message doInDB(DB db) throws MongoException, DataAccessException { - Sort sort = new Sort(MessageDocumentFields.LAST_MODIFIED_TIME, MessageDocumentFields.SEQUENCE); - Query query = groupIdQuery(groupId).with(sort); - MessageDocument document = mongoTemplate.findAndRemove(query, MessageDocument.class, collectionName); - Message message = null; - if (document != null) { - message = document.getMessage(); - updateGroup(groupId, lastModifiedUpdate()); - } - return message; - } - }); + Sort sort = new Sort(MessageDocumentFields.LAST_MODIFIED_TIME, MessageDocumentFields.SEQUENCE); + Query query = groupIdQuery(groupId).with(sort); + MessageDocument document = mongoTemplate.findAndRemove(query, MessageDocument.class, collectionName); + Message message = null; + if (document != null) { + message = document.getMessage(); + updateGroup(groupId, lastModifiedUpdate()); + } + return message; } @Override @@ -260,24 +238,18 @@ public class ConfigurableMongoDbMessageStore extends AbstractConfigurableMongoDb @Override public Iterator iterator() { - return this.mongoTemplate.executeInSession(new DbCallback>() { + List messageGroups = new ArrayList(); - @Override - public Iterator doInDB(DB db) throws MongoException, DataAccessException { - List messageGroups = new ArrayList(); + Query query = Query.query(Criteria.where(MessageDocumentFields.GROUP_ID).exists(true)); + @SuppressWarnings("rawtypes") + List groupIds = mongoTemplate.getCollection(collectionName) + .distinct(MessageDocumentFields.GROUP_ID, query.getQueryObject()); - Query query = Query.query(Criteria.where(MessageDocumentFields.GROUP_ID).exists(true)); - @SuppressWarnings("rawtypes") - List groupIds = mongoTemplate.getCollection(collectionName) - .distinct(MessageDocumentFields.GROUP_ID, query.getQueryObject()); + for (Object groupId : groupIds) { + messageGroups.add(getMessageGroup(groupId)); + } - for (Object groupId : groupIds) { - messageGroups.add(getMessageGroup(groupId)); - } - - return messageGroups.iterator(); - } - }); + return messageGroups.iterator(); } @Override @@ -364,7 +336,8 @@ public class ConfigurableMongoDbMessageStore extends AbstractConfigurableMongoDb } private static Query groupOrderQuery(Object groupId) { - Sort sort = new Sort(Sort.Direction.DESC, MessageDocumentFields.LAST_MODIFIED_TIME, MessageDocumentFields.SEQUENCE); + Sort sort = new Sort(Sort.Direction.DESC, MessageDocumentFields.LAST_MODIFIED_TIME, + MessageDocumentFields.SEQUENCE); return groupIdQuery(groupId).with(sort); } diff --git a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/MongoDbMessageStore.java b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/MongoDbMessageStore.java index 63024a8fda..7b78153462 100644 --- a/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/MongoDbMessageStore.java +++ b/spring-integration-mongodb/src/main/java/org/springframework/integration/mongodb/store/MongoDbMessageStore.java @@ -38,14 +38,12 @@ import org.springframework.core.convert.converter.Converter; import org.springframework.core.convert.converter.GenericConverter; import org.springframework.core.serializer.support.DeserializingConverter; import org.springframework.core.serializer.support.SerializingConverter; -import org.springframework.dao.DataAccessException; import org.springframework.data.annotation.Id; import org.springframework.data.annotation.Transient; import org.springframework.data.convert.WritingConverter; import org.springframework.data.domain.Sort; import org.springframework.data.mapping.context.MappingContext; import org.springframework.data.mongodb.MongoDbFactory; -import org.springframework.data.mongodb.core.DbCallback; import org.springframework.data.mongodb.core.FindAndModifyOptions; import org.springframework.data.mongodb.core.IndexOperations; import org.springframework.data.mongodb.core.MongoTemplate; @@ -78,9 +76,7 @@ import org.springframework.util.StringUtils; import com.mongodb.BasicDBList; import com.mongodb.BasicDBObject; -import com.mongodb.DB; import com.mongodb.DBObject; -import com.mongodb.MongoException; /** @@ -141,7 +137,6 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore /** * Create a MongoDbMessageStore using the provided {@link MongoDbFactory}.and the default collection name. - * * @param mongoDbFactory The mongodb factory. */ public MongoDbMessageStore(MongoDbFactory mongoDbFactory) { @@ -150,7 +145,6 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore /** * Create a MongoDbMessageStore using the provided {@link MongoDbFactory} and collection name. - * * @param mongoDbFactory The mongodb factory. * @param collectionName The collection name. */ @@ -196,41 +190,39 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore } private void addMessageDocument(final MessageWrapper document) { - this.template.executeInSession(new DbCallback() { - @Override - public Void doInDB(DB db) throws MongoException, DataAccessException { - Message message = document.getMessage(); - if (message.getHeaders().containsKey(SAVED_KEY)) { - Message saved = getMessage(message.getHeaders().getId()); - if (saved != null) { - if (saved.equals(message)) { - return null; - } // We need to save it under its own id - } - } - - final long createdDate = document.get_Group_timestamp() == 0 ? System.currentTimeMillis() : document.get_Group_timestamp(); - - Message result = getMessageBuilderFactory().fromMessage(message).setHeader(SAVED_KEY, Boolean.TRUE) - .setHeader(CREATED_DATE_KEY, createdDate).build(); - - @SuppressWarnings("unchecked") - Map innerMap = (Map) new DirectFieldAccessor(result.getHeaders()).getPropertyValue("headers"); - // using reflection to set ID since it is immutable through MessageHeaders - innerMap.put(MessageHeaders.ID, message.getHeaders().get(MessageHeaders.ID)); - innerMap.put(MessageHeaders.TIMESTAMP, message.getHeaders().get(MessageHeaders.TIMESTAMP)); - - document.set_Group_timestamp(createdDate); - template.insert(document, collectionName); - return null; + Message message = document.getMessage(); + if (message.getHeaders().containsKey(SAVED_KEY)) { + Message saved = getMessage(message.getHeaders().getId()); + if (saved != null) { + if (saved.equals(message)) { + return; + } // We need to save it under its own id } - }); + } + + final long createdDate = document.get_Group_timestamp() == 0 + ? System.currentTimeMillis() + : document.get_Group_timestamp(); + + Message result = getMessageBuilderFactory().fromMessage(message).setHeader(SAVED_KEY, Boolean.TRUE) + .setHeader(CREATED_DATE_KEY, createdDate).build(); + + @SuppressWarnings("unchecked") + Map innerMap = + (Map) new DirectFieldAccessor(result.getHeaders()).getPropertyValue("headers"); + // using reflection to set ID since it is immutable through MessageHeaders + innerMap.put(MessageHeaders.ID, message.getHeaders().get(MessageHeaders.ID)); + innerMap.put(MessageHeaders.TIMESTAMP, message.getHeaders().get(MessageHeaders.TIMESTAMP)); + + document.set_Group_timestamp(createdDate); + template.insert(document, collectionName); } @Override public Message getMessage(UUID id) { Assert.notNull(id, "'id' must not be null"); - MessageWrapper messageWrapper = this.template.findOne(whereMessageIdIs(id), MessageWrapper.class, this.collectionName); + MessageWrapper messageWrapper = + this.template.findOne(whereMessageIdIs(id), MessageWrapper.class, this.collectionName); return (messageWrapper != null) ? messageWrapper.getMessage() : null; } @@ -243,8 +235,9 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore @Override public Message removeMessage(UUID id) { Assert.notNull(id, "'id' must not be null"); - MessageWrapper messageWrapper = this.template.findAndRemove(whereMessageIdIs(id), MessageWrapper.class, this.collectionName); - return (messageWrapper != null) ? messageWrapper.getMessage() : null; + MessageWrapper messageWrapper = + this.template.findAndRemove(whereMessageIdIs(id), MessageWrapper.class, this.collectionName); + return (messageWrapper != null ? messageWrapper.getMessage() : null); } @Override @@ -282,36 +275,31 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore public MessageGroup addMessageToGroup(final Object groupId, final Message message) { Assert.notNull(groupId, "'groupId' must not be null"); Assert.notNull(message, "'message' must not be null"); - return this.template.executeInSession(new DbCallback() { + Query query = whereGroupIdOrder(groupId); + MessageWrapper messageDocument = template.findOne(query, MessageWrapper.class, collectionName); - @Override - public MessageGroup doInDB(DB db) throws MongoException, DataAccessException { - Query query = whereGroupIdOrder(groupId); - MessageWrapper messageDocument = template.findOne(query, MessageWrapper.class, collectionName); + long createdTime = 0; + int lastReleasedSequence = 0; + boolean complete = false; - long createdTime = 0; - int lastReleasedSequence = 0; - boolean complete = false; - - if (messageDocument != null) { - createdTime = messageDocument.get_Group_timestamp(); - lastReleasedSequence = messageDocument.get_LastReleasedSequenceNumber(); - complete = messageDocument.get_Group_complete(); - } + if (messageDocument != null) { + createdTime = messageDocument.get_Group_timestamp(); + lastReleasedSequence = messageDocument.get_LastReleasedSequenceNumber(); + complete = messageDocument.get_Group_complete(); + } - MessageWrapper wrapper = new MessageWrapper(message); - wrapper.set_GroupId(groupId); - wrapper.set_Group_timestamp(createdTime == 0 ? System.currentTimeMillis() : createdTime); - wrapper.set_Group_update_timestamp(System.currentTimeMillis()); - wrapper.set_Group_complete(complete); - wrapper.set_LastReleasedSequenceNumber(lastReleasedSequence); - wrapper.setSequence(getNextId()); + MessageWrapper wrapper = new MessageWrapper(message); + wrapper.set_GroupId(groupId); + wrapper.set_Group_timestamp(createdTime == 0 ? System.currentTimeMillis() : createdTime); + wrapper.set_Group_update_timestamp(System.currentTimeMillis()); + wrapper.set_Group_complete(complete); + wrapper.set_LastReleasedSequenceNumber(lastReleasedSequence); + wrapper.setSequence(getNextId()); + + addMessageDocument(wrapper); + return getMessageGroup(groupId); - addMessageDocument(wrapper); - return getMessageGroup(groupId); - } - }); } @Override @@ -319,16 +307,10 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore Assert.notNull(groupId, "'groupId' must not be null"); Assert.notNull(messageToRemove, "'messageToRemove' must not be null"); - return this.template.executeInSession(new DbCallback() { - - @Override - public MessageGroup doInDB(DB db) throws MongoException, DataAccessException { - template.findAndRemove(whereMessageIdIsAndGroupIdIs(messageToRemove.getHeaders().getId(), groupId), - MessageWrapper.class, collectionName); - updateGroup(groupId, lastModifiedUpdate()); - return getMessageGroup(groupId); - } - }); + template.findAndRemove(whereMessageIdIsAndGroupIdIs(messageToRemove.getHeaders().getId(), groupId), + MessageWrapper.class, collectionName); + updateGroup(groupId, lastModifiedUpdate()); + return getMessageGroup(groupId); } @Override @@ -338,41 +320,32 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore @Override public Iterator iterator() { - return this.template.executeInSession(new DbCallback>() { + List messageGroups = new ArrayList(); - @Override - public Iterator doInDB(DB db) throws MongoException, DataAccessException { - List messageGroups = new ArrayList(); + Query query = Query.query(Criteria.where(GROUP_ID_KEY).exists(true)); - Query query = Query.query(Criteria.where(GROUP_ID_KEY).exists(true)); - @SuppressWarnings("rawtypes") - List groupIds = template.getCollection(collectionName) - .distinct(GROUP_ID_KEY, query.getQueryObject()); + @SuppressWarnings("rawtypes") + List groupIds = template.getCollection(collectionName) + .distinct(GROUP_ID_KEY, query.getQueryObject()); - for (Object groupId : groupIds) { - messageGroups.add(getMessageGroup(groupId)); - } + for (Object groupId : groupIds) { + messageGroups.add(getMessageGroup(groupId)); + } - return messageGroups.iterator(); - } - }); + return messageGroups.iterator(); } + @Override public Message pollMessageFromGroup(final Object groupId) { Assert.notNull(groupId, "'groupId' must not be null"); - return this.template.executeInSession(new DbCallback>() { - @Override - public Message doInDB(DB db) throws MongoException, DataAccessException { - Query query = whereGroupIdIs(groupId).with(new Sort(GROUP_UPDATE_TIMESTAMP_KEY, SEQUENCE)); - MessageWrapper messageWrapper = template.findAndRemove(query, MessageWrapper.class, collectionName); - Message message = null; - if (messageWrapper != null) { - message = messageWrapper.getMessage(); - } - updateGroup(groupId, lastModifiedUpdate()); - return message; - } - }); + Query query = whereGroupIdIs(groupId).with(new Sort(GROUP_UPDATE_TIMESTAMP_KEY, SEQUENCE)); + MessageWrapper messageWrapper = template.findAndRemove(query, MessageWrapper.class, collectionName); + Message message = null; + if (messageWrapper != null) { + message = messageWrapper.getMessage(); + } + updateGroup(groupId, lastModifiedUpdate()); + return message; } @Override @@ -585,7 +558,8 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore @SuppressWarnings("unchecked") private static void enhanceHeaders(MessageHeaders messageHeaders, Map headers) { - Map innerMap = (Map) new DirectFieldAccessor(messageHeaders).getPropertyValue("headers"); + Map innerMap = + (Map) new DirectFieldAccessor(messageHeaders).getPropertyValue("headers"); // using reflection to set ID and TIMESTAMP since they are immutable through MessageHeaders innerMap.put(MessageHeaders.ID, headers.get(MessageHeaders.ID)); innerMap.put(MessageHeaders.TIMESTAMP, headers.get(MessageHeaders.TIMESTAMP)); @@ -635,9 +609,11 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore public GenericMessage convert(DBObject source) { @SuppressWarnings("unchecked") - Map headers = MongoDbMessageStore.this.converter.normalizeHeaders((Map) source.get("headers")); + Map headers = + MongoDbMessageStore.this.converter.normalizeHeaders((Map) source.get("headers")); - GenericMessage message = new GenericMessage(MongoDbMessageStore.this.converter.extractPayload(source), headers); + GenericMessage message = + new GenericMessage(MongoDbMessageStore.this.converter.extractPayload(source), headers); enhanceHeaders(message.getHeaders(), headers); return message; } @@ -669,9 +645,12 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore public Object convert(Object source, TypeDescriptor sourceType, TypeDescriptor targetType) { DBObject dbObject = (DBObject) source; @SuppressWarnings("unchecked") - Map headers = MongoDbMessageStore.this.converter.normalizeHeaders((Map) dbObject.get("headers")); + Map headers = + MongoDbMessageStore.this.converter.normalizeHeaders((Map) dbObject.get("headers")); - return MutableMessageBuilder.withPayload(MongoDbMessageStore.this.converter.extractPayload(dbObject)).copyHeaders(headers).build(); + return MutableMessageBuilder.withPayload(MongoDbMessageStore.this.converter.extractPayload(dbObject)) + .copyHeaders(headers) + .build(); } } @@ -680,7 +659,8 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore @Override public AdviceMessage convert(DBObject source) { @SuppressWarnings("unchecked") - Map headers = MongoDbMessageStore.this.converter.normalizeHeaders((Map) source.get("headers")); + Map headers = + MongoDbMessageStore.this.converter.normalizeHeaders((Map) source.get("headers")); Message inputMessage = null; @@ -696,7 +676,8 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore } } - AdviceMessage message = new AdviceMessage(MongoDbMessageStore.this.converter.extractPayload(source), headers, inputMessage); + AdviceMessage message = + new AdviceMessage(MongoDbMessageStore.this.converter.extractPayload(source), headers, inputMessage); enhanceHeaders(message.getHeaders(), headers); return message; @@ -711,7 +692,8 @@ public class MongoDbMessageStore extends AbstractMessageGroupStore @Override public ErrorMessage convert(DBObject source) { @SuppressWarnings("unchecked") - Map headers = MongoDbMessageStore.this.converter.normalizeHeaders((Map) source.get("headers")); + Map headers = + MongoDbMessageStore.this.converter.normalizeHeaders((Map) source.get("headers")); Object payload = this.deserializingConverter.convert((byte[]) source.get("payload")); ErrorMessage message = new ErrorMessage((Throwable) payload, headers);