From a5d4e3ab48e5bb1918cfa2cec28cc1962e6c17de Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Tue, 29 Oct 2013 20:50:20 +0200 Subject: [PATCH] INT-3181: Improve JdbcChannelMS usingIdCache doc JIRA: https://jira.springsource.org/browse/INT-3181 INT-3181: Fix JdbcChannelMS#pollMessageFromGroup * Add `doRemoveMessageFromGroup` and check its result. If it is `false` then just return `null`, not polled `Message` * Add `AbstractTxTimeoutMessageStoreTests#testInt3181ConcurrentPolling` test and implement it for some DBs * Make a note regarding MVCC as important in the Doc Minor Polishing --- .../jdbc/store/JdbcChannelMessageStore.java | 23 +++++++---- .../AbstractTxTimeoutMessageStoreTests.java | 30 +++++++++++++- .../DerbyTxTimeoutMessageStoreTests.java | 7 ++++ .../HsqlTxTimeoutMessageStoreTests.java | 7 ++++ .../MySqlTxTimeoutMessageStoreTests.java | 7 ++++ .../TxTimeoutMessageStoreTests-context.xml | 39 +++++++++++++++++++ src/reference/docbook/jdbc.xml | 9 +++++ 7 files changed, 113 insertions(+), 9 deletions(-) diff --git a/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/store/JdbcChannelMessageStore.java b/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/store/JdbcChannelMessageStore.java index 7e0e655f6d..51cd0cdb20 100644 --- a/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/store/JdbcChannelMessageStore.java +++ b/spring-integration-jdbc/src/main/java/org/springframework/integration/jdbc/store/JdbcChannelMessageStore.java @@ -16,7 +16,6 @@ package org.springframework.integration.jdbc.store; import java.sql.PreparedStatement; import java.sql.SQLException; import java.sql.Types; -import java.util.Collections; import java.util.HashMap; import java.util.HashSet; import java.util.Iterator; @@ -24,7 +23,6 @@ import java.util.List; import java.util.Map; import java.util.Set; import java.util.UUID; -import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReadWriteLock; import java.util.concurrent.locks.ReentrantReadWriteLock; @@ -33,6 +31,7 @@ import javax.sql.DataSource; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; + import org.springframework.beans.DirectFieldAccessor; import org.springframework.beans.factory.InitializingBean; import org.springframework.core.serializer.Deserializer; @@ -42,12 +41,12 @@ import org.springframework.core.serializer.support.SerializingConverter; import org.springframework.integration.Message; import org.springframework.integration.MessageHeaders; import org.springframework.integration.jdbc.JdbcMessageStore; +import org.springframework.integration.jdbc.store.channel.ChannelMessageStoreQueryProvider; import org.springframework.integration.jdbc.store.channel.DerbyChannelMessageStoreQueryProvider; import org.springframework.integration.jdbc.store.channel.MessageRowMapper; import org.springframework.integration.jdbc.store.channel.MySqlChannelMessageStoreQueryProvider; import org.springframework.integration.jdbc.store.channel.OracleChannelMessageStoreQueryProvider; import org.springframework.integration.jdbc.store.channel.PostgresChannelMessageStoreQueryProvider; -import org.springframework.integration.jdbc.store.channel.ChannelMessageStoreQueryProvider; import org.springframework.integration.store.AbstractMessageGroupStore; import org.springframework.integration.store.MessageGroup; import org.springframework.integration.store.MessageGroupStore; @@ -592,7 +591,9 @@ public class JdbcChannelMessageStore extends AbstractMessageGroupStore implement final Message polledMessage = this.doPollForMessage(key); if (polledMessage != null){ - this.removeMessageFromGroup(groupId, polledMessage); + if (!this.doRemoveMessageFromGroup(groupId, polledMessage)) { + return null; + } } return polledMessage; @@ -607,18 +608,26 @@ public class JdbcChannelMessageStore extends AbstractMessageGroupStore implement */ public MessageGroup removeMessageFromGroup(Object groupId, Message messageToRemove) { + this.doRemoveMessageFromGroup(groupId, messageToRemove); + + return getMessageGroup(groupId); + } + + private boolean doRemoveMessageFromGroup(Object groupId, Message messageToRemove) { final UUID id = messageToRemove.getHeaders().getId(); int updated = jdbcTemplate.update(getQuery(channelMessageStoreQueryProvider.getDeleteMessageQuery()), new Object[] { getKey(id), getKey(groupId), region }, new int[] { Types.VARCHAR, Types.VARCHAR, Types.VARCHAR }); - if (updated != 0) { + boolean result = updated != 0; + if (result) { logger.debug(String.format("Message with id '%s' was deleted.", id)); - } else { + } + else { logger.warn(String.format("Message with id '%s' was not deleted.", id)); } - return getMessageGroup(groupId); + return result; } /** diff --git a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/store/channel/AbstractTxTimeoutMessageStoreTests.java b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/store/channel/AbstractTxTimeoutMessageStoreTests.java index 712b5bff52..0982d30416 100644 --- a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/store/channel/AbstractTxTimeoutMessageStoreTests.java +++ b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/store/channel/AbstractTxTimeoutMessageStoreTests.java @@ -12,22 +12,27 @@ */ package org.springframework.integration.jdbc.store.channel; +import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertTrue; import java.util.concurrent.Callable; import java.util.concurrent.CompletionService; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorCompletionService; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; -import javax.sql.DataSource; +import java.util.concurrent.atomic.AtomicInteger; -import org.junit.Assert; +import javax.sql.DataSource; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; +import org.junit.Assert; + 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.jdbc.store.JdbcChannelMessageStore; @@ -62,8 +67,19 @@ abstract class AbstractTxTimeoutMessageStoreTests { protected TestService testService; @Autowired + @Qualifier("store") protected JdbcChannelMessageStore jdbcChannelMessageStore; + @Autowired + private MessageChannel first; + + @Autowired + private CountDownLatch successfulLatch; + + @Autowired + private AtomicInteger errorAtomicInteger; + + public void test() throws InterruptedException { int maxMessages = 10; @@ -157,4 +173,14 @@ abstract class AbstractTxTimeoutMessageStoreTests { assertTrue(executorService.awaitTermination(5, TimeUnit.SECONDS)); } + public void testInt3181ConcurrentPolling() throws InterruptedException { + for (int i = 0; i < 10; i++) { + this.first.send(new GenericMessage("test")); + } + + assertTrue(this.successfulLatch.await(5, TimeUnit.SECONDS)); + + assertEquals(0, errorAtomicInteger.get()); + } + } diff --git a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/store/channel/DerbyTxTimeoutMessageStoreTests.java b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/store/channel/DerbyTxTimeoutMessageStoreTests.java index f48528226d..b954331242 100644 --- a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/store/channel/DerbyTxTimeoutMessageStoreTests.java +++ b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/store/channel/DerbyTxTimeoutMessageStoreTests.java @@ -15,6 +15,7 @@ package org.springframework.integration.jdbc.store.channel; import org.junit.Ignore; import org.junit.Test; import org.junit.runner.RunWith; + import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; @@ -36,4 +37,10 @@ public class DerbyTxTimeoutMessageStoreTests extends AbstractTxTimeoutMessageSto super.test(); } + @Test + @Override + public void testInt3181ConcurrentPolling() throws InterruptedException { + super.testInt3181ConcurrentPolling(); + } + } diff --git a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/store/channel/HsqlTxTimeoutMessageStoreTests.java b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/store/channel/HsqlTxTimeoutMessageStoreTests.java index 803a31af9f..f88d4fec3d 100644 --- a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/store/channel/HsqlTxTimeoutMessageStoreTests.java +++ b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/store/channel/HsqlTxTimeoutMessageStoreTests.java @@ -16,6 +16,7 @@ import java.util.concurrent.ExecutionException; import org.junit.Test; import org.junit.runner.RunWith; + import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; @@ -40,4 +41,10 @@ public class HsqlTxTimeoutMessageStoreTests extends AbstractTxTimeoutMessageStor super.testInt2993IdCacheConcurrency(); } + @Test + @Override + public void testInt3181ConcurrentPolling() throws InterruptedException { + super.testInt3181ConcurrentPolling(); + } + } diff --git a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/store/channel/MySqlTxTimeoutMessageStoreTests.java b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/store/channel/MySqlTxTimeoutMessageStoreTests.java index b7f816d06f..7eaa27a7f7 100644 --- a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/store/channel/MySqlTxTimeoutMessageStoreTests.java +++ b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/store/channel/MySqlTxTimeoutMessageStoreTests.java @@ -15,6 +15,7 @@ package org.springframework.integration.jdbc.store.channel; import org.junit.Ignore; import org.junit.Test; import org.junit.runner.RunWith; + import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; @@ -35,4 +36,10 @@ public class MySqlTxTimeoutMessageStoreTests extends AbstractTxTimeoutMessageSto super.test(); } + @Test + @Override + public void testInt3181ConcurrentPolling() throws InterruptedException { + super.testInt3181ConcurrentPolling(); + } + } diff --git a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/store/channel/TxTimeoutMessageStoreTests-context.xml b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/store/channel/TxTimeoutMessageStoreTests-context.xml index b4f04f17f0..7276d88ed6 100644 --- a/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/store/channel/TxTimeoutMessageStoreTests-context.xml +++ b/spring-integration-jdbc/src/test/java/org/springframework/integration/jdbc/store/channel/TxTimeoutMessageStoreTests-context.xml @@ -54,5 +54,44 @@ + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + diff --git a/src/reference/docbook/jdbc.xml b/src/reference/docbook/jdbc.xml index e01bed77ec..7c067aa6e9 100644 --- a/src/reference/docbook/jdbc.xml +++ b/src/reference/docbook/jdbc.xml @@ -416,6 +416,7 @@ to configure the associated Poller with a TaskExecutor reference. + Keep in mind, though, that if you use a JDBC backed Message Channel and you are planning on polling the channel and consequently the message @@ -426,6 +427,13 @@ threads, may not materialize as expected. For example Apache Derby is problematic in that regard. + + To achieve better JDBC queue throughput, and avoid issues when different threads may poll the same + Message from the queue, it is important + to set the usingIdCache property of JdbcChannelMessageStore to true + when using databases that do not support MVCC: + + @@ -459,6 +467,7 @@ …]]> +
Initializing the Database