From 88165c888126f09c18022bbf2843df3fdf480fad Mon Sep 17 00:00:00 2001 From: Gunnar Hillert Date: Wed, 10 Apr 2013 00:10:14 -0400 Subject: [PATCH] INT-2980 Jdbc Message Store: Polling Wrong Message * Add Join to JDBC Query * Add MySql-specific tests * INT-2987 - MessageStore MySQL - Support Fractional Seconds - Update MySql Connector version to 5.1.24 - Add DDL scripts for MySql versions 5.6.4 and higher INT-2980 - Check also for Regions * Add test INT-2980/INT-2987 - Add Documentation INT-2980 - Doc: Add info to What's New Section --- build.gradle | 2 +- spring-integration-jdbc/build.gradle | 2 +- .../integration/jdbc/JdbcMessageStore.java | 10 +- .../jdbc/schema-drop-mysql-5_6_4.sql | 6 + .../integration/jdbc/schema-mysql-5_6_4.sql | 29 + .../src/main/sql/mysql-5_6_4.properties | 14 + .../src/main/sql/mysql-5_6_4.vpp | 4 + .../jdbc/JdbcMessageStoreTests-context.xml | 2 +- .../jdbc/JdbcMessageStoreTests.java | 81 +++ .../jdbc/mysql/CommonMySql-context.xml | 21 + ...ssageStoreMultipleChannelTests-context.xml | 99 ++++ ...lJdbcMessageStoreMultipleChannelTests.java | 175 ++++++ .../MySqlJdbcMessageStoreTests-context.xml | 21 + .../mysql/MySqlJdbcMessageStoreTests.java | 520 ++++++++++++++++++ src/reference/docbook/jdbc.xml | 29 + src/reference/docbook/whats-new.xml | 12 +- 16 files changed, 1020 insertions(+), 7 deletions(-) create mode 100644 spring-integration-jdbc/src/main/resources/org/springframework/integration/jdbc/schema-drop-mysql-5_6_4.sql create mode 100644 spring-integration-jdbc/src/main/resources/org/springframework/integration/jdbc/schema-mysql-5_6_4.sql create mode 100644 spring-integration-jdbc/src/main/sql/mysql-5_6_4.properties create mode 100644 spring-integration-jdbc/src/main/sql/mysql-5_6_4.vpp create mode 100644 spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/mysql/CommonMySql-context.xml create mode 100644 spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/mysql/MySqlJdbcMessageStoreMultipleChannelTests-context.xml create mode 100644 spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/mysql/MySqlJdbcMessageStoreMultipleChannelTests.java create mode 100644 spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/mysql/MySqlJdbcMessageStoreTests-context.xml create mode 100644 spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/mysql/MySqlJdbcMessageStoreTests.java diff --git a/build.gradle b/build.gradle index c18d6eb0e9..9a6e589527 100644 --- a/build.gradle +++ b/build.gradle @@ -302,7 +302,7 @@ project('spring-integration-jdbc') { testCompile "org.powermock:powermock-api-mockito:1.5" testCompile "postgresql:postgresql:9.1-901-1.jdbc4" - testCompile "mysql:mysql-connector-java:5.1.21" + testCompile "mysql:mysql-connector-java:5.1.24" testCompile "commons-dbcp:commons-dbcp:1.4" //testCompile "com.oracle:ojdbc6:11.2.0.3" diff --git a/spring-integration-jdbc/build.gradle b/spring-integration-jdbc/build.gradle index 72cc936eef..88049c8b29 100644 --- a/spring-integration-jdbc/build.gradle +++ b/spring-integration-jdbc/build.gradle @@ -24,7 +24,7 @@ task generateSql { classpath: configurations.vpp.asPath) doLast { - ['hsqldb', 'h2', 'db2', 'derby', 'mysql', + ['hsqldb', 'h2', 'db2', 'derby', 'mysql', 'mysql-5_6_4', 'oracle10g', 'postgresql', 'sqlserver', 'sybase'].each { dbType -> ant.vppcopy(todir: generatedResourcesDir, overwrite: 'true') { config { 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 046c9ef23e..190f5f5682 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 @@ -115,11 +115,15 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa POLL_FROM_GROUP("SELECT %PREFIX%MESSAGE.MESSAGE_ID, %PREFIX%MESSAGE.MESSAGE_BYTES from %PREFIX%MESSAGE " + "where %PREFIX%MESSAGE.MESSAGE_ID = " + - "(SELECT min(MESSAGE_ID) from %PREFIX%MESSAGE where CREATED_DATE = " + + "(SELECT min(m.MESSAGE_ID) from %PREFIX%MESSAGE m " + + "join %PREFIX%GROUP_TO_MESSAGE on m.MESSAGE_ID = %PREFIX%GROUP_TO_MESSAGE.MESSAGE_ID " + + "where CREATED_DATE = " + "(SELECT min(CREATED_DATE) from %PREFIX%MESSAGE, %PREFIX%GROUP_TO_MESSAGE " + "where %PREFIX%MESSAGE.MESSAGE_ID = %PREFIX%GROUP_TO_MESSAGE.MESSAGE_ID " + "and %PREFIX%GROUP_TO_MESSAGE.GROUP_KEY = ? " + - "and %PREFIX%MESSAGE.REGION = ?))"), + "and %PREFIX%MESSAGE.REGION = ?) " + + "and %PREFIX%GROUP_TO_MESSAGE.GROUP_KEY = ? " + + "and m.REGION = ?)"), GET_GROUP_INFO("SELECT COMPLETE, LAST_RELEASED_SEQUENCE, CREATED_DATE, UPDATED_DATE" + " from %PREFIX%MESSAGE_GROUP where GROUP_KEY = ? and REGION=?"), @@ -606,7 +610,7 @@ public class JdbcMessageStore extends AbstractMessageGroupStore implements Messa * @return a message; could be null if query produced no Messages */ protected Message doPollForMessage(String groupIdKey) { - List> messages = jdbcTemplate.query(getQuery(Query.POLL_FROM_GROUP), new Object[] { groupIdKey, region }, mapper); + List> messages = jdbcTemplate.query(getQuery(Query.POLL_FROM_GROUP), new Object[] { groupIdKey, region, groupIdKey, region }, mapper); Assert.isTrue(messages.size() == 0 || messages.size() == 1); if (messages.size() > 0){ return messages.get(0); diff --git a/spring-integration-jdbc/src/main/resources/org/springframework/integration/jdbc/schema-drop-mysql-5_6_4.sql b/spring-integration-jdbc/src/main/resources/org/springframework/integration/jdbc/schema-drop-mysql-5_6_4.sql new file mode 100644 index 0000000000..fce4d8c950 --- /dev/null +++ b/spring-integration-jdbc/src/main/resources/org/springframework/integration/jdbc/schema-drop-mysql-5_6_4.sql @@ -0,0 +1,6 @@ +-- Autogenerated: do not edit this file + +DROP TABLE IF EXISTS INT_MESSAGE ; +DROP TABLE IF EXISTS INT_MESSAGE_GROUP ; +DROP TABLE IF EXISTS INT_GROUP_TO_MESSAGE ; +DROP INDEX IF EXISTS INT_MESSAGE_IX1 ; diff --git a/spring-integration-jdbc/src/main/resources/org/springframework/integration/jdbc/schema-mysql-5_6_4.sql b/spring-integration-jdbc/src/main/resources/org/springframework/integration/jdbc/schema-mysql-5_6_4.sql new file mode 100644 index 0000000000..e3e0b79c44 --- /dev/null +++ b/spring-integration-jdbc/src/main/resources/org/springframework/integration/jdbc/schema-mysql-5_6_4.sql @@ -0,0 +1,29 @@ +-- Autogenerated: do not edit this file + +CREATE TABLE INT_MESSAGE ( + MESSAGE_ID CHAR(36), + REGION VARCHAR(100), + CREATED_DATE DATETIME(6) NOT NULL, + MESSAGE_BYTES BLOB, + constraint MESSAGE_PK primary key (MESSAGE_ID, REGION) +) ENGINE=InnoDB; + +CREATE INDEX INT_MESSAGE_IX1 ON INT_MESSAGE (CREATED_DATE); + +CREATE TABLE INT_GROUP_TO_MESSAGE ( + GROUP_KEY CHAR(36), + MESSAGE_ID CHAR(36), + REGION VARCHAR(100), + constraint GROUP_TO_MESSAG_PK primary key (GROUP_KEY, MESSAGE_ID, REGION) +) ENGINE=InnoDB; + +CREATE TABLE INT_MESSAGE_GROUP ( + GROUP_KEY CHAR(36), + REGION VARCHAR(100), + MARKED BIGINT, + COMPLETE BIGINT, + LAST_RELEASED_SEQUENCE BIGINT, + CREATED_DATE DATETIME(6) NOT NULL, + UPDATED_DATE DATETIME(6) DEFAULT NULL, + constraint MESSAGE_GROUP_PK primary key (GROUP_KEY, REGION) +) ENGINE=InnoDB; \ No newline at end of file diff --git a/spring-integration-jdbc/src/main/sql/mysql-5_6_4.properties b/spring-integration-jdbc/src/main/sql/mysql-5_6_4.properties new file mode 100644 index 0000000000..1a5c2c4046 --- /dev/null +++ b/spring-integration-jdbc/src/main/sql/mysql-5_6_4.properties @@ -0,0 +1,14 @@ +platform=mysql +# SQL language oddities +BIGINT = BIGINT +IDENTITY = +GENERATED = +VOODOO = ENGINE=InnoDB +IFEXISTSBEFORE = IF EXISTS +DOUBLE = DOUBLE PRECISION +BLOB = BLOB +CLOB = TEXT +TIMESTAMP = DATETIME(6) +VARCHAR = VARCHAR +# for generating drop statements... +SEQUENCE = TABLE diff --git a/spring-integration-jdbc/src/main/sql/mysql-5_6_4.vpp b/spring-integration-jdbc/src/main/sql/mysql-5_6_4.vpp new file mode 100644 index 0000000000..fd65fad6e3 --- /dev/null +++ b/spring-integration-jdbc/src/main/sql/mysql-5_6_4.vpp @@ -0,0 +1,4 @@ +#macro (sequence $name $value)CREATE TABLE ${name} (ID BIGINT NOT NULL) ENGINE=MYISAM; +INSERT INTO ${name} values(0); +#end +#macro (notnull $name $type)MODIFY COLUMN ${name} ${type} NOT NULL#end diff --git a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/JdbcMessageStoreTests-context.xml b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/JdbcMessageStoreTests-context.xml index a11f031e2e..26bb28c8e1 100644 --- a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/JdbcMessageStoreTests-context.xml +++ b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/JdbcMessageStoreTests-context.xml @@ -6,7 +6,7 @@ http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd"> - + 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 88548bbdfb..56e7245139 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 @@ -35,6 +35,8 @@ import java.util.UUID; import javax.sql.DataSource; +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; import org.junit.Before; import org.junit.Test; import org.junit.runner.RunWith; @@ -42,6 +44,7 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.core.serializer.Deserializer; import org.springframework.core.serializer.Serializer; import org.springframework.integration.Message; +import org.springframework.integration.MessageHeaders; import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.history.MessageHistory; import org.springframework.integration.message.GenericMessage; @@ -67,6 +70,8 @@ import org.springframework.transaction.annotation.Transactional; @RunWith(SpringJUnit4ClassRunner.class) public class JdbcMessageStoreTests { + private static final Log LOG = LogFactory.getLog(JdbcMessageStoreTests.class); + @Autowired private DataSource dataSource; @@ -389,4 +394,80 @@ public class JdbcMessageStoreTests { assertEquals(1, group.size()); } + @Test + @Transactional + public void testSameMessageToMultipleGroups() throws Exception { + + final String group1Id = "group1"; + final String group2Id = "group2"; + + final Message message = MessageBuilder.withPayload("foo").build(); + + final MessageBuilder builder1 = MessageBuilder.fromMessage(message); + final MessageBuilder builder2 = MessageBuilder.fromMessage(message); + + builder1.setSequenceNumber(1); + builder2.setSequenceNumber(2); + + final Message message1 = builder1.build(); + final Message message2 = builder2.build(); + + messageStore.addMessageToGroup(group1Id, message1); + messageStore.addMessageToGroup(group2Id, message2); + + final Message messageFromGroup1 = messageStore.pollMessageFromGroup(group1Id); + final Message messageFromGroup2 = messageStore.pollMessageFromGroup(group2Id); + + assertNotNull(messageFromGroup1); + assertNotNull(messageFromGroup2); + + LOG.info("messageFromGroup1: " + messageFromGroup1.getHeaders().getId() + "; Sequence #: " + messageFromGroup1.getHeaders().getSequenceNumber()); + LOG.info("messageFromGroup2: " + messageFromGroup2.getHeaders().getId() + "; Sequence #: " + messageFromGroup2.getHeaders().getSequenceNumber()); + + assertEquals(Integer.valueOf(1), (Integer) messageFromGroup1.getHeaders().get(MessageHeaders.SEQUENCE_NUMBER)); + assertEquals(Integer.valueOf(2), (Integer) messageFromGroup2.getHeaders().get(MessageHeaders.SEQUENCE_NUMBER)); + + } + + @Test + @Transactional + public void testSameMessageAndGroupToMultipleRegions() throws Exception { + + final String groupId = "myGroup"; + final String region1 = "region1"; + final String region2 = "region2"; + + final JdbcMessageStore messageStore1 = new JdbcMessageStore(dataSource); + messageStore1.setRegion(region1); + + final JdbcMessageStore messageStore2 = new JdbcMessageStore(dataSource); + messageStore1.setRegion(region2); + + final Message message = MessageBuilder.withPayload("foo").build(); + + final MessageBuilder builder1 = MessageBuilder.fromMessage(message); + final MessageBuilder builder2 = MessageBuilder.fromMessage(message); + + builder1.setSequenceNumber(1); + builder2.setSequenceNumber(2); + + final Message message1 = builder1.build(); + final Message message2 = builder2.build(); + + messageStore1.addMessageToGroup(groupId, message1); + messageStore2.addMessageToGroup(groupId, message2); + + final Message messageFromRegion1 = messageStore1.pollMessageFromGroup(groupId); + final Message messageFromRegion2 = messageStore2.pollMessageFromGroup(groupId); + + assertNotNull(messageFromRegion1); + assertNotNull(messageFromRegion2); + + LOG.info("messageFromRegion1: " + messageFromRegion1.getHeaders().getId() + "; Sequence #: " + messageFromRegion1.getHeaders().getSequenceNumber()); + LOG.info("messageFromRegion2: " + messageFromRegion2.getHeaders().getId() + "; Sequence #: " + messageFromRegion2.getHeaders().getSequenceNumber()); + + assertEquals(Integer.valueOf(1), (Integer) messageFromRegion1.getHeaders().get(MessageHeaders.SEQUENCE_NUMBER)); + assertEquals(Integer.valueOf(2), (Integer) messageFromRegion2.getHeaders().get(MessageHeaders.SEQUENCE_NUMBER)); + + } } diff --git a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/mysql/CommonMySql-context.xml b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/mysql/CommonMySql-context.xml new file mode 100644 index 0000000000..968cc852d5 --- /dev/null +++ b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/mysql/CommonMySql-context.xml @@ -0,0 +1,21 @@ + + + + + + + + + + + + + + + + + diff --git a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/mysql/MySqlJdbcMessageStoreMultipleChannelTests-context.xml b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/mysql/MySqlJdbcMessageStoreMultipleChannelTests-context.xml new file mode 100644 index 0000000000..46144dec2d --- /dev/null +++ b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/mysql/MySqlJdbcMessageStoreMultipleChannelTests-context.xml @@ -0,0 +1,99 @@ + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + diff --git a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/mysql/MySqlJdbcMessageStoreMultipleChannelTests.java b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/mysql/MySqlJdbcMessageStoreMultipleChannelTests.java new file mode 100644 index 0000000000..315db03992 --- /dev/null +++ b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/mysql/MySqlJdbcMessageStoreMultipleChannelTests.java @@ -0,0 +1,175 @@ +/* + * Copyright 2002-2013 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. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.integration.jdbc.mysql; + +import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertTrue; + +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; + +import javax.sql.DataSource; + +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; +import org.junit.After; +import org.junit.Before; +import org.junit.Ignore; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.integration.Message; +import org.springframework.integration.MessageChannel; +import org.springframework.integration.channel.QueueChannel; +import org.springframework.integration.support.MessageBuilder; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.test.annotation.DirtiesContext; +import org.springframework.test.annotation.DirtiesContext.ClassMode; +import org.springframework.test.context.ContextConfiguration; +import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; +import org.springframework.transaction.PlatformTransactionManager; +import org.springframework.transaction.TransactionStatus; +import org.springframework.transaction.support.TransactionCallback; +import org.springframework.transaction.support.TransactionTemplate; + +/** + * This test was created to reproduce INT-2980. + * + * @author Gunnar Hillert + * + */ +@ContextConfiguration +@RunWith(SpringJUnit4ClassRunner.class) +@DirtiesContext(classMode=ClassMode.AFTER_EACH_TEST_METHOD) +@Ignore +public class MySqlJdbcMessageStoreMultipleChannelTests { + + private static final Log LOG = LogFactory.getLog(MySqlJdbcMessageStoreMultipleChannelTests.class); + + private static final CountDownLatch countDownLatch1 = new CountDownLatch(1); + private static final CountDownLatch countDownLatch2 = new CountDownLatch(1); + + private static AtomicBoolean success = new AtomicBoolean(true); + + @Autowired + @Qualifier("requestChannel") + private MessageChannel requestChannel; + + @Autowired + @Qualifier("errorChannel") + private QueueChannel errorChannel; + + @Autowired + private PlatformTransactionManager transactionManager; + + private JdbcTemplate jdbcTemplate; + + @Autowired + private DataSource dataSource; + + @Before + public void beforeTest() { + this.jdbcTemplate = new JdbcTemplate(dataSource); + } + + @After + public void afterTest() { + new TransactionTemplate(this.transactionManager).execute(new TransactionCallback() { + public Void doInTransaction(TransactionStatus status) { + final int deletedGroupToMessageRows = jdbcTemplate.update("delete from INT_GROUP_TO_MESSAGE"); + final int deletedMessages = jdbcTemplate.update("delete from INT_MESSAGE"); + final int deletedMessageGroups = jdbcTemplate.update("delete from INT_MESSAGE_GROUP"); + + LOG.info(String.format("Cleaning Database - Deleted Messages: %s, " + + "Deleted GroupToMessage Rows: %s, Deleted Message Groups: %s", + deletedMessages, deletedGroupToMessageRows, deletedMessageGroups)); + + return null; + } + }); + } + + @Test + public void testSendAndActivateTransactionalSend() throws Exception { + + new TransactionTemplate(this.transactionManager).execute(new TransactionCallback() { + public Void doInTransaction(TransactionStatus status) { + requestChannel.send(MessageBuilder.withPayload("Hello ").build()); + return null; + } + }); + + assertTrue("countDownLatch1 was " + countDownLatch1.getCount(), countDownLatch1.await(10000, TimeUnit.MILLISECONDS)); + assertTrue("countDownLatch2 was " + countDownLatch2.getCount(), countDownLatch2.await(10000, TimeUnit.MILLISECONDS)); + + assertTrue("Wrong Sequence Number handled.", success.get()); + assertNull(errorChannel.receive(0)); + } + + public static class Splitter { + + public Splitter() { + super(); + } + + public List duplicate(Message message) { + ArrayList res = new ArrayList(); + res.add(message); + res.add(message); + + System.out.println("Split Complete"); + + return res; + } + } + + public static class ServiceActivator { + + public ServiceActivator() { + super(); + } + + public void first(Message message ) { + + int sequenceNumber = message.getHeaders().getSequenceNumber(); + + LOG.info("First handling sequence number: " + sequenceNumber + "; Message ID: " + message.getHeaders().getId()); + + if (sequenceNumber != 1) { + success.set(false); + } + + countDownLatch1.countDown(); + } + + public void second(Message message ) { + + int sequenceNumber = message.getHeaders().getSequenceNumber(); + LOG.info("Second handling sequence number: " + sequenceNumber + "; Message ID: " + message.getHeaders().getId()); + + if (sequenceNumber != 2) { + success.set(false); + } + + countDownLatch2.countDown(); + } + } +} diff --git a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/mysql/MySqlJdbcMessageStoreTests-context.xml b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/mysql/MySqlJdbcMessageStoreTests-context.xml new file mode 100644 index 0000000000..968cc852d5 --- /dev/null +++ b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/mysql/MySqlJdbcMessageStoreTests-context.xml @@ -0,0 +1,21 @@ + + + + + + + + + + + + + + + + + diff --git a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/mysql/MySqlJdbcMessageStoreTests.java b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/mysql/MySqlJdbcMessageStoreTests.java new file mode 100644 index 0000000000..091bffba17 --- /dev/null +++ b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/mysql/MySqlJdbcMessageStoreTests.java @@ -0,0 +1,520 @@ +/* + * Copyright 2002-2013 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. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.integration.jdbc.mysql; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertNotSame; +import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertSame; +import static org.junit.Assert.assertThat; +import static org.junit.Assert.assertTrue; +import static org.springframework.integration.test.matcher.PayloadAndHeaderMatcher.sameExceptIgnorableHeaders; + +import java.io.BufferedReader; +import java.io.IOException; +import java.io.InputStream; +import java.io.InputStreamReader; +import java.io.OutputStream; +import java.util.Properties; +import java.util.UUID; + +import javax.sql.DataSource; + +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; +import org.junit.After; +import org.junit.Before; +import org.junit.Ignore; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.core.serializer.Deserializer; +import org.springframework.core.serializer.Serializer; +import org.springframework.integration.Message; +import org.springframework.integration.MessageHeaders; +import org.springframework.integration.channel.DirectChannel; +import org.springframework.integration.history.MessageHistory; +import org.springframework.integration.jdbc.JdbcMessageStore; +import org.springframework.integration.message.GenericMessage; +import org.springframework.integration.store.MessageGroup; +import org.springframework.integration.store.MessageGroupStore; +import org.springframework.integration.store.MessageGroupStore.MessageGroupCallback; +import org.springframework.integration.support.MessageBuilder; +import org.springframework.integration.util.UUIDConverter; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.test.annotation.DirtiesContext; +import org.springframework.test.annotation.DirtiesContext.ClassMode; +import org.springframework.test.annotation.Repeat; +import org.springframework.test.annotation.Rollback; +import org.springframework.test.context.ContextConfiguration; +import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; +import org.springframework.transaction.PlatformTransactionManager; +import org.springframework.transaction.TransactionStatus; +import org.springframework.transaction.annotation.Transactional; +import org.springframework.transaction.support.TransactionCallback; +import org.springframework.transaction.support.TransactionTemplate; + +/** + * Based on the test for Derby: + * + * {@link org.springframework.integration.jdbc.JdbcMessageStoreTests} + * + * This tests requires at least MySql 5.6.4 as it uses the fractional second support + * in that version. For more information, please see: + * + * http://dev.mysql.com/doc/refman/5.6/en/fractional-seconds.html + * + * Also, please make sure you are using the respective DDL scripts: + * + * schema-mysql-5_6_4.sql + * + * @author Gunnar Hillert + */ +@ContextConfiguration +@RunWith(SpringJUnit4ClassRunner.class) +@DirtiesContext(classMode=ClassMode.AFTER_EACH_TEST_METHOD) +@Ignore +public class MySqlJdbcMessageStoreTests { + + private static final Log LOG = LogFactory.getLog(MySqlJdbcMessageStoreTests.class); + + @Autowired + private DataSource dataSource; + + private JdbcMessageStore messageStore; + + @Autowired + private PlatformTransactionManager transactionManager; + + @Before + public void init() { + messageStore = new JdbcMessageStore(dataSource); + messageStore.setRegion("JdbcMessageStoreTests"); + } + + @After + public void afterTest() { + final JdbcTemplate jdbcTemplate = new JdbcTemplate(dataSource); + new TransactionTemplate(this.transactionManager).execute(new TransactionCallback() { + public Void doInTransaction(TransactionStatus status) { + final int deletedGroupToMessageRows = jdbcTemplate.update("delete from INT_GROUP_TO_MESSAGE"); + final int deletedMessages = jdbcTemplate.update("delete from INT_MESSAGE"); + final int deletedMessageGroups = jdbcTemplate.update("delete from INT_MESSAGE_GROUP"); + + LOG.info(String.format("Cleaning Database - Deleted Messages: %s, " + + "Deleted GroupToMessage Rows: %s, Deleted Message Groups: %s", + deletedMessages, deletedGroupToMessageRows, deletedMessageGroups)); + return null; + } + }); + } + + @Test + @Transactional + public void testGetNonExistent() throws Exception { + Message result = messageStore.getMessage(UUID.randomUUID()); + assertNull(result); + } + + @Test + @Transactional + public void testAddAndGet() throws Exception { + Message message = MessageBuilder.withPayload("foo").build(); + Message saved = messageStore.addMessage(message); + assertNotNull(messageStore.getMessage(message.getHeaders().getId())); + Message result = messageStore.getMessage(saved.getHeaders().getId()); + assertNotNull(result); + assertThat(saved, sameExceptIgnorableHeaders(result)); + assertNotNull(result.getHeaders().get(JdbcMessageStore.SAVED_KEY)); + assertNotNull(result.getHeaders().get(JdbcMessageStore.CREATED_DATE_KEY)); + } + + @Test + @Transactional + public void testWithMessageHistory() throws Exception{ + + Message message = new GenericMessage("Hello"); + DirectChannel fooChannel = new DirectChannel(); + fooChannel.setBeanName("fooChannel"); + DirectChannel barChannel = new DirectChannel(); + barChannel.setBeanName("barChannel"); + + message = MessageHistory.write(message, fooChannel); + message = MessageHistory.write(message, barChannel); + messageStore.addMessage(message); + message = messageStore.getMessage(message.getHeaders().getId()); + MessageHistory messageHistory = MessageHistory.read(message); + assertNotNull(messageHistory); + assertEquals(2, messageHistory.size()); + Properties fooChannelHistory = messageHistory.get(0); + assertEquals("fooChannel", fooChannelHistory.get("name")); + assertEquals("channel", fooChannelHistory.get("type")); + } + + @Test + @Transactional + public void testSize() throws Exception { + Message message = MessageBuilder.withPayload("foo").build(); + messageStore.addMessage(message); + assertEquals(1, messageStore.getMessageCount()); + } + + @Test + @Transactional + public void testSerializer() throws Exception { + // N.B. these serializers are not realistic (just for test purposes) + messageStore.setSerializer(new Serializer>() { + public void serialize(Message object, OutputStream outputStream) throws IOException { + outputStream.write(((Message) object).getPayload().toString().getBytes()); + outputStream.flush(); + } + }); + messageStore.setDeserializer(new Deserializer>() { + public GenericMessage deserialize(InputStream inputStream) throws IOException { + BufferedReader reader = new BufferedReader(new InputStreamReader(inputStream)); + return new GenericMessage(reader.readLine()); + } + }); + Message message = MessageBuilder.withPayload("foo").build(); + Message saved = messageStore.addMessage(message); + assertNotNull(messageStore.getMessage(message.getHeaders().getId())); + Message result = messageStore.getMessage(saved.getHeaders().getId()); + assertNotNull(result); + assertEquals("foo", result.getPayload()); + } + + @Test + @Transactional + public void testAddAndGetWithDifferentRegion() throws Exception { + Message message = MessageBuilder.withPayload("foo").build(); + Message saved = messageStore.addMessage(message); + messageStore.setRegion("FOO"); + Message result = messageStore.getMessage(saved.getHeaders().getId()); + assertNull(result); + } + + @Test + @Transactional + public void testAddAndUpdate() throws Exception { + Message message = MessageBuilder.withPayload("foo").setCorrelationId("X").build(); + message = messageStore.addMessage(message); + message = MessageBuilder.fromMessage(message).setCorrelationId("Y").build(); + message = messageStore.addMessage(message); + assertEquals("Y", messageStore.getMessage(message.getHeaders().getId()).getHeaders().getCorrelationId()); + } + + @Test + @Transactional + public void testAddAndUpdateAlreadySaved() throws Exception { + Message message = MessageBuilder.withPayload("foo").build(); + message = messageStore.addMessage(message); + Message result = messageStore.addMessage(message); + assertSame(message, result); + } + + @Test + @Transactional + public void testAddAndUpdateAlreadySavedAndCopied() throws Exception { + Message message = MessageBuilder.withPayload("foo").build(); + Message saved = messageStore.addMessage(message); + Message copy = MessageBuilder.fromMessage(saved).build(); + Message result = messageStore.addMessage(copy); + assertEquals(copy, result); + assertEquals(saved, result); + assertNotNull(messageStore.getMessage(saved.getHeaders().getId())); + } + + @Test + @Transactional + public void testAddAndUpdateWithChange() throws Exception { + Message message = MessageBuilder.withPayload("foo").build(); + Message saved = messageStore.addMessage(message); + Message copy = MessageBuilder.fromMessage(saved).setHeader("newHeader", 1).build(); + Message result = messageStore.addMessage(copy); + assertNotSame(saved, result); + assertThat(saved, sameExceptIgnorableHeaders(result, JdbcMessageStore.CREATED_DATE_KEY, "newHeader")); + assertNotNull(messageStore.getMessage(saved.getHeaders().getId())); + } + + @Test + @Transactional + public void testAddAndRemoveMessageGroup() throws Exception { + Message message = MessageBuilder.withPayload("foo").build(); + message = messageStore.addMessage(message); + assertNotNull(messageStore.removeMessage(message.getHeaders().getId())); + } + + @Test + @Transactional + public void testAddAndGetMessageGroup() throws Exception { + String groupId = "X"; + Message message = MessageBuilder.withPayload("foo").setCorrelationId(groupId).build(); + long now = System.currentTimeMillis(); + messageStore.addMessageToGroup(groupId, message); + MessageGroup group = messageStore.getMessageGroup(groupId); + assertEquals(1, group.size()); + assertTrue("Timestamp too early: " + group.getTimestamp() + "<" + now, group.getTimestamp() >= now); + } + + @Test + @Transactional + public void testAddAndRemoveMessageFromMessageGroup() throws Exception { + String groupId = "X"; + Message message = MessageBuilder.withPayload("foo").setCorrelationId(groupId).build(); + messageStore.addMessageToGroup(groupId, message); + messageStore.removeMessageFromGroup(groupId, message); + MessageGroup group = messageStore.getMessageGroup(groupId); + assertEquals(0, group.size()); + } + + @Test + @Transactional + public void testRemoveMessageGroup() throws Exception { + JdbcTemplate template = new JdbcTemplate(dataSource); + template.afterPropertiesSet(); + String groupId = "X"; + + Message message = MessageBuilder.withPayload("foo").setCorrelationId(groupId).build(); + messageStore.addMessageToGroup(groupId, message); + messageStore.removeMessageGroup(groupId); + MessageGroup group = messageStore.getMessageGroup(groupId); + assertEquals(0, group.size()); + + String uuidGroupId = UUIDConverter.getUUID(groupId).toString(); + assertTrue(template.queryForList( + "SELECT * from INT_GROUP_TO_MESSAGE where GROUP_KEY = '" + uuidGroupId + "'").size() == 0); + } + + @Test + @Transactional + public void testCompleteMessageGroup() throws Exception { + String groupId = "X"; + Message message = MessageBuilder.withPayload("foo").setCorrelationId(groupId).build(); + messageStore.addMessageToGroup(groupId, message); + messageStore.completeGroup(groupId); + MessageGroup group = messageStore.getMessageGroup(groupId); + assertTrue(group.isComplete()); + assertEquals(1, group.size()); + } + + @Test + @Transactional + public void testUpdateLastReleasedSequence() throws Exception { + String groupId = "X"; + Message message = MessageBuilder.withPayload("foo").setCorrelationId(groupId).build(); + messageStore.addMessageToGroup(groupId, message); + messageStore.setLastReleasedSequenceNumberForGroup(groupId, 5); + MessageGroup group = messageStore.getMessageGroup(groupId); + assertEquals(5, group.getLastReleasedMessageSequenceNumber()); + } + + @Test + @Transactional + public void testMessageGroupCount() throws Exception { + String groupId = "X"; + Message message = MessageBuilder.withPayload("foo").build(); + messageStore.addMessageToGroup(groupId, message); + assertEquals(1, messageStore.getMessageGroupCount()); + } + + @Test + @Transactional + public void testMessageGroupSizes() throws Exception { + String groupId = "X"; + Message message = MessageBuilder.withPayload("foo").build(); + messageStore.addMessageToGroup(groupId, message); + assertEquals(1, messageStore.getMessageCountForAllMessageGroups()); + } + + @Test + @Transactional + public void testOrderInMessageGroup() throws Exception { + String groupId = "X"; + + messageStore.addMessageToGroup(groupId, MessageBuilder.withPayload("foo").setCorrelationId(groupId).build()); + Thread.sleep(1); + messageStore.addMessageToGroup(groupId, MessageBuilder.withPayload("bar").setCorrelationId(groupId).build()); + MessageGroup group = messageStore.getMessageGroup(groupId); + assertEquals(2, group.size()); + assertEquals("foo", messageStore.pollMessageFromGroup(groupId).getPayload()); + assertEquals("bar", messageStore.pollMessageFromGroup(groupId).getPayload()); + } + + @Test + @Transactional + public void testExpireMessageGroupOnCreateOnly() throws Exception { + String groupId = "X"; + Message message = MessageBuilder.withPayload("foo").setCorrelationId(groupId).build(); + messageStore.addMessageToGroup(groupId, message); + messageStore.registerMessageGroupExpiryCallback(new MessageGroupCallback() { + public void execute(MessageGroupStore messageGroupStore, MessageGroup group) { + messageGroupStore.removeMessageGroup(group.getGroupId()); + } + }); + Thread.sleep(1000); + messageStore.expireMessageGroups(2000); + MessageGroup group = messageStore.getMessageGroup(groupId); + assertEquals(1, group.size()); + messageStore.addMessageToGroup(groupId, MessageBuilder.withPayload("bar").setCorrelationId(groupId).build()); + Thread.sleep(2001); + messageStore.expireMessageGroups(2000); + group = messageStore.getMessageGroup(groupId); + assertEquals(0, group.size()); + } + + @Test + @Transactional + public void testExpireMessageGroupOnIdleOnly() throws Exception { + String groupId = "X"; + Message message = MessageBuilder.withPayload("foo").setCorrelationId(groupId).build(); + messageStore.setTimeoutOnIdle(true); + messageStore.addMessageToGroup(groupId, message); + messageStore.registerMessageGroupExpiryCallback(new MessageGroupCallback() { + public void execute(MessageGroupStore messageGroupStore, MessageGroup group) { + messageGroupStore.removeMessageGroup(group.getGroupId()); + } + }); + Thread.sleep(1000); + messageStore.expireMessageGroups(2000); + MessageGroup group = messageStore.getMessageGroup(groupId); + assertEquals(1, group.size()); + Thread.sleep(2000); + messageStore.addMessageToGroup(groupId, MessageBuilder.withPayload("bar").setCorrelationId(groupId).build()); + group = messageStore.getMessageGroup(groupId); + assertEquals(2, group.size()); + Thread.sleep(2000); + messageStore.expireMessageGroups(2000); + group = messageStore.getMessageGroup(groupId); + assertEquals(0, group.size()); + } + + @Test + @Transactional + public void testMessagePollingFromTheGroup() throws Exception { + + final String groupX = "X"; + + messageStore.addMessageToGroup(groupX, MessageBuilder.withPayload("foo").setCorrelationId(groupX).build()); + Thread.sleep(100); + messageStore.addMessageToGroup(groupX, MessageBuilder.withPayload("bar").setCorrelationId(groupX).build()); + Thread.sleep(100); + messageStore.addMessageToGroup(groupX, MessageBuilder.withPayload("baz").setCorrelationId(groupX).build()); + Thread.sleep(100); + messageStore.addMessageToGroup("Y", MessageBuilder.withPayload("barA").setCorrelationId(groupX).build()); + Thread.sleep(100); + messageStore.addMessageToGroup("Y", MessageBuilder.withPayload("bazA").setCorrelationId(groupX).build()); + Thread.sleep(100); + MessageGroup group = messageStore.getMessageGroup(groupX); + assertEquals(3, group.size()); + + Message message1 = messageStore.pollMessageFromGroup(groupX); + assertNotNull(message1); + assertEquals("foo", message1.getPayload()); + + group = messageStore.getMessageGroup(groupX); + assertEquals(2, group.size()); + + Message message2 = messageStore.pollMessageFromGroup(groupX); + assertNotNull(message2); + assertEquals("bar", message2.getPayload()); + + group = messageStore.getMessageGroup(groupX); + assertEquals(1, group.size()); + } + + @Test + @Transactional + @Rollback(false) + @Repeat(20) + public void testSameMessageToMultipleGroups() throws Exception { + + final String group1Id = "group1"; + final String group2Id = "group2"; + + final Message message = MessageBuilder.withPayload("foo").build(); + + final MessageBuilder builder1 = MessageBuilder.fromMessage(message); + final MessageBuilder builder2 = MessageBuilder.fromMessage(message); + + builder1.setSequenceNumber(1); + builder2.setSequenceNumber(2); + + final Message message1 = builder1.build(); + final Message message2 = builder2.build(); + + messageStore.addMessageToGroup(group1Id, message1); + messageStore.addMessageToGroup(group2Id, message2); + + final Message messageFromGroup1 = messageStore.pollMessageFromGroup(group1Id); + final Message messageFromGroup2 = messageStore.pollMessageFromGroup(group2Id); + + assertNotNull(messageFromGroup1); + assertNotNull(messageFromGroup2); + + LOG.info("messageFromGroup1: " + messageFromGroup1.getHeaders().getId() + "; Sequence #: " + messageFromGroup1.getHeaders().getSequenceNumber()); + LOG.info("messageFromGroup2: " + messageFromGroup2.getHeaders().getId() + "; Sequence #: " + messageFromGroup2.getHeaders().getSequenceNumber()); + + assertEquals(Integer.valueOf(1), (Integer) messageFromGroup1.getHeaders().get(MessageHeaders.SEQUENCE_NUMBER)); + assertEquals(Integer.valueOf(2), (Integer) messageFromGroup2.getHeaders().get(MessageHeaders.SEQUENCE_NUMBER)); + + } + + @Test + @Transactional + @Rollback(false) + @Repeat(20) + public void testSameMessageAndGroupToMultipleRegions() throws Exception { + + final String groupId = "myGroup"; + final String region1 = "region1"; + final String region2 = "region2"; + + final JdbcMessageStore messageStore1 = new JdbcMessageStore(dataSource); + messageStore1.setRegion(region1); + + final JdbcMessageStore messageStore2 = new JdbcMessageStore(dataSource); + messageStore1.setRegion(region2); + + final Message message = MessageBuilder.withPayload("foo").build(); + + final MessageBuilder builder1 = MessageBuilder.fromMessage(message); + final MessageBuilder builder2 = MessageBuilder.fromMessage(message); + + builder1.setSequenceNumber(1); + builder2.setSequenceNumber(2); + + final Message message1 = builder1.build(); + final Message message2 = builder2.build(); + + messageStore1.addMessageToGroup(groupId, message1); + messageStore2.addMessageToGroup(groupId, message2); + + final Message messageFromRegion1 = messageStore1.pollMessageFromGroup(groupId); + final Message messageFromRegion2 = messageStore2.pollMessageFromGroup(groupId); + + assertNotNull(messageFromRegion1); + assertNotNull(messageFromRegion2); + + LOG.info("messageFromRegion1: " + messageFromRegion1.getHeaders().getId() + "; Sequence #: " + messageFromRegion1.getHeaders().getSequenceNumber()); + LOG.info("messageFromRegion2: " + messageFromRegion2.getHeaders().getId() + "; Sequence #: " + messageFromRegion2.getHeaders().getSequenceNumber()); + + assertEquals(Integer.valueOf(1), (Integer) messageFromRegion1.getHeaders().get(MessageHeaders.SEQUENCE_NUMBER)); + assertEquals(Integer.valueOf(2), (Integer) messageFromRegion2.getHeaders().get(MessageHeaders.SEQUENCE_NUMBER)); + + } +} diff --git a/src/reference/docbook/jdbc.xml b/src/reference/docbook/jdbc.xml index 424df5048b..8ff1cb9361 100644 --- a/src/reference/docbook/jdbc.xml +++ b/src/reference/docbook/jdbc.xml @@ -334,6 +334,35 @@ and a prefix for the table names in the queries generated by the store. The table name prefix defaults to "INT_". + + + If you plan on using MySQL, + please use MySQL version 5.6.4 or higher, if + possible. Prior versions do not support fractional seconds + for temporal data types. Because of that, messages may not arrive in + the precise FIFO order when polling from such a MySQL Message Store. + + + Therefore, starting with Spring Integration 3.0, + we provide an additional set of DDL scripts for MySQL version + 5.6.4 or higher: + + + schema-drop-mysql-5_6_4.sql + schema-mysql-5_6_4.sql + + + For more information, please see: + + + + + + Also important, please ensure that you use an up-to-date version + of the JDBC driver for MySQL (Connector/J), e.g. version + 5.1.24 or higher. + +
Backing Message Channels diff --git a/src/reference/docbook/whats-new.xml b/src/reference/docbook/whats-new.xml index 7430c0e9ed..8dffebfff2 100644 --- a/src/reference/docbook/whats-new.xml +++ b/src/reference/docbook/whats-new.xml @@ -116,6 +116,16 @@ files.
+
+ JDBC Message Store Improvements + + Spring Integration 3.0 adds a new set of DDL + scripts for MySQL version 5.6.4 and higher. + Now MySQL supports fractional + seconds and is thus improving the FIFO ordering when + polling from a MySQL-based Message Store. For more information, + please see . + +
-