diff --git a/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/JdbcMessageStore.java b/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/JdbcMessageStore.java index b34dee8763..6ef9206afe 100644 --- a/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/JdbcMessageStore.java +++ b/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/JdbcMessageStore.java @@ -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 org/springframework/integration/jdbc/schema-*.sql, where * 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> serializer) { - this.serializer = new SerializingConverter((OutputStreamingConverter) serializer); + public void setSerializer(Serializer/*>*/ 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> deserializer) { - this.deserializer = new DeserializingConverter((InputStreamingConverter) deserializer); + public void setDeserializer(Deserializer/*>*/ deserializer) { + this.deserializer = new DeserializingConverter(deserializer); } /** diff --git a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/JdbcMessageStoreTests.java b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/JdbcMessageStoreTests.java index 9eca461a5d..0aab2a4e2a 100644 --- a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/JdbcMessageStoreTests.java +++ b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/JdbcMessageStoreTests.java @@ -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>() { - public void convert(Message object, OutputStream outputStream) throws IOException { - outputStream.write(object.getPayload().toString().getBytes()); + messageStore.setSerializer(new Serializer/*>>*/() { + public void serialize(/*Message*/ Object object, OutputStream outputStream) throws IOException { + outputStream.write(((Message) object).getPayload().toString().getBytes()); outputStream.flush(); } }); - messageStore.setDeserializer(new InputStreamingConverter>() { - public Message convert(InputStream inputStream) throws IOException { + messageStore.setDeserializer(new Deserializer/*>*/() { + public Message deserialize(InputStream inputStream) throws IOException { BufferedReader reader = new BufferedReader(new InputStreamReader(inputStream)); return new GenericMessage(reader.readLine()); } diff --git a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/config/JdbcMessageStoreParserTests.java b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/config/JdbcMessageStoreParserTests.java index e63637ef7c..e683f87892 100644 --- a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/config/JdbcMessageStoreParserTests.java +++ b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/config/JdbcMessageStoreParserTests.java @@ -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, Deserializer { + + private final Serializer targetSerializer = new DefaultSerializer(); + + private final Deserializer 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); } }