From 071c101c2806ad369b068127f86818ec757dc4b7 Mon Sep 17 00:00:00 2001 From: Mark Paluch Date: Wed, 12 Oct 2016 11:43:21 +0200 Subject: [PATCH] DATACASS-344 - Upgrade Cassandra Driver to 3.1.1. MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit We are now compatible with Datastax’ Cassandra Driver 3.1.1 and support the newly introduced configuration options MaxQueueSize for PoolingOptions. This change also removes test assertions for PoolTimeoutMillis as this option was deprecated and made unusable with 3.1.1. --- .travis.yml | 5 +- pom.xml | 2 +- .../config/PoolingOptionsFactoryBean.java | 57 +++++++++++++++++++ .../config/xml/CassandraCqlClusterParser.java | 3 + .../cassandra/config/spring-cql-1.0.xsd | 10 ++++ .../PoolingOptionsFactoryBeanUnitTests.java | 36 ++++++++++++ .../test/integration/CassandraRule.java | 12 ++-- ...paceCreatingXmlConfigIntegrationTests.java | 1 - .../config/xml/XmlConfigIntegrationTests.java | 3 +- .../xml/XmlConfigIntegrationTests-context.xml | 2 +- .../cassandra/config/spring-cassandra-1.0.xsd | 10 ++++ .../CassandraTypeMappingIntegrationTest.java | 15 ----- 12 files changed, 130 insertions(+), 26 deletions(-) diff --git a/.travis.yml b/.travis.yml index 82c93f55c..d5394aaac 100644 --- a/.travis.yml +++ b/.travis.yml @@ -4,6 +4,9 @@ jdk: env: matrix: - PROFILE=ci + - PROFILE=ci CASSANDRA_DRIVER_VERSION=3.0.3 + - PROFILE=ci CASSANDRA_DRIVER_VERSION=3.1.0 + - PROFILE=ci CASSANDRA_DRIVER_VERSION=3.1.1 - PROFILE=spring41-next - PROFILE=spring42 - PROFILE=spring42-next @@ -29,4 +32,4 @@ sudo: false before_install: - sed -i.bak -e 's|https://nexus.codehaus.org/snapshots/|https://oss.sonatype.org/content/repositories/codehaus-snapshots/|g' ~/.m2/settings.xml install: if [ ! -z ${CASSANDRA_VERSION} ]; then ./setup-cassandra.sh; fi; -script: mvn clean install -P${PROFILE} -Dmaven.javadoc.skip=true +script: mvn clean install -P${PROFILE} -Dcassandra-driver.version=${CASSANDRA_DRIVER_VERSION:-3.1.1} -Dmaven.javadoc.skip=true diff --git a/pom.xml b/pom.xml index 318a960e4..6813a1ac2 100644 --- a/pom.xml +++ b/pom.xml @@ -70,7 +70,7 @@ 17001 17000 3.9 - 3.1.0 + 3.1.1 spring-data-cassandra 1.0 2.16 diff --git a/spring-cql/src/main/java/org/springframework/cassandra/config/PoolingOptionsFactoryBean.java b/spring-cql/src/main/java/org/springframework/cassandra/config/PoolingOptionsFactoryBean.java index 55706dcaf..5b55237ac 100644 --- a/spring-cql/src/main/java/org/springframework/cassandra/config/PoolingOptionsFactoryBean.java +++ b/spring-cql/src/main/java/org/springframework/cassandra/config/PoolingOptionsFactoryBean.java @@ -15,6 +15,10 @@ */ package org.springframework.cassandra.config; + +import static org.springframework.util.ReflectionUtils.invokeMethod; + +import java.lang.reflect.Method; import java.util.concurrent.Executor; import org.springframework.beans.factory.FactoryBean; @@ -22,6 +26,7 @@ import org.springframework.beans.factory.InitializingBean; import com.datastax.driver.core.HostDistance; import com.datastax.driver.core.PoolingOptions; +import org.springframework.util.ReflectionUtils; /** * Spring {@link FactoryBean} for the Cassandra Java driver {@link PoolingOptions}. @@ -37,6 +42,18 @@ import com.datastax.driver.core.PoolingOptions; @SuppressWarnings("unused") public class PoolingOptionsFactoryBean implements FactoryBean, InitializingBean { + private static final PoolingOptions DEFAULT = new PoolingOptions(); + + private static final Method SET_MAX_QUEUE_SIZE; + private static final Method GET_MAX_QUEUE_SIZE; + + static { + SET_MAX_QUEUE_SIZE = ReflectionUtils + .findMethod(PoolingOptions.class, "setMaxQueueSize", int.class); + GET_MAX_QUEUE_SIZE = ReflectionUtils + .findMethod(PoolingOptions.class, "getMaxQueueSize"); + } + private Executor initializationExecutor; private Integer heartbeatIntervalSeconds; @@ -45,7 +62,12 @@ public class PoolingOptionsFactoryBean implements FactoryBean, I private Integer localMaxConnections; private Integer localMaxSimultaneousRequests; private Integer localMinSimultaneousRequests; + + // Deprecated since Cassandra Driver 3.1.1 private Integer poolTimeoutMilliseconds; + + // Available since Cassandra Driver 3.1.1 + private int maxQueueSize; private Integer remoteCoreConnections; private Integer remoteMaxConnections; private Integer remoteMaxSimultaneousRequests; @@ -78,6 +100,23 @@ public class PoolingOptionsFactoryBean implements FactoryBean, I if (poolTimeoutMilliseconds != null) { poolingOptions.setPoolTimeoutMillis(poolTimeoutMilliseconds); } + + if (!isDefaultMaxQueueSize() && SET_MAX_QUEUE_SIZE != null) { + invokeMethod(SET_MAX_QUEUE_SIZE, poolingOptions, maxQueueSize); + } + } + + private boolean isDefaultMaxQueueSize() { + + if(GET_MAX_QUEUE_SIZE != null){ + + Integer defaultMaxQueueSize = (Integer) invokeMethod(GET_MAX_QUEUE_SIZE, poolingOptions); + if(defaultMaxQueueSize.intValue() == maxQueueSize){ + return true; + } + } + + return false; } /* @@ -247,6 +286,24 @@ public class PoolingOptionsFactoryBean implements FactoryBean, I return poolTimeoutMilliseconds; } + /** + * Sets the maximum number of requests that get enqueued if no connection is available. + * + * @param maxQueueSize maximum number of requests that get enqueued if no connection is available. + */ + public void setMaxQueueSize(Integer maxQueueSize) { + this.maxQueueSize = maxQueueSize; + } + + /** + * Gets the maximum number of requests that get enqueued if no connection is available. + * + * @return the {@code maxQueueSize}. + */ + public Integer getMaxQueueSize() { + return maxQueueSize; + } + /** * Sets the core number of connections per host for the {@link HostDistance#LOCAL} scope. * diff --git a/spring-cql/src/main/java/org/springframework/cassandra/config/xml/CassandraCqlClusterParser.java b/spring-cql/src/main/java/org/springframework/cassandra/config/xml/CassandraCqlClusterParser.java index d73373789..d553a5ba6 100644 --- a/spring-cql/src/main/java/org/springframework/cassandra/config/xml/CassandraCqlClusterParser.java +++ b/spring-cql/src/main/java/org/springframework/cassandra/config/xml/CassandraCqlClusterParser.java @@ -48,6 +48,7 @@ import com.datastax.driver.core.SocketOptions; * @author Matthew T. Adams * @author David Webb * @author John Blum + * @author Mark Paluch */ public class CassandraCqlClusterParser extends AbstractBeanDefinitionParser { @@ -95,6 +96,7 @@ public class CassandraCqlClusterParser extends AbstractBeanDefinitionParser { addOptionalPropertyReference(builder, "hostStateListener", element, "host-state-listener-ref"); addOptionalPropertyReference(builder, "latencyTracker", element, "latency-tracker-ref"); addOptionalPropertyReference(builder, "loadBalancingPolicy", element, "load-balancing-policy-ref"); + addOptionalPropertyReference(builder, "nettyOptions", element, "netty-options-ref"); addOptionalPropertyReference(builder, "reconnectionPolicy", element, "reconnection-policy-ref"); addOptionalPropertyReference(builder, "retryPolicy", element, "retry-policy-ref"); addOptionalPropertyReference(builder, "speculativeExecutionPolicy", element, "speculative-execution-policy-ref"); @@ -138,6 +140,7 @@ public class CassandraCqlClusterParser extends AbstractBeanDefinitionParser { addOptionalPropertyValue(poolingOptionsBuilder, "heartbeatIntervalSeconds", element, "heartbeat-interval-seconds"); addOptionalPropertyValue(poolingOptionsBuilder, "idleTimeoutSeconds", element, "idle-timeout-seconds"); addOptionalPropertyValue(poolingOptionsBuilder, "poolTimeoutMilliseconds", element, "pool-timeout-milliseconds"); + addOptionalPropertyValue(poolingOptionsBuilder, "maxQueueSize", element, "max-queue-size"); // parse child elements for (Element subElement : DomUtils.getChildElements(element)) { diff --git a/spring-cql/src/main/resources/org/springframework/cassandra/config/spring-cql-1.0.xsd b/spring-cql/src/main/resources/org/springframework/cassandra/config/spring-cql-1.0.xsd index dd33c6cf0..530fdd7e7 100644 --- a/spring-cql/src/main/resources/org/springframework/cassandra/config/spring-cql-1.0.xsd +++ b/spring-cql/src/main/resources/org/springframework/cassandra/config/spring-cql-1.0.xsd @@ -263,6 +263,16 @@ LoadBalancingPolicy implementation. + + + + + DATACASS-176 * @see DATACASS-298 */ @@ -77,6 +82,7 @@ public class PoolingOptionsFactoryBeanUnitTests { /** * @see DATACASS-298 + * @see DATACASS-344 */ @Test public void setAndGetFactoryBeanProperties() { @@ -196,6 +202,36 @@ public class PoolingOptionsFactoryBeanUnitTests { verify(poolingOptionsSpy, never()).setNewConnectionThreshold(eq(HostDistance.REMOTE), eq(5)); } + /** + * @see DATACASS-344 + */ + @Test + public void afterPropertiesSetInitializesMaxQueueSize() throws Exception { + + Method setMaxQueueSize = ReflectionUtils + .findMethod(PoolingOptions.class, "setMaxQueueSize", int.class); + + Method getMaxQueueSize = ReflectionUtils + .findMethod(PoolingOptions.class, "getMaxQueueSize"); + + assumeNotNull(setMaxQueueSize); + + PoolingOptionsFactoryBean poolingOptionsFactoryBean = new PoolingOptionsFactoryBean() { + @Override + PoolingOptions newPoolingOptions() { + return poolingOptionsSpy; + } + }; + + poolingOptionsFactoryBean.setMaxQueueSize(1234); + + poolingOptionsFactoryBean.afterPropertiesSet(); + + assertThat(poolingOptionsFactoryBean.getObject(), is(sameInstance(poolingOptionsSpy))); + assertThat(poolingOptionsFactoryBean.getObjectType(), is(equalTo((Class) poolingOptionsSpy.getClass()))); + assertThat(invokeMethod(getMaxQueueSize, poolingOptionsSpy), is(equalTo((Object) 1234))); + } + /** * This particular test case is technically an integration test since it uses an actual instance of a DataStax Java * driver class type... {@link PoolingOptions}! The max values should be set before setting core values. Otherwise the diff --git a/spring-cql/src/test/java/org/springframework/cassandra/test/integration/CassandraRule.java b/spring-cql/src/test/java/org/springframework/cassandra/test/integration/CassandraRule.java index 97a527039..7182dfe1d 100644 --- a/spring-cql/src/test/java/org/springframework/cassandra/test/integration/CassandraRule.java +++ b/spring-cql/src/test/java/org/springframework/cassandra/test/integration/CassandraRule.java @@ -16,6 +16,7 @@ package org.springframework.cassandra.test.integration; +import static org.apache.cassandra.db.marshal.CompositeType.build; import static org.springframework.cassandra.test.integration.CassandraRule.InvocationMode.*; import java.util.ArrayList; @@ -317,11 +318,12 @@ public class CassandraRule extends ExternalResource { QueryOptions queryOptions = new QueryOptions(); queryOptions.setRefreshSchemaIntervalMillis(0); - cluster = new Cluster.Builder().addContactPoints(hostIp).// - withPort(port).// - withQueryOptions(queryOptions).// - withNettyOptions(FastShutdownNettyOptions.INSTANCE).// - build(); + cluster = new Cluster.Builder().addContactPoints(hostIp) // + .withPort(port) // + .withMaxSchemaAgreementWaitSeconds(3) // + .withQueryOptions(queryOptions) // + .withNettyOptions(FastShutdownNettyOptions.INSTANCE) // + .build(); } else { cluster = parent.cluster; cassandraPort = parent.cassandraPort; diff --git a/spring-cql/src/test/java/org/springframework/cassandra/test/integration/config/xml/PropertyPlaceholderNamespaceCreatingXmlConfigIntegrationTests.java b/spring-cql/src/test/java/org/springframework/cassandra/test/integration/config/xml/PropertyPlaceholderNamespaceCreatingXmlConfigIntegrationTests.java index 1570df810..544c4ac05 100644 --- a/spring-cql/src/test/java/org/springframework/cassandra/test/integration/config/xml/PropertyPlaceholderNamespaceCreatingXmlConfigIntegrationTests.java +++ b/spring-cql/src/test/java/org/springframework/cassandra/test/integration/config/xml/PropertyPlaceholderNamespaceCreatingXmlConfigIntegrationTests.java @@ -69,7 +69,6 @@ public class PropertyPlaceholderNamespaceCreatingXmlConfigIntegrationTests assertThat(poolingOptions, is(notNullValue(PoolingOptions.class))); assertThat(poolingOptions.getHeartbeatIntervalSeconds(), is(equalTo(60))); assertThat(poolingOptions.getIdleTimeoutSeconds(), is(equalTo(180))); - assertThat(poolingOptions.getPoolTimeoutMillis(), is(equalTo(30000))); assertThat(poolingOptions.getCoreConnectionsPerHost(HostDistance.LOCAL), is(equalTo(4))); assertThat(poolingOptions.getMaxConnectionsPerHost(HostDistance.LOCAL), is(equalTo(8))); assertThat(poolingOptions.getMaxRequestsPerConnection(HostDistance.LOCAL), is(equalTo(20))); diff --git a/spring-cql/src/test/java/org/springframework/cassandra/test/integration/config/xml/XmlConfigIntegrationTests.java b/spring-cql/src/test/java/org/springframework/cassandra/test/integration/config/xml/XmlConfigIntegrationTests.java index 766c10a1b..792b2f2cd 100644 --- a/spring-cql/src/test/java/org/springframework/cassandra/test/integration/config/xml/XmlConfigIntegrationTests.java +++ b/spring-cql/src/test/java/org/springframework/cassandra/test/integration/config/xml/XmlConfigIntegrationTests.java @@ -97,7 +97,7 @@ public class XmlConfigIntegrationTests extends AbstractEmbeddedCassandraIntegrat public void clusterConfigurationIsCorrect() { assertThat(cluster.getConfiguration().getPolicies().getAddressTranslator(), is(equalTo(addressTranslator))); assertThat(cluster.getClusterName(), is(equalTo("skynet"))); - assertThat(cluster.getConfiguration().getProtocolOptions().getMaxSchemaAgreementWaitSeconds(), is(equalTo(30))); + assertThat(cluster.getConfiguration().getProtocolOptions().getMaxSchemaAgreementWaitSeconds(), is(equalTo(2))); assertThat(cluster.getConfiguration().getPolicies().getSpeculativeExecutionPolicy(), is(equalTo(speculativeExecutionPolicy))); @@ -123,7 +123,6 @@ public class XmlConfigIntegrationTests extends AbstractEmbeddedCassandraIntegrat assertThat(poolingOptions.getHeartbeatIntervalSeconds(), is(equalTo(60))); assertThat(poolingOptions.getIdleTimeoutSeconds(), is(equalTo(300))); assertThat(poolingOptions.getInitializationExecutor(), is(equalTo(executor))); - assertThat(poolingOptions.getPoolTimeoutMillis(), is(equalTo(15000))); assertThat(poolingOptions.getCoreConnectionsPerHost(HostDistance.LOCAL), is(equalTo(2))); assertThat(poolingOptions.getMaxConnectionsPerHost(HostDistance.LOCAL), is(equalTo(8))); assertThat(poolingOptions.getMaxRequestsPerConnection(HostDistance.LOCAL), is(equalTo(100))); diff --git a/spring-cql/src/test/resources/org/springframework/cassandra/test/integration/config/xml/XmlConfigIntegrationTests-context.xml b/spring-cql/src/test/resources/org/springframework/cassandra/test/integration/config/xml/XmlConfigIntegrationTests-context.xml index 47db2c364..5fc254343 100644 --- a/spring-cql/src/test/resources/org/springframework/cassandra/test/integration/config/xml/XmlConfigIntegrationTests-context.xml +++ b/spring-cql/src/test/resources/org/springframework/cassandra/test/integration/config/xml/XmlConfigIntegrationTests-context.xml @@ -30,7 +30,7 @@ heartbeat-interval-seconds="60" initialization-executor-ref="testExecutor" idle-timeout-seconds="300" - max-schema-agreement-wait-seconds="30" + max-schema-agreement-wait-seconds="2" pool-timeout-milliseconds="15000" speculative-execution-policy-ref="testSpeculativeExecutionPolicy" timestamp-generator-ref="testTimestampGenerator"> diff --git a/spring-data-cassandra/src/main/resources/org/springframework/data/cassandra/config/spring-cassandra-1.0.xsd b/spring-data-cassandra/src/main/resources/org/springframework/data/cassandra/config/spring-cassandra-1.0.xsd index 4c0df3639..3edcf567b 100644 --- a/spring-data-cassandra/src/main/resources/org/springframework/data/cassandra/config/spring-cassandra-1.0.xsd +++ b/spring-data-cassandra/src/main/resources/org/springframework/data/cassandra/config/spring-cassandra-1.0.xsd @@ -272,6 +272,16 @@ LoadBalancingPolicy implementation. + + + + +