From 10a01b13bcea03f99c6deb9642fde292f1717c63 Mon Sep 17 00:00:00 2001 From: John Blum Date: Wed, 10 Mar 2021 08:56:45 -0800 Subject: [PATCH] Move TransactonListener/Writer handling logic to the AbstractBasicCacheFactoryBean class. Resolves gh-493. --- .../AbstractBasicCacheFactoryBean.java | 73 ++++++++++++++ .../data/gemfire/CacheFactoryBean.java | 99 ++++--------------- 2 files changed, 92 insertions(+), 80 deletions(-) diff --git a/spring-data-geode/src/main/java/org/springframework/data/gemfire/AbstractBasicCacheFactoryBean.java b/spring-data-geode/src/main/java/org/springframework/data/gemfire/AbstractBasicCacheFactoryBean.java index 94d51d23..03169c90 100644 --- a/spring-data-geode/src/main/java/org/springframework/data/gemfire/AbstractBasicCacheFactoryBean.java +++ b/spring-data-geode/src/main/java/org/springframework/data/gemfire/AbstractBasicCacheFactoryBean.java @@ -21,6 +21,8 @@ import static org.springframework.data.gemfire.util.RuntimeExceptionFactory.newR import java.io.File; import java.io.IOException; import java.io.InputStream; +import java.util.List; +import java.util.Objects; import java.util.Optional; import java.util.Properties; @@ -30,6 +32,8 @@ import org.apache.geode.cache.Cache; import org.apache.geode.cache.CacheFactory; import org.apache.geode.cache.GemFireCache; import org.apache.geode.cache.Region; +import org.apache.geode.cache.TransactionListener; +import org.apache.geode.cache.TransactionWriter; import org.apache.geode.cache.client.ClientCache; import org.apache.geode.cache.client.ClientCacheFactory; @@ -47,6 +51,7 @@ import org.springframework.data.gemfire.config.annotation.ClientCacheConfigurer; import org.springframework.data.gemfire.config.annotation.PeerCacheConfigurer; import org.springframework.data.gemfire.support.AbstractFactoryBeanSupport; import org.springframework.data.gemfire.support.GemfireBeanFactoryLocator; +import org.springframework.data.gemfire.util.CollectionUtils; import org.springframework.lang.NonNull; import org.springframework.lang.Nullable; import org.springframework.util.Assert; @@ -112,10 +117,14 @@ public abstract class AbstractBasicCacheFactoryBean extends AbstractFactoryBeanS private GemFireCache cache; + private List transactionListeners; + private Properties properties; private Resource cacheXml; + private TransactionWriter transactionWriter; + /** * Gets a reference to the configured {@link GemfireBeanFactoryLocator} used to resolve Spring bean references * in Apache Geode native configuration metadata (e.g. {@literal cache.xml}). @@ -444,6 +453,53 @@ public abstract class AbstractBasicCacheFactoryBean extends AbstractFactoryBeanS return this.properties; } + /** + * Configures the cache (transaction manager) with a {@link List} of {@link TransactionListener TransactionListeners} + * implemented by applications to listen for and receive transaction events after a transaction is processed + * (i.e. committed or rolled back). + * + * @param transactionListeners {@link List} of application-defined {@link TransactionListener TransactionListeners} + * registered with the cache to listen for and receive transaction events. + * @see org.apache.geode.cache.TransactionListener + */ + public void setTransactionListeners(@NonNull List transactionListeners) { + this.transactionListeners = transactionListeners; + } + + /** + * Returns the {@link List} of configured, application-defined {@link TransactionListener TransactionListeners} + * registered with the cache (transaction manager) to enable applications to receive transaction events after a + * transaction is processed (i.e. committed or rolled back). + * + * @return a {@link List} of application-defined {@link TransactionListener TransactionListeners} registered with + * the cache (transaction manager) to listen for and receive transaction events. + * @see org.apache.geode.cache.TransactionListener + */ + public @NonNull List getTransactionListeners() { + return CollectionUtils.nullSafeList(this.transactionListeners); + } + + /** + * Configures a {@link TransactionWriter} implemented by the application to receive transaction events and perform + * a action, like a veto. + * + * @param transactionWriter {@link TransactionWriter} receiving transaction events. + * @see org.apache.geode.cache.TransactionWriter + */ + public void setTransactionWriter(@Nullable TransactionWriter transactionWriter) { + this.transactionWriter = transactionWriter; + } + + /** + * Return the configured {@link TransactionWriter} used to process and handle transaction events. + * + * @return the configured {@link TransactionWriter}. + * @see org.apache.geode.cache.TransactionWriter + */ + public @Nullable TransactionWriter getTransactionWriter() { + return this.transactionWriter; + } + /** * Sets a boolean value used to determine whether to enable the {@link GemfireBeanFactoryLocator}. * @@ -665,6 +721,23 @@ public abstract class AbstractBasicCacheFactoryBean extends AbstractFactoryBeanS return cache; } + protected GemFireCache registerTransactionListeners(GemFireCache cache) { + + CollectionUtils.nullSafeCollection(getTransactionListeners()).stream() + .filter(Objects::nonNull) + .forEach(transactionListener -> cache.getCacheTransactionManager().addListener(transactionListener)); + + return cache; + } + + protected GemFireCache registerTransactionWriter(GemFireCache cache) { + + Optional.ofNullable(getTransactionWriter()) + .ifPresent(transactionWriter -> cache.getCacheTransactionManager().setWriter(transactionWriter)); + + return cache; + } + /** * Resolves the Apache Geode {@link Properties} used to configure the {@link Cache}. * diff --git a/spring-data-geode/src/main/java/org/springframework/data/gemfire/CacheFactoryBean.java b/spring-data-geode/src/main/java/org/springframework/data/gemfire/CacheFactoryBean.java index 407f0312..03416328 100644 --- a/spring-data-geode/src/main/java/org/springframework/data/gemfire/CacheFactoryBean.java +++ b/spring-data-geode/src/main/java/org/springframework/data/gemfire/CacheFactoryBean.java @@ -17,7 +17,6 @@ package org.springframework.data.gemfire; import static org.springframework.data.gemfire.GemfireUtils.apacheGeodeProductName; import static org.springframework.data.gemfire.GemfireUtils.apacheGeodeVersion; -import static org.springframework.data.gemfire.util.ArrayUtils.nullSafeArray; import static org.springframework.data.gemfire.util.CollectionUtils.nullSafeList; import static org.springframework.data.gemfire.util.RuntimeExceptionFactory.newRuntimeException; @@ -33,8 +32,6 @@ import org.apache.geode.cache.Cache; import org.apache.geode.cache.CacheClosedException; import org.apache.geode.cache.CacheFactory; import org.apache.geode.cache.GemFireCache; -import org.apache.geode.cache.TransactionListener; -import org.apache.geode.cache.TransactionWriter; import org.apache.geode.cache.util.GatewayConflictResolver; import org.apache.geode.distributed.DistributedSystem; import org.apache.geode.internal.datasource.ConfigProperty; @@ -48,6 +45,7 @@ import org.springframework.data.gemfire.util.ArrayUtils; import org.springframework.data.gemfire.util.CollectionUtils; import org.springframework.data.gemfire.util.SpringUtils; import org.springframework.lang.NonNull; +import org.springframework.lang.Nullable; import org.springframework.util.Assert; /** @@ -88,8 +86,6 @@ public class CacheFactoryBean extends AbstractPdxConfigurableCacheFactoryBean { private List jndiDataSources; - private List transactionListeners; - private final PeerCacheConfigurer compositePeerCacheConfigurer = (beanName, bean) -> nullSafeList(peerCacheConfigurers).forEach(peerCacheConfigurer -> peerCacheConfigurer.configure(beanName, bean)); @@ -98,8 +94,6 @@ public class CacheFactoryBean extends AbstractPdxConfigurableCacheFactoryBean { private org.apache.geode.security.SecurityManager securityManager; - private TransactionWriter transactionWriter; - /** * Applies the composite {@link PeerCacheConfigurer PeerCacheConfigurers} to this {@link CacheFactoryBean} * before the {@link Cache peer cache} is created. @@ -394,22 +388,6 @@ public class CacheFactoryBean extends AbstractPdxConfigurableCacheFactoryBean { return cache; } - private GemFireCache registerTransactionListeners(GemFireCache cache) { - - CollectionUtils.nullSafeCollection(getTransactionListeners()) - .forEach(transactionListener -> cache.getCacheTransactionManager().addListener(transactionListener)); - - return cache; - } - - private GemFireCache registerTransactionWriter(GemFireCache cache) { - - Optional.ofNullable(getTransactionWriter()) - .ifPresent(transactionWriter -> cache.getCacheTransactionManager().setWriter(transactionWriter)); - - return cache; - } - /** * Returns a reference to the Composite {@link PeerCacheConfigurer} used to apply additional configuration * to this {@link CacheFactoryBean} on Spring container initialization. @@ -417,7 +395,7 @@ public class CacheFactoryBean extends AbstractPdxConfigurableCacheFactoryBean { * @return the Composite {@link PeerCacheConfigurer}. * @see org.springframework.data.gemfire.config.annotation.PeerCacheConfigurer */ - public PeerCacheConfigurer getCompositePeerCacheConfigurer() { + public @NonNull PeerCacheConfigurer getCompositePeerCacheConfigurer() { return this.compositePeerCacheConfigurer; } @@ -427,7 +405,7 @@ public class CacheFactoryBean extends AbstractPdxConfigurableCacheFactoryBean { * @param enableAutoReconnect a boolean value to enable/disable auto-reconnect functionality. * @since GemFire 8.0 */ - public void setEnableAutoReconnect(Boolean enableAutoReconnect) { + public void setEnableAutoReconnect(@Nullable Boolean enableAutoReconnect) { this.enableAutoReconnect = enableAutoReconnect; } @@ -437,7 +415,7 @@ public class CacheFactoryBean extends AbstractPdxConfigurableCacheFactoryBean { * @return a boolean value indicating whether auto-reconnect was specified (non-null) and whether it was enabled * or not. */ - public Boolean getEnableAutoReconnect() { + public @Nullable Boolean getEnableAutoReconnect() { return this.enableAutoReconnect; } @@ -447,14 +425,14 @@ public class CacheFactoryBean extends AbstractPdxConfigurableCacheFactoryBean { * compatibility with Gemfire 6 compatibility. This must be an instance of * {@link org.apache.geode.cache.util.GatewayConflictResolver} */ - public void setGatewayConflictResolver(GatewayConflictResolver gatewayConflictResolver) { + public void setGatewayConflictResolver(@Nullable GatewayConflictResolver gatewayConflictResolver) { this.gatewayConflictResolver = gatewayConflictResolver; } /** * @return the gatewayConflictResolver */ - public GatewayConflictResolver getGatewayConflictResolver() { + public @Nullable GatewayConflictResolver getGatewayConflictResolver() { return this.gatewayConflictResolver; } @@ -477,14 +455,14 @@ public class CacheFactoryBean extends AbstractPdxConfigurableCacheFactoryBean { * * @param lockLease an integer value indicating the object lock lease timeout. */ - public void setLockLease(Integer lockLease) { + public void setLockLease(@Nullable Integer lockLease) { this.lockLease = lockLease; } /** * @return the lockLease */ - public Integer getLockLease() { + public @Nullable Integer getLockLease() { return this.lockLease; } @@ -493,14 +471,14 @@ public class CacheFactoryBean extends AbstractPdxConfigurableCacheFactoryBean { * * @param lockTimeout an integer value specifying the object lock request timeout. */ - public void setLockTimeout(Integer lockTimeout) { + public void setLockTimeout(@Nullable Integer lockTimeout) { this.lockTimeout = lockTimeout; } /** * @return the lockTimeout */ - public Integer getLockTimeout() { + public @Nullable Integer getLockTimeout() { return this.lockTimeout; } @@ -512,14 +490,14 @@ public class CacheFactoryBean extends AbstractPdxConfigurableCacheFactoryBean { * @param messageSyncInterval an integer value specifying the number of seconds in which the primary server * sends messages to secondary servers. */ - public void setMessageSyncInterval(Integer messageSyncInterval) { + public void setMessageSyncInterval(@Nullable Integer messageSyncInterval) { this.messageSyncInterval = messageSyncInterval; } /** * @return the messageSyncInterval */ - public Integer getMessageSyncInterval() { + public @Nullable Integer getMessageSyncInterval() { return this.messageSyncInterval; } @@ -533,7 +511,7 @@ public class CacheFactoryBean extends AbstractPdxConfigurableCacheFactoryBean { * @see #setPeerCacheConfigurers(List) */ public void setPeerCacheConfigurers(PeerCacheConfigurer... peerCacheConfigurers) { - setPeerCacheConfigurers(Arrays.asList(nullSafeArray(peerCacheConfigurers, PeerCacheConfigurer.class))); + setPeerCacheConfigurers(Arrays.asList(ArrayUtils.nullSafeArray(peerCacheConfigurers, PeerCacheConfigurer.class))); } /** @@ -553,14 +531,14 @@ public class CacheFactoryBean extends AbstractPdxConfigurableCacheFactoryBean { * * @param searchTimeout an integer value indicating the netSearch timeout value. */ - public void setSearchTimeout(Integer searchTimeout) { + public void setSearchTimeout(@Nullable Integer searchTimeout) { this.searchTimeout = searchTimeout; } /** * @return the searchTimeout */ - public Integer getSearchTimeout() { + public @Nullable Integer getSearchTimeout() { return this.searchTimeout; } @@ -570,7 +548,7 @@ public class CacheFactoryBean extends AbstractPdxConfigurableCacheFactoryBean { * @param securityManager {@link org.apache.geode.security.SecurityManager} used to secure this cache. * @see org.apache.geode.security.SecurityManager */ - public void setSecurityManager(SecurityManager securityManager) { + public void setSecurityManager(@Nullable SecurityManager securityManager) { this.securityManager = securityManager; } @@ -580,49 +558,10 @@ public class CacheFactoryBean extends AbstractPdxConfigurableCacheFactoryBean { * @return the {@link org.apache.geode.security.SecurityManager} used to secure this cache. * @see org.apache.geode.security.SecurityManager */ - public SecurityManager getSecurityManager() { + public @Nullable SecurityManager getSecurityManager() { return this.securityManager; } - /** - * Sets the list of TransactionListeners used to configure the Cache to receive transaction events after - * the transaction is processed (committed, rolled back). - * - * @param transactionListeners the list of GemFire TransactionListeners listening for transaction events. - * @see org.apache.geode.cache.TransactionListener - */ - public void setTransactionListeners(List transactionListeners) { - this.transactionListeners = transactionListeners; - } - - /** - * @return the transactionListeners - */ - public List getTransactionListeners() { - return this.transactionListeners; - } - - /** - * Sets the {@link TransactionWriter} used to configure the cache for handling transaction events, such as to veto - * the transaction or update an external DB before the commit. - * - * @param transactionWriter configured {@link TransactionWriter} callback receiving transaction events. - * @see org.apache.geode.cache.TransactionWriter - */ - public void setTransactionWriter(TransactionWriter transactionWriter) { - this.transactionWriter = transactionWriter; - } - - /** - * Return the configured {@link TransactionWriter} used to process and handle transaction events. - * - * @return the configured {@link TransactionWriter}. - * @see org.apache.geode.cache.TransactionWriter - */ - public TransactionWriter getTransactionWriter() { - return this.transactionWriter; - } - /** * Sets the state of the {@literal use-shared-configuration} Pivotal GemFire/Apache Geode * distribution configuration setting. @@ -630,7 +569,7 @@ public class CacheFactoryBean extends AbstractPdxConfigurableCacheFactoryBean { * @param useSharedConfiguration boolean value to set the {@literal use-shared-configuration} * Pivotal GemFire/Apache Geode distribution configuration setting. */ - public void setUseClusterConfiguration(Boolean useSharedConfiguration) { + public void setUseClusterConfiguration(@Nullable Boolean useSharedConfiguration) { this.useClusterConfiguration = useSharedConfiguration; } @@ -641,7 +580,7 @@ public class CacheFactoryBean extends AbstractPdxConfigurableCacheFactoryBean { * @return the current boolean value for the {@literal use-shared-configuration} * Pivotal GemFire/Apache Geode distribution configuration setting. */ - public Boolean getUseClusterConfiguration() { + public @Nullable Boolean getUseClusterConfiguration() { return this.useClusterConfiguration; }