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
This commit is contained in:
committed by
Gary Russell
parent
749896b812
commit
88165c8881
@@ -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"
|
||||
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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<Message<?>> messages = jdbcTemplate.query(getQuery(Query.POLL_FROM_GROUP), new Object[] { groupIdKey, region }, mapper);
|
||||
List<Message<?>> 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);
|
||||
|
||||
@@ -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 ;
|
||||
@@ -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;
|
||||
14
spring-integration-jdbc/src/main/sql/mysql-5_6_4.properties
Normal file
14
spring-integration-jdbc/src/main/sql/mysql-5_6_4.properties
Normal file
@@ -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
|
||||
4
spring-integration-jdbc/src/main/sql/mysql-5_6_4.vpp
Normal file
4
spring-integration-jdbc/src/main/sql/mysql-5_6_4.vpp
Normal file
@@ -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
|
||||
@@ -6,7 +6,7 @@
|
||||
http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd">
|
||||
|
||||
<jdbc:embedded-database id="dataSource" type="DERBY"/>
|
||||
|
||||
|
||||
<jdbc:initialize-database data-source="dataSource" ignore-failures="DROPS">
|
||||
<jdbc:script location="${int.drop.script}" />
|
||||
<jdbc:script location="${int.schema.script}" />
|
||||
|
||||
@@ -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<String> message = MessageBuilder.withPayload("foo").build();
|
||||
|
||||
final MessageBuilder<String> builder1 = MessageBuilder.fromMessage(message);
|
||||
final MessageBuilder<String> 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<String> message = MessageBuilder.withPayload("foo").build();
|
||||
|
||||
final MessageBuilder<String> builder1 = MessageBuilder.fromMessage(message);
|
||||
final MessageBuilder<String> 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));
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,21 @@
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<beans xmlns="http://www.springframework.org/schema/beans"
|
||||
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
|
||||
xmlns:jdbc="http://www.springframework.org/schema/jdbc"
|
||||
xsi:schemaLocation="http://www.springframework.org/schema/jdbc http://www.springframework.org/schema/jdbc/spring-jdbc.xsd
|
||||
http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd">
|
||||
|
||||
<bean id="dataSource" class="org.apache.commons.dbcp.BasicDataSource" destroy-method="close">
|
||||
<property name="driverClassName" value="com.mysql.jdbc.Driver"/>
|
||||
<property name="url" value="jdbc:mysql://localhost:3306/int30"/>
|
||||
<property name="username" value="root"/>
|
||||
<property name="password" value="root"/>
|
||||
<property name="maxActive" value="10"/>
|
||||
<property name="defaultAutoCommit" value="false"/>
|
||||
</bean>
|
||||
|
||||
<bean id="transactionManager" class="org.springframework.jdbc.datasource.DataSourceTransactionManager">
|
||||
<property name="dataSource" ref="dataSource" />
|
||||
</bean>
|
||||
|
||||
</beans>
|
||||
@@ -0,0 +1,99 @@
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<beans xmlns="http://www.springframework.org/schema/beans"
|
||||
xmlns:beans="http://www.springframework.org/schema/beans"
|
||||
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
|
||||
xmlns:int="http://www.springframework.org/schema/integration"
|
||||
xmlns:jdbc="http://www.springframework.org/schema/jdbc"
|
||||
xmlns:int-jdbc="http://www.springframework.org/schema/integration/jdbc"
|
||||
xmlns:task="http://www.springframework.org/schema/task"
|
||||
xmlns:tx="http://www.springframework.org/schema/tx"
|
||||
xsi:schemaLocation="http://www.springframework.org/schema/integration http://www.springframework.org/schema/integration/spring-integration.xsd
|
||||
http://www.springframework.org/schema/task http://www.springframework.org/schema/task/spring-task.xsd
|
||||
http://www.springframework.org/schema/jdbc http://www.springframework.org/schema/jdbc/spring-jdbc.xsd
|
||||
http://www.springframework.org/schema/tx http://www.springframework.org/schema/tx/spring-tx.xsd
|
||||
http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd
|
||||
http://www.springframework.org/schema/integration/jdbc http://www.springframework.org/schema/integration/jdbc/spring-integration-jdbc.xsd">
|
||||
|
||||
<tx:annotation-driven/>
|
||||
<int:annotation-config/>
|
||||
|
||||
<bean id="dataSource" class="org.apache.commons.dbcp.BasicDataSource" destroy-method="close">
|
||||
<property name="driverClassName" value="com.mysql.jdbc.Driver"/>
|
||||
<property name="url" value="jdbc:mysql://localhost:3306/int30"/>
|
||||
<property name="username" value="root"/>
|
||||
<property name="password" value="root"/>
|
||||
<property name="maxActive" value="10"/>
|
||||
<property name="defaultAutoCommit" value="false"/>
|
||||
</bean>
|
||||
|
||||
<int-jdbc:message-store id="messageStore" data-source="dataSource" region="MessageStoreMultipleChannelTests"/>
|
||||
|
||||
<bean id="placeholderProperties" class="org.springframework.beans.factory.config.PropertyPlaceholderConfigurer">
|
||||
<property name="location" value="classpath:int-${ENVIRONMENT:derby}.properties" />
|
||||
<property name="systemPropertiesModeName" value="SYSTEM_PROPERTIES_MODE_OVERRIDE" />
|
||||
<property name="ignoreUnresolvablePlaceholders" value="true" />
|
||||
<property name="order" value="1" />
|
||||
</bean>
|
||||
|
||||
<bean id="transactionManager" class="org.springframework.jdbc.datasource.DataSourceTransactionManager">
|
||||
<property name="dataSource" ref="dataSource" />
|
||||
</bean>
|
||||
|
||||
<task:executor id="pool" pool-size="1"
|
||||
queue-capacity="20" keep-alive="120"/>
|
||||
|
||||
<int:poller id="defaultPoller" default="true" fixed-rate="2000"
|
||||
max-messages-per-poll="1" task-executor="pool">
|
||||
<int:transactional />
|
||||
</int:poller>
|
||||
|
||||
<!-- Start of Flow -->
|
||||
|
||||
<int:channel id="requestChannel">
|
||||
<int:queue message-store="messageStore"/>
|
||||
</int:channel>
|
||||
|
||||
<bean id="splitter" class="org.springframework.integration.jdbc.mysql.MySqlJdbcMessageStoreMultipleChannelTests$Splitter"/>
|
||||
|
||||
<int:splitter input-channel="requestChannel" output-channel="afterSplitChannel"
|
||||
ref="splitter" method="duplicate" apply-sequence="true">
|
||||
</int:splitter>
|
||||
|
||||
<!-- Must not use the message store -->
|
||||
<int:channel id="afterSplitChannel"/>
|
||||
|
||||
<int:header-value-router input-channel="afterSplitChannel" header-name="sequenceNumber">
|
||||
<int:mapping value="1" channel="firstChannel" />
|
||||
<int:mapping value="2" channel="secondChannel" />
|
||||
</int:header-value-router>
|
||||
|
||||
<int:channel id="firstChannel">
|
||||
<int:queue message-store="messageStore"/>
|
||||
<int:interceptors>
|
||||
<int:wire-tap channel="loggit"/>
|
||||
</int:interceptors>
|
||||
</int:channel>
|
||||
|
||||
<int:channel id="secondChannel">
|
||||
<int:queue message-store="messageStore"/>
|
||||
<int:interceptors>
|
||||
<int:wire-tap channel="loggit"/>
|
||||
</int:interceptors>
|
||||
</int:channel>
|
||||
|
||||
<int:logging-channel-adapter id="loggit" log-full-message="true"/>
|
||||
|
||||
<beans:bean id="serviceActivator"
|
||||
class="org.springframework.integration.jdbc.mysql.MySqlJdbcMessageStoreMultipleChannelTests$ServiceActivator"/>
|
||||
|
||||
<int:service-activator id="serviceActivator1" input-channel="firstChannel"
|
||||
ref="serviceActivator" method="first" />
|
||||
|
||||
<int:service-activator id="serviceActivator2" input-channel="secondChannel"
|
||||
ref="serviceActivator" method="second" />
|
||||
|
||||
<int:channel id="errorChannel">
|
||||
<int:queue/>
|
||||
</int:channel>
|
||||
|
||||
</beans>
|
||||
@@ -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<Void>() {
|
||||
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<Void>() {
|
||||
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<Object> duplicate(Message<?> message) {
|
||||
ArrayList<Object> res = new ArrayList<Object>();
|
||||
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();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,21 @@
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<beans xmlns="http://www.springframework.org/schema/beans"
|
||||
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
|
||||
xmlns:jdbc="http://www.springframework.org/schema/jdbc"
|
||||
xsi:schemaLocation="http://www.springframework.org/schema/jdbc http://www.springframework.org/schema/jdbc/spring-jdbc.xsd
|
||||
http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd">
|
||||
|
||||
<bean id="dataSource" class="org.apache.commons.dbcp.BasicDataSource" destroy-method="close">
|
||||
<property name="driverClassName" value="com.mysql.jdbc.Driver"/>
|
||||
<property name="url" value="jdbc:mysql://localhost:3306/int30"/>
|
||||
<property name="username" value="root"/>
|
||||
<property name="password" value="root"/>
|
||||
<property name="maxActive" value="10"/>
|
||||
<property name="defaultAutoCommit" value="false"/>
|
||||
</bean>
|
||||
|
||||
<bean id="transactionManager" class="org.springframework.jdbc.datasource.DataSourceTransactionManager">
|
||||
<property name="dataSource" ref="dataSource" />
|
||||
</bean>
|
||||
|
||||
</beans>
|
||||
@@ -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<Void>() {
|
||||
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<String> message = MessageBuilder.withPayload("foo").build();
|
||||
Message<String> 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<String>("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<String> 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<Message<?>>() {
|
||||
public void serialize(Message<?> object, OutputStream outputStream) throws IOException {
|
||||
outputStream.write(((Message<?>) object).getPayload().toString().getBytes());
|
||||
outputStream.flush();
|
||||
}
|
||||
});
|
||||
messageStore.setDeserializer(new Deserializer<GenericMessage<String>>() {
|
||||
public GenericMessage<String> deserialize(InputStream inputStream) throws IOException {
|
||||
BufferedReader reader = new BufferedReader(new InputStreamReader(inputStream));
|
||||
return new GenericMessage<String>(reader.readLine());
|
||||
}
|
||||
});
|
||||
Message<String> message = MessageBuilder.withPayload("foo").build();
|
||||
Message<String> 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<String> message = MessageBuilder.withPayload("foo").build();
|
||||
Message<String> saved = messageStore.addMessage(message);
|
||||
messageStore.setRegion("FOO");
|
||||
Message<?> result = messageStore.getMessage(saved.getHeaders().getId());
|
||||
assertNull(result);
|
||||
}
|
||||
|
||||
@Test
|
||||
@Transactional
|
||||
public void testAddAndUpdate() throws Exception {
|
||||
Message<String> 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<String> message = MessageBuilder.withPayload("foo").build();
|
||||
message = messageStore.addMessage(message);
|
||||
Message<String> result = messageStore.addMessage(message);
|
||||
assertSame(message, result);
|
||||
}
|
||||
|
||||
@Test
|
||||
@Transactional
|
||||
public void testAddAndUpdateAlreadySavedAndCopied() throws Exception {
|
||||
Message<String> message = MessageBuilder.withPayload("foo").build();
|
||||
Message<String> saved = messageStore.addMessage(message);
|
||||
Message<String> copy = MessageBuilder.fromMessage(saved).build();
|
||||
Message<String> result = messageStore.addMessage(copy);
|
||||
assertEquals(copy, result);
|
||||
assertEquals(saved, result);
|
||||
assertNotNull(messageStore.getMessage(saved.getHeaders().getId()));
|
||||
}
|
||||
|
||||
@Test
|
||||
@Transactional
|
||||
public void testAddAndUpdateWithChange() throws Exception {
|
||||
Message<String> message = MessageBuilder.withPayload("foo").build();
|
||||
Message<String> saved = messageStore.addMessage(message);
|
||||
Message<String> copy = MessageBuilder.fromMessage(saved).setHeader("newHeader", 1).build();
|
||||
Message<String> 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<String> 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<String> 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<String> 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<String> 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<String> 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<String> 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<String> message = MessageBuilder.withPayload("foo").build();
|
||||
messageStore.addMessageToGroup(groupId, message);
|
||||
assertEquals(1, messageStore.getMessageGroupCount());
|
||||
}
|
||||
|
||||
@Test
|
||||
@Transactional
|
||||
public void testMessageGroupSizes() throws Exception {
|
||||
String groupId = "X";
|
||||
Message<String> 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<String> 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<String> 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<String> message = MessageBuilder.withPayload("foo").build();
|
||||
|
||||
final MessageBuilder<String> builder1 = MessageBuilder.fromMessage(message);
|
||||
final MessageBuilder<String> 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<String> message = MessageBuilder.withPayload("foo").build();
|
||||
|
||||
final MessageBuilder<String> builder1 = MessageBuilder.fromMessage(message);
|
||||
final MessageBuilder<String> 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));
|
||||
|
||||
}
|
||||
}
|
||||
@@ -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_".
|
||||
</para>
|
||||
<note>
|
||||
<para>
|
||||
If you plan on using <emphasis role="bold">MySQL</emphasis>,
|
||||
please use MySQL version <emphasis>5.6.4</emphasis> or higher, if
|
||||
possible. Prior versions do not support <emphasis>fractional seconds</emphasis>
|
||||
for temporal data types. Because of that, messages may not arrive in
|
||||
the precise FIFO order when polling from such a MySQL Message Store.
|
||||
</para>
|
||||
<para>
|
||||
Therefore, starting with <emphasis>Spring Integration 3.0</emphasis>,
|
||||
we provide an additional set of DDL scripts for MySQL version
|
||||
<emphasis>5.6.4</emphasis> or higher:
|
||||
</para>
|
||||
<itemizedlist>
|
||||
<listitem>schema-drop-mysql-5_6_4.sql</listitem>
|
||||
<listitem>schema-mysql-5_6_4.sql</listitem>
|
||||
</itemizedlist>
|
||||
<para>
|
||||
For more information, please see:
|
||||
</para>
|
||||
<para>
|
||||
<ulink url="http://dev.mysql.com/doc/refman/5.6/en/fractional-seconds.html"></ulink>
|
||||
</para>
|
||||
<para>
|
||||
Also important, please ensure that you use an up-to-date version
|
||||
of the JDBC driver for MySQL (Connector/J), e.g. version
|
||||
<emphasis>5.1.24</emphasis> or higher.
|
||||
</para>
|
||||
</note>
|
||||
</section>
|
||||
<section id="jdbc-message-store-channels">
|
||||
<title>Backing Message Channels</title>
|
||||
|
||||
@@ -116,6 +116,16 @@
|
||||
files.
|
||||
</para>
|
||||
</section>
|
||||
<section id="3.0-jdbc-mysql-v5_6_4">
|
||||
<title>JDBC Message Store Improvements</title>
|
||||
<para>
|
||||
<emphasis>Spring Integration 3.0</emphasis> adds a new set of DDL
|
||||
scripts for <emphasis>MySQL</emphasis> version 5.6.4 and higher.
|
||||
Now <emphasis>MySQL</emphasis> supports <emphasis>fractional
|
||||
seconds</emphasis> and is thus improving the FIFO ordering when
|
||||
polling from a MySQL-based Message Store. For more information,
|
||||
please see <xref linkend="jdbc-message-store-generic"/>.
|
||||
</para>
|
||||
</section>
|
||||
</section>
|
||||
|
||||
</chapter>
|
||||
|
||||
Reference in New Issue
Block a user