From 24fcea3f9253ecaf67941f1a372e74e4a5f366f4 Mon Sep 17 00:00:00 2001 From: Mark Paluch Date: Thu, 11 Jul 2019 15:33:35 +0200 Subject: [PATCH] DATACASS-318 - Allow configuration of CodecRegistry in CassandraClusterFactoryBean. --- .../config/CassandraClusterFactoryBean.java | 12 ++++++ .../CassandraClusterFactoryBeanUnitTests.java | 40 ++++++++----------- 2 files changed, 29 insertions(+), 23 deletions(-) diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/config/CassandraClusterFactoryBean.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/config/CassandraClusterFactoryBean.java index 5f78ce29f..c24a2fe2b 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/config/CassandraClusterFactoryBean.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/config/CassandraClusterFactoryBean.java @@ -132,6 +132,7 @@ public class CassandraClusterFactoryBean private @Nullable AddressTranslator addressTranslator; private @Nullable AuthProvider authProvider; private @Nullable CompressionType compressionType; + private @Nullable CodecRegistry codecRegistry; private @Nullable Host.StateListener hostStateListener; private @Nullable LatencyTracker latencyTracker; private @Nullable LoadBalancingPolicy loadBalancingPolicy; @@ -179,6 +180,7 @@ public class CassandraClusterFactoryBean .withMaxSchemaAgreementWaitSeconds(this.maxSchemaAgreementWaitSeconds).withPort(this.port); Optional.ofNullable(this.addressTranslator).ifPresent(clusterBuilder::withAddressTranslator); + Optional.ofNullable(this.codecRegistry).ifPresent(clusterBuilder::withCodecRegistry); Optional.ofNullable(this.loadBalancingPolicy).ifPresent(clusterBuilder::withLoadBalancingPolicy); Optional.ofNullable(this.nettyOptions).ifPresent(clusterBuilder::withNettyOptions); Optional.ofNullable(this.poolingOptions).ifPresent(clusterBuilder::withPoolingOptions); @@ -410,6 +412,16 @@ public class CassandraClusterFactoryBean this.compressionType = compressionType; } + /** + * Set the {@link CodecRegistry}. Default uses {@link CodecRegistry#DEFAULT_INSTANCE}. + * + * @param codecRegistry the {@link CodecRegistry} used by the new cluster. + * @since 2.2 + */ + public void setCodecRegistry(@Nullable CodecRegistry codecRegistry) { + this.codecRegistry = codecRegistry; + } + /** * Set the {@link PoolingOptions} to configure the connection pooling behavior. * diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/config/CassandraClusterFactoryBeanUnitTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/config/CassandraClusterFactoryBeanUnitTests.java index 90605adc9..d28f05712 100755 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/config/CassandraClusterFactoryBeanUnitTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/config/CassandraClusterFactoryBeanUnitTests.java @@ -15,35 +15,18 @@ */ package org.springframework.data.cassandra.config; -import static org.assertj.core.api.Assertions.assertThat; -import static org.mockito.ArgumentMatchers.anyInt; -import static org.mockito.ArgumentMatchers.eq; -import static org.mockito.ArgumentMatchers.isA; +import static org.assertj.core.api.Assertions.*; +import static org.mockito.ArgumentMatchers.*; +import static org.mockito.Mockito.*; import static org.mockito.Mockito.anyString; -import static org.mockito.Mockito.mock; -import static org.mockito.Mockito.spy; -import static org.mockito.Mockito.times; -import static org.mockito.Mockito.verify; -import static org.mockito.Mockito.when; import org.junit.Test; import org.springframework.data.cassandra.support.IntegrationTestNettyOptions; import org.springframework.test.util.ReflectionTestUtils; -import com.datastax.driver.core.AuthProvider; -import com.datastax.driver.core.Cluster; -import com.datastax.driver.core.Configuration; -import com.datastax.driver.core.PlainTextAuthProvider; -import com.datastax.driver.core.PoolingOptions; -import com.datastax.driver.core.ProtocolOptions; +import com.datastax.driver.core.*; import com.datastax.driver.core.ProtocolOptions.Compression; -import com.datastax.driver.core.ProtocolVersion; -import com.datastax.driver.core.QueryOptions; -import com.datastax.driver.core.RemoteEndpointAwareJdkSSLOptions; -import com.datastax.driver.core.SSLOptions; -import com.datastax.driver.core.SocketOptions; -import com.datastax.driver.core.TimestampGenerator; import com.datastax.driver.core.policies.AddressTranslator; import com.datastax.driver.core.policies.ExponentialReconnectionPolicy; import com.datastax.driver.core.policies.LoadBalancingPolicy; @@ -108,6 +91,18 @@ public class CassandraClusterFactoryBeanUnitTests { assertThat(getConfiguration(bean).getProtocolOptions().getCompression()).isEqualTo(Compression.SNAPPY); } + @Test // DATACASS-318 + public void shouldSetCodecRegistry() throws Exception { + + CodecRegistry codecRegistry = new CodecRegistry(); + + CassandraClusterFactoryBean bean = new CassandraClusterFactoryBean(); + bean.setCodecRegistry(codecRegistry); + bean.afterPropertiesSet(); + + assertThat(getConfiguration(bean).getCodecRegistry()).isSameAs(codecRegistry); + } + @Test // DATACASS-226 public void shouldSetPoolingOptions() throws Exception { @@ -150,7 +145,6 @@ public class CassandraClusterFactoryBeanUnitTests { CassandraClusterFactoryBean bean = new CassandraClusterFactoryBean(); bean.afterPropertiesSet(); - assertThat(getConfiguration(bean).getQueryOptions()).isNotNull(); assertThat(getConfiguration(bean).getQueryOptions()).isEqualTo(new QueryOptions()); } @@ -327,7 +321,7 @@ public class CassandraClusterFactoryBeanUnitTests { bean.setClusterName("XYZ"); bean.afterPropertiesSet(); - verify(bean,times(1)).newClusterBuilder(); + verify(bean, times(1)).newClusterBuilder(); verify(mockClusterBuilder, times(1)).withClusterName(eq("XYZ")); }