INT-3689: Upgrade to Spring Data Fowler

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

Remove usage of deprecated `mongoTemplate.executeInSession`.
Since MongoDB doesn't provide the `member pinning` guarantee with `requestStart/requestDone`,
this feature has been removed in MongoDB 3.0: https://github.com/mongodb/specifications/blob/master/source/server-selection/server-selection.rst#what-happened-to-pinning
This commit is contained in:
Artem Bilan
2015-03-27 11:39:27 +02:00
committed by Gary Russell
parent efef65c52f
commit a60a8733cc
4 changed files with 170 additions and 226 deletions

View File

@@ -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'

View File

@@ -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<Void>() {
@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<String, Object> innerMap = (Map<String, Object>) 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<String, Object> innerMap = (Map<String, Object>) 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 {

View File

@@ -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<MessageGroupCallback> 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<MessageGroup>() {
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<MessageGroup>() {
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<Message<?>>() {
@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<MessageGroup> iterator() {
return this.mongoTemplate.executeInSession(new DbCallback<Iterator<MessageGroup>>() {
List<MessageGroup> messageGroups = new ArrayList<MessageGroup>();
@Override
public Iterator<MessageGroup> doInDB(DB db) throws MongoException, DataAccessException {
List<MessageGroup> messageGroups = new ArrayList<MessageGroup>();
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);
}

View File

@@ -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<Void>() {
@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<String, Object> innerMap = (Map<String, Object>) 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<String, Object> innerMap =
(Map<String, Object>) 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<MessageGroup>() {
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<MessageGroup>() {
@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<MessageGroup> iterator() {
return this.template.executeInSession(new DbCallback<Iterator<MessageGroup>>() {
List<MessageGroup> messageGroups = new ArrayList<MessageGroup>();
@Override
public Iterator<MessageGroup> doInDB(DB db) throws MongoException, DataAccessException {
List<MessageGroup> messageGroups = new ArrayList<MessageGroup>();
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<Message<?>>() {
@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<String, Object> headers) {
Map<String, Object> innerMap = (Map<String, Object>) new DirectFieldAccessor(messageHeaders).getPropertyValue("headers");
Map<String, Object> innerMap =
(Map<String, Object>) 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<String, Object> headers = MongoDbMessageStore.this.converter.normalizeHeaders((Map<String, Object>) source.get("headers"));
Map<String, Object> headers =
MongoDbMessageStore.this.converter.normalizeHeaders((Map<String, Object>) source.get("headers"));
GenericMessage<?> message = new GenericMessage<Object>(MongoDbMessageStore.this.converter.extractPayload(source), headers);
GenericMessage<?> message =
new GenericMessage<Object>(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<String, Object> headers = MongoDbMessageStore.this.converter.normalizeHeaders((Map<String, Object>) dbObject.get("headers"));
Map<String, Object> headers =
MongoDbMessageStore.this.converter.normalizeHeaders((Map<String, Object>) 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<String, Object> headers = MongoDbMessageStore.this.converter.normalizeHeaders((Map<String, Object>) source.get("headers"));
Map<String, Object> headers =
MongoDbMessageStore.this.converter.normalizeHeaders((Map<String, Object>) 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<String, Object> headers = MongoDbMessageStore.this.converter.normalizeHeaders((Map<String, Object>) source.get("headers"));
Map<String, Object> headers =
MongoDbMessageStore.this.converter.normalizeHeaders((Map<String, Object>) source.get("headers"));
Object payload = this.deserializingConverter.convert((byte[]) source.get("payload"));
ErrorMessage message = new ErrorMessage((Throwable) payload, headers);