INT-1486 updated JDBC module for refactored Serialization code
This commit is contained in:
@@ -13,13 +13,24 @@
|
||||
|
||||
package org.springframework.integration.jdbc;
|
||||
|
||||
import java.sql.PreparedStatement;
|
||||
import java.sql.ResultSet;
|
||||
import java.sql.SQLException;
|
||||
import java.sql.Timestamp;
|
||||
import java.sql.Types;
|
||||
import java.util.Iterator;
|
||||
import java.util.List;
|
||||
import java.util.UUID;
|
||||
|
||||
import javax.sql.DataSource;
|
||||
|
||||
import org.apache.commons.logging.Log;
|
||||
import org.apache.commons.logging.LogFactory;
|
||||
|
||||
import org.springframework.commons.serializer.Deserializer;
|
||||
import org.springframework.commons.serializer.DeserializingConverter;
|
||||
import org.springframework.commons.serializer.InputStreamingConverter;
|
||||
import org.springframework.commons.serializer.OutputStreamingConverter;
|
||||
import org.springframework.commons.serializer.Serializer;
|
||||
import org.springframework.commons.serializer.SerializingConverter;
|
||||
import org.springframework.commons.serializer.java.JavaStreamingConverter;
|
||||
import org.springframework.integration.Message;
|
||||
import org.springframework.integration.store.AbstractMessageGroupStore;
|
||||
import org.springframework.integration.store.MessageGroup;
|
||||
@@ -27,18 +38,16 @@ import org.springframework.integration.store.MessageStore;
|
||||
import org.springframework.integration.store.SimpleMessageGroup;
|
||||
import org.springframework.integration.support.MessageBuilder;
|
||||
import org.springframework.integration.util.UUIDConverter;
|
||||
import org.springframework.jdbc.core.*;
|
||||
import org.springframework.jdbc.core.JdbcOperations;
|
||||
import org.springframework.jdbc.core.JdbcTemplate;
|
||||
import org.springframework.jdbc.core.PreparedStatementSetter;
|
||||
import org.springframework.jdbc.core.RowMapper;
|
||||
import org.springframework.jdbc.core.SingleColumnRowMapper;
|
||||
import org.springframework.jdbc.support.lob.DefaultLobHandler;
|
||||
import org.springframework.jdbc.support.lob.LobHandler;
|
||||
import org.springframework.util.Assert;
|
||||
import org.springframework.util.StringUtils;
|
||||
|
||||
import javax.sql.DataSource;
|
||||
import java.sql.*;
|
||||
import java.util.Iterator;
|
||||
import java.util.List;
|
||||
import java.util.UUID;
|
||||
|
||||
/**
|
||||
* Implementation of {@link MessageStore} using a relational database via JDBC. SQL scripts to create the necessary
|
||||
* tables are packaged as <code>org/springframework/integration/jdbc/schema-*.sql</code>, where <code>*</code> is the
|
||||
@@ -113,9 +122,8 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
|
||||
* Convenient constructor for configuration use.
|
||||
*/
|
||||
public JdbcMessageStore() {
|
||||
JavaStreamingConverter converter = new JavaStreamingConverter();
|
||||
deserializer = new DeserializingConverter(converter);
|
||||
serializer = new SerializingConverter(converter);
|
||||
deserializer = new DeserializingConverter();
|
||||
serializer = new SerializingConverter();
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -194,8 +202,8 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
|
||||
* @param serializer the serializer to set
|
||||
*/
|
||||
@SuppressWarnings("unchecked")
|
||||
public void setSerializer(OutputStreamingConverter<? super Message<?>> serializer) {
|
||||
this.serializer = new SerializingConverter((OutputStreamingConverter<Object>) serializer);
|
||||
public void setSerializer(Serializer/*<? super Message<?>>*/ serializer) {
|
||||
this.serializer = new SerializingConverter(serializer);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -204,8 +212,8 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa
|
||||
* @param deserializer the deserializer to set
|
||||
*/
|
||||
@SuppressWarnings("unchecked")
|
||||
public void setDeserializer(InputStreamingConverter<? super Message<?>> deserializer) {
|
||||
this.deserializer = new DeserializingConverter((InputStreamingConverter<Object>) deserializer);
|
||||
public void setDeserializer(Deserializer/*<? super Message<?>>*/ deserializer) {
|
||||
this.deserializer = new DeserializingConverter(deserializer);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -37,9 +37,10 @@ import javax.sql.DataSource;
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
import org.junit.runner.RunWith;
|
||||
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.commons.serializer.InputStreamingConverter;
|
||||
import org.springframework.commons.serializer.OutputStreamingConverter;
|
||||
import org.springframework.commons.serializer.Deserializer;
|
||||
import org.springframework.commons.serializer.Serializer;
|
||||
import org.springframework.integration.Message;
|
||||
import org.springframework.integration.message.GenericMessage;
|
||||
import org.springframework.integration.store.MessageGroup;
|
||||
@@ -88,14 +89,14 @@ public class JdbcMessageStoreTests {
|
||||
@Transactional
|
||||
public void testSerializer() throws Exception {
|
||||
// N.B. these serializers are not realistic (just for test purposes)
|
||||
messageStore.setSerializer(new OutputStreamingConverter<Message<?>>() {
|
||||
public void convert(Message<?> object, OutputStream outputStream) throws IOException {
|
||||
outputStream.write(object.getPayload().toString().getBytes());
|
||||
messageStore.setSerializer(new Serializer/*<Message<?>><Message<?>>*/() {
|
||||
public void serialize(/*Message<?>*/ Object object, OutputStream outputStream) throws IOException {
|
||||
outputStream.write(((Message<?>) object).getPayload().toString().getBytes());
|
||||
outputStream.flush();
|
||||
}
|
||||
});
|
||||
messageStore.setDeserializer(new InputStreamingConverter<Message<?>>() {
|
||||
public Message<?> convert(InputStream inputStream) throws IOException {
|
||||
messageStore.setDeserializer(new Deserializer/*<Message<?>>*/() {
|
||||
public Message<?> deserialize(InputStream inputStream) throws IOException {
|
||||
BufferedReader reader = new BufferedReader(new InputStreamReader(inputStream));
|
||||
return new GenericMessage<String>(reader.readLine());
|
||||
}
|
||||
|
||||
@@ -9,7 +9,11 @@ import java.io.OutputStream;
|
||||
|
||||
import org.junit.After;
|
||||
import org.junit.Test;
|
||||
import org.springframework.commons.serializer.java.JavaStreamingConverter;
|
||||
|
||||
import org.springframework.commons.serializer.DefaultDeserializer;
|
||||
import org.springframework.commons.serializer.DefaultSerializer;
|
||||
import org.springframework.commons.serializer.Deserializer;
|
||||
import org.springframework.commons.serializer.Serializer;
|
||||
import org.springframework.context.support.ClassPathXmlApplicationContext;
|
||||
import org.springframework.integration.Message;
|
||||
import org.springframework.integration.jdbc.JdbcMessageStore;
|
||||
@@ -41,9 +45,9 @@ public class JdbcMessageStoreParserTests {
|
||||
public void testSimpleMessageStoreWithSerializer() {
|
||||
setUp("serializerJdbcMessageStore.xml", getClass());
|
||||
MessageStore store = context.getBean("messageStore", MessageStore.class);
|
||||
Object serializer = TestUtils.getPropertyValue(store, "serializer.streamingConverter");
|
||||
Object serializer = TestUtils.getPropertyValue(store, "serializer.serializer");
|
||||
assertTrue(serializer instanceof EnhancedSerializer);
|
||||
Object deserializer = TestUtils.getPropertyValue(store, "deserializer.streamingConverter");
|
||||
Object deserializer = TestUtils.getPropertyValue(store, "deserializer.deserializer");
|
||||
assertTrue(deserializer instanceof EnhancedSerializer);
|
||||
}
|
||||
|
||||
@@ -67,17 +71,22 @@ public class JdbcMessageStoreParserTests {
|
||||
context = new ClassPathXmlApplicationContext(name, cls);
|
||||
}
|
||||
|
||||
public static class EnhancedSerializer extends JavaStreamingConverter {
|
||||
@Override
|
||||
public Object convert(InputStream inputStream) throws IOException {
|
||||
Message<?> message = (Message<?>) super.convert(inputStream);
|
||||
|
||||
public static class EnhancedSerializer implements Serializer<Object>, Deserializer<Object> {
|
||||
|
||||
private final Serializer<Object> targetSerializer = new DefaultSerializer();
|
||||
|
||||
private final Deserializer<Object> targetDeserializer = new DefaultDeserializer();
|
||||
|
||||
public Object deserialize(InputStream inputStream) throws IOException {
|
||||
Message<?> message = (Message<?>) targetDeserializer.deserialize(inputStream);
|
||||
return message;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void convert(Object object, OutputStream outputStream) throws IOException {
|
||||
public void serialize(Object object, OutputStream outputStream) throws IOException {
|
||||
Message<?> message = (Message<?>) object;
|
||||
super.convert(MessageBuilder.fromMessage(message).setHeader("serializer", "CUSTOM").build(), outputStream);
|
||||
message = MessageBuilder.fromMessage(message).setHeader("serializer", "CUSTOM").build();
|
||||
targetSerializer.serialize(message, outputStream);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user