updated MongoMessageStore to use a MongoDbFactory rather than direct Mongo instance
This commit is contained in:
@@ -23,16 +23,13 @@ import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.UUID;
|
||||
|
||||
import org.bson.types.ObjectId;
|
||||
import org.springframework.beans.DirectFieldAccessor;
|
||||
import org.springframework.beans.factory.BeanClassLoaderAware;
|
||||
import org.springframework.core.convert.converter.Converter;
|
||||
import org.springframework.data.mapping.context.MappingContext;
|
||||
import org.springframework.data.mongodb.MongoDbFactory;
|
||||
import org.springframework.data.mongodb.core.MongoTemplate;
|
||||
import org.springframework.data.mongodb.core.SimpleMongoDbFactory;
|
||||
import org.springframework.data.mongodb.core.convert.MappingMongoConverter;
|
||||
import org.springframework.data.mongodb.core.convert.MongoConverter;
|
||||
import org.springframework.data.mongodb.core.mapping.MongoMappingContext;
|
||||
import org.springframework.data.mongodb.core.mapping.MongoPersistentEntity;
|
||||
import org.springframework.data.mongodb.core.mapping.MongoPersistentProperty;
|
||||
@@ -45,7 +42,6 @@ import org.springframework.util.Assert;
|
||||
import org.springframework.util.ClassUtils;
|
||||
|
||||
import com.mongodb.DBObject;
|
||||
import com.mongodb.Mongo;
|
||||
|
||||
/**
|
||||
* @author Mark Fisher
|
||||
@@ -55,39 +51,26 @@ public class MongoMessageStore implements MessageStore, BeanClassLoaderAware {
|
||||
|
||||
private final static String DEFAULT_COLLECTION_NAME = "messages";
|
||||
|
||||
private final MongoTemplate template;
|
||||
|
||||
private volatile String collectionName = DEFAULT_COLLECTION_NAME;
|
||||
|
||||
//private final MongoConverter mongoConverter = new MessageReadingMongoConverter();
|
||||
private final MongoConverter mongoConverter = null;
|
||||
private final MongoTemplate template;
|
||||
|
||||
private final String collectionName;
|
||||
|
||||
private volatile ClassLoader classLoader = ClassUtils.getDefaultClassLoader();
|
||||
|
||||
|
||||
public MongoMessageStore(Mongo mongo, String databaseName) {
|
||||
Assert.notNull(mongo, "mongo must not be null");
|
||||
Assert.hasText(databaseName, "databaseName must not be empty");
|
||||
MongoDbFactory mongoFactory = new SimpleMongoDbFactory(mongo, databaseName);
|
||||
MessageReadingMongoConverter converter = new MessageReadingMongoConverter(mongoFactory, new MongoMappingContext());
|
||||
this.template = new MongoTemplate(mongoFactory, converter);
|
||||
// this.template.setDefaultCollectionName(DEFAULT_COLLECTION_NAME);
|
||||
// this.template.createCollection(DEFAULT_COLLECTION_NAME);
|
||||
public MongoMessageStore(MongoDbFactory mongoDbFactory) {
|
||||
this(mongoDbFactory, null);
|
||||
}
|
||||
|
||||
|
||||
public void setCollectionName(String collectionName) {
|
||||
Assert.hasText(collectionName, "collectionName must not be empty");
|
||||
this.collectionName = collectionName;
|
||||
public MongoMessageStore(MongoDbFactory mongoDbFactory, String collectionName) {
|
||||
Assert.notNull(mongoDbFactory, "mongoDbFactory must not be null");
|
||||
MessageReadingMongoConverter converter = new MessageReadingMongoConverter(mongoDbFactory, new MongoMappingContext());
|
||||
this.template = new MongoTemplate(mongoDbFactory, converter);
|
||||
this.collectionName = (collectionName != null) ? collectionName : DEFAULT_COLLECTION_NAME;
|
||||
//this.template.createCollection(collectionName);
|
||||
}
|
||||
|
||||
// public void setUsername(String username) {
|
||||
// template.setUsername(username);
|
||||
// }
|
||||
//
|
||||
// public void setPassword(String password) {
|
||||
// template.setPassword(password);
|
||||
// }
|
||||
|
||||
public void setBeanClassLoader(ClassLoader classLoader) {
|
||||
Assert.notNull(classLoader, "classLoader must not be null");
|
||||
@@ -119,11 +102,9 @@ public class MongoMessageStore implements MessageStore, BeanClassLoaderAware {
|
||||
|
||||
private class MessageReadingMongoConverter extends MappingMongoConverter {
|
||||
|
||||
public MessageReadingMongoConverter(
|
||||
MongoDbFactory mongoDbFactory,
|
||||
public MessageReadingMongoConverter(MongoDbFactory mongoDbFactory,
|
||||
MappingContext<? extends MongoPersistentEntity<?>, MongoPersistentProperty> mappingContext) {
|
||||
super(mongoDbFactory, mappingContext);
|
||||
// TODO Auto-generated constructor stub
|
||||
}
|
||||
|
||||
|
||||
@@ -140,7 +121,7 @@ public class MongoMessageStore implements MessageStore, BeanClassLoaderAware {
|
||||
@Override
|
||||
public void write(Object source, DBObject target) {
|
||||
if (source instanceof Message) {
|
||||
String payloadType = ((Message<?>) source).getPayload().getClass().getName();
|
||||
//String payloadType = ((Message<?>) source).getPayload().getClass().getName();
|
||||
//target.put("_payloadType", payloadType);
|
||||
target.put("_id", ((Message<?>) source).getHeaders().getId().toString());
|
||||
//target.put("_id", ((Message<?>) source).getHeaders().getId().toString());
|
||||
|
||||
@@ -17,6 +17,8 @@
|
||||
package org.springframework.integration.mongodb.store;
|
||||
|
||||
import org.junit.Test;
|
||||
import org.springframework.data.mongodb.MongoDbFactory;
|
||||
import org.springframework.data.mongodb.core.SimpleMongoDbFactory;
|
||||
import org.springframework.integration.Message;
|
||||
import org.springframework.integration.mongodb.rules.MongodbAvailable;
|
||||
import org.springframework.integration.mongodb.rules.MongodbAvailableTests;
|
||||
@@ -36,8 +38,8 @@ public class MongoClaimCheckIntegrationTests extends MongodbAvailableTests{
|
||||
@Test
|
||||
@MongodbAvailable
|
||||
public void stringPayload() throws Exception {
|
||||
Mongo mongo = new Mongo();
|
||||
MongoMessageStore messageStore = new MongoMessageStore(mongo, "test");
|
||||
MongoDbFactory mongoDbFactory = new SimpleMongoDbFactory(new Mongo(), "test");
|
||||
MongoMessageStore messageStore = new MongoMessageStore(mongoDbFactory);
|
||||
ClaimCheckInTransformer checkin = new ClaimCheckInTransformer(messageStore);
|
||||
ClaimCheckOutTransformer checkout = new ClaimCheckOutTransformer(messageStore);
|
||||
Message<?> originalMessage = MessageBuilder.withPayload("test1").build();
|
||||
@@ -55,8 +57,8 @@ public class MongoClaimCheckIntegrationTests extends MongodbAvailableTests{
|
||||
@Test
|
||||
@MongodbAvailable
|
||||
public void objectPayload() throws Exception {
|
||||
Mongo mongo = new Mongo();
|
||||
MongoMessageStore messageStore = new MongoMessageStore(mongo, "test");
|
||||
MongoDbFactory mongoDbFactory = new SimpleMongoDbFactory(new Mongo(), "test");
|
||||
MongoMessageStore messageStore = new MongoMessageStore(mongoDbFactory);
|
||||
ClaimCheckInTransformer checkin = new ClaimCheckInTransformer(messageStore);
|
||||
ClaimCheckOutTransformer checkout = new ClaimCheckOutTransformer(messageStore);
|
||||
Beverage payload = new Beverage();
|
||||
|
||||
@@ -17,6 +17,9 @@
|
||||
package org.springframework.integration.mongodb.store;
|
||||
|
||||
import org.junit.Test;
|
||||
|
||||
import org.springframework.data.mongodb.MongoDbFactory;
|
||||
import org.springframework.data.mongodb.core.SimpleMongoDbFactory;
|
||||
import org.springframework.integration.Message;
|
||||
import org.springframework.integration.mongodb.rules.MongodbAvailable;
|
||||
import org.springframework.integration.mongodb.rules.MongodbAvailableTests;
|
||||
@@ -35,43 +38,52 @@ public class MongoMessageStoreTests extends MongodbAvailableTests{
|
||||
@Test
|
||||
@MongodbAvailable
|
||||
public void addGetWithStringPayload() throws Exception {
|
||||
Mongo mongo = new Mongo();
|
||||
MongoMessageStore store = new MongoMessageStore(mongo, "test");
|
||||
Message<?> message = MessageBuilder.withPayload("Hello").build();
|
||||
System.out.println(message);
|
||||
store.addMessage(message);
|
||||
assertNotNull(store.getMessage(message.getHeaders().getId()));
|
||||
MongoDbFactory mongoDbFactory = new SimpleMongoDbFactory(new Mongo(), "test");
|
||||
MongoMessageStore store = new MongoMessageStore(mongoDbFactory);
|
||||
Message<?> messageToStore = MessageBuilder.withPayload("Hello").build();
|
||||
//System.out.println("before: " + messageToStore);
|
||||
store.addMessage(messageToStore);
|
||||
Message<?> retrievedMessage = store.getMessage(messageToStore.getHeaders().getId());
|
||||
//System.out.println("after: " + retrievedMessage);
|
||||
assertNotNull(retrievedMessage);
|
||||
}
|
||||
|
||||
|
||||
@Test
|
||||
@MongodbAvailable
|
||||
public void addGetWithObjectDefaultConstructorPayload() throws Exception {
|
||||
Mongo mongo = new Mongo();
|
||||
MongoMessageStore store = new MongoMessageStore(mongo, "test");
|
||||
MongoDbFactory mongoDbFactory = new SimpleMongoDbFactory(new Mongo(), "test");
|
||||
MongoMessageStore store = new MongoMessageStore(mongoDbFactory);
|
||||
Person p = new Person();
|
||||
p.setFname("John");
|
||||
p.setLname("Doe");
|
||||
|
||||
Message<?> message = MessageBuilder.withPayload(p).build();
|
||||
System.out.println(message);
|
||||
store.addMessage(message);
|
||||
Message<?> m = store.getMessage(message.getHeaders().getId());
|
||||
assertNotNull(m);
|
||||
Message<?> messageToStore = MessageBuilder.withPayload(p).build();
|
||||
//System.out.println("before: " + messageToStore);
|
||||
store.addMessage(messageToStore);
|
||||
Message<?> retrievedMessage = store.getMessage(messageToStore.getHeaders().getId());
|
||||
//System.out.println("after: " + retrievedMessage);
|
||||
assertNotNull(retrievedMessage);
|
||||
}
|
||||
|
||||
public static class Person{
|
||||
|
||||
|
||||
public static class Person {
|
||||
|
||||
private String fname;
|
||||
|
||||
private String lname;
|
||||
|
||||
public String getFname() {
|
||||
return fname;
|
||||
}
|
||||
|
||||
public void setFname(String fname) {
|
||||
this.fname = fname;
|
||||
}
|
||||
|
||||
public String getLname() {
|
||||
return lname;
|
||||
}
|
||||
|
||||
public void setLname(String lname) {
|
||||
this.lname = lname;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user