diff --git a/spring-data-geode/src/main/java/org/springframework/data/gemfire/AbstractResolvableCacheFactoryBean.java b/spring-data-geode/src/main/java/org/springframework/data/gemfire/AbstractResolvableCacheFactoryBean.java new file mode 100644 index 00000000..1013a82f --- /dev/null +++ b/spring-data-geode/src/main/java/org/springframework/data/gemfire/AbstractResolvableCacheFactoryBean.java @@ -0,0 +1,241 @@ +/* + * Copyright 2020 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 + * + * https://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.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.RuntimeExceptionFactory.newRuntimeException; + +import java.util.Optional; +import java.util.Properties; + +import org.apache.geode.cache.CacheClosedException; +import org.apache.geode.cache.GemFireCache; +import org.apache.geode.distributed.DistributedSystem; + +import org.springframework.lang.NonNull; + +/** + * Abstract base class encapsulating logic to resolve or create a {@link GemFireCache cache} instance. + * + * @author John Blum + * @see java.util.Properties + * @see org.apache.geode.cache.GemFireCache + * @see org.apache.geode.distributed.DistributedMember + * @see org.apache.geode.distributed.DistributedSystem + * @see org.springframework.data.gemfire.AbstractPdxConfigurableCacheFactoryBean + * @since 2.5.0 + */ +public abstract class AbstractResolvableCacheFactoryBean extends AbstractPdxConfigurableCacheFactoryBean { + + private volatile String cacheResolutionMessagePrefix; + + /** + * @inheritDoc + */ + @Override + protected GemFireCache doGetObject() { + return init(); + } + + /** + * Initializes a {@link GemFireCache}. + * + * @return a reference to the initialized {@link GemFireCache}. + * @see org.apache.geode.cache.GemFireCache + * @see #setCache(GemFireCache) + * @see #resolveCache() + * @see #getCache() + */ + protected GemFireCache init() { + + ClassLoader currentThreadContextClassLoader = Thread.currentThread().getContextClassLoader(); + + try { + // Use Spring Bean ClassLoader to load Spring configured Apache Geode classes + Thread.currentThread().setContextClassLoader(getBeanClassLoader()); + + setCache(resolveCache()); + + logCacheInitialization(); + + return getCache(); + } + catch (Exception cause) { + throw newRuntimeException(cause, "Error occurred while initializing the cache"); + } + finally { + Thread.currentThread().setContextClassLoader(currentThreadContextClassLoader); + } + } + + @SuppressWarnings("deprecation") + private void logCacheInitialization() { + + getOptionalCache().ifPresent(cache -> { + + Optional.ofNullable(cache.getDistributedSystem()) + .map(DistributedSystem::getDistributedMember) + .ifPresent(member -> { + + String message = "Connected to Distributed System [%1$s] as Member [%2$s] in Group(s) [%3$s]" + + " with Role(s) [%4$s] on Host [%5$s] having PID [%6$d]"; + + logInfo(() -> String.format(message, + cache.getDistributedSystem().getName(), member.getId(), member.getGroups(), + member.getRoles(), member.getHost(), member.getProcessId())); + }); + + logInfo(() -> String.format("%1$s %2$s version [%3$s] Cache [%4$s]", this.cacheResolutionMessagePrefix, + apacheGeodeProductName(), apacheGeodeVersion(), cache.getName())); + + }); + } + + /** + * Resolves a {@link GemFireCache} by attempting to lookup an existing {@link GemFireCache} instance in the JVM, + * first. If an existing {@link GemFireCache} could not be found, then this method proceeds in attempting to + * create a new {@link GemFireCache} instance. + * + * @param parameterized {@link Class} type extending {@link GemFireCache}. + * @return the resolved {@link GemFireCache}. + * @see org.apache.geode.cache.client.ClientCache + * @see org.apache.geode.cache.GemFireCache + * @see org.apache.geode.cache.Cache + * @see #fetchCache() + * @see #resolveProperties() + * @see #createFactory(java.util.Properties) + * @see #initializeFactory(Object) + * @see #configureFactory(Object) + * @see #postProcess(Object) + * @see #createCache(Object) + * @see #postProcess(GemFireCache) + */ + protected T resolveCache() { + + try { + + this.cacheResolutionMessagePrefix = "Found existing"; + + return fetchCache(); + } + catch (CacheClosedException cause) { + + this.cacheResolutionMessagePrefix = "Created new"; + + Properties gemfireProperties = resolveProperties(); + + Object factory = createFactory(gemfireProperties); + + factory = initializeFactory(factory); + factory = configureFactory(factory); + factory = postProcess(factory); + + T cache = createCache(factory); + + cache = postProcess(cache); + + return cache; + } + } + + /** + * Constructs a new cache factory initialized with the given Apache Geode {@link Properties} + * used to construct, configure and initialize a new {@link GemFireCache}. + * + * @param gemfireProperties {@link Properties} used by the cache factory to configure the {@link GemFireCache}; + * must not be {@literal null} + * @return a new cache factory initialized with the given Apache Geode {@link Properties}. + * @see org.apache.geode.cache.client.ClientCacheFactory + * @see org.apache.geode.cache.CacheFactory + * @see java.util.Properties + * @see #resolveProperties() + */ + protected abstract @NonNull Object createFactory(@NonNull Properties gemfireProperties); + + /** + * Configures the cache factory used to create the {@link GemFireCache}. + * + * @param factory cache factory to configure; must not be {@literal null}. + * @return the given cache factory. + * @see #createFactory(Properties) + */ + protected @NonNull Object configureFactory(@NonNull Object factory) { + return factory; + } + + /** + * @inheritDoc + */ + @Override + protected Object initializeFactory(Object factory) { + return super.initializeFactory(factory); + } + + /** + * Post process the cache factory used to create the {@link GemFireCache}. + * + * @param factory cache factory to post process; must not be {@literal null}. + * @return the post processed cache factory. + * @see org.apache.geode.cache.client.ClientCacheFactory + * @see org.apache.geode.cache.CacheFactory + * @see #createFactory(Properties) + */ + protected @NonNull Object postProcess(@NonNull Object factory) { + return factory; + } + + /** + * Creates a new {@link GemFireCache} instance using the provided {@link Object factory}. + * + * @param {@link Class Subtype} of {@link GemFireCache}. + * @param factory factory used to create the {@link GemFireCache}. + * @return a new instance of {@link GemFireCache} created by the provided {@link Object factory}. + * @see org.apache.geode.cache.client.ClientCacheFactory#create() + * @see org.apache.geode.cache.CacheFactory#create() + * @see org.apache.geode.cache.GemFireCache + */ + protected abstract @NonNull T createCache(@NonNull Object factory); + + /** + * Post process the {@link GemFireCache} by loading any {@literal cache.xml} file, applying custom settings + * specified in SDG XML configuration metadata, and registering appropriate Transaction Listeners, Writer + * and JVM Heap configuration. + * + * @param parameterized {@link Class} type extending {@link GemFireCache}. + * @param cache {@link GemFireCache} to post process. + * @return the given {@link GemFireCache}. + * @see #loadCacheXml(GemFireCache) + * @see org.apache.geode.cache.Cache#loadCacheXml(java.io.InputStream) + * @see #configureHeapPercentages(org.apache.geode.cache.GemFireCache) + * @see #configureOffHeapPercentages(GemFireCache) + * @see #registerTransactionListeners(org.apache.geode.cache.GemFireCache) + * @see #registerTransactionWriter(org.apache.geode.cache.GemFireCache) + */ + protected @NonNull T postProcess(@NonNull T cache) { + + loadCacheXml(cache); + + Optional.ofNullable(getCopyOnRead()).ifPresent(cache::setCopyOnRead); + + configureHeapPercentages(cache); + configureOffHeapPercentages(cache); + registerTransactionListeners(cache); + registerTransactionWriter(cache); + + return 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 03416328..0a05b433 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 @@ -15,10 +15,7 @@ */ 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.CollectionUtils.nullSafeList; -import static org.springframework.data.gemfire.util.RuntimeExceptionFactory.newRuntimeException; import java.util.ArrayList; import java.util.Arrays; @@ -29,11 +26,9 @@ import java.util.Properties; import java.util.stream.StreamSupport; 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.util.GatewayConflictResolver; -import org.apache.geode.distributed.DistributedSystem; import org.apache.geode.internal.datasource.ConfigProperty; import org.apache.geode.internal.jndi.JNDIInvoker; import org.apache.geode.pdx.PdxSerializer; @@ -63,14 +58,11 @@ import org.springframework.util.Assert; * @see org.apache.geode.cache.GemFireCache * @see org.apache.geode.pdx.PdxSerializer * @see org.apache.geode.security.SecurityManager - * @see org.apache.geode.distributed.DistributedMember - * @see org.apache.geode.distributed.DistributedSystem - * @see org.springframework.beans.factory.BeanFactory * @see org.springframework.beans.factory.FactoryBean * @see org.springframework.data.gemfire.config.annotation.PeerCacheConfigurer */ @SuppressWarnings("unused") -public class CacheFactoryBean extends AbstractPdxConfigurableCacheFactoryBean { +public class CacheFactoryBean extends AbstractResolvableCacheFactoryBean { private Boolean enableAutoReconnect; private Boolean useClusterConfiguration; @@ -90,8 +82,6 @@ public class CacheFactoryBean extends AbstractPdxConfigurableCacheFactoryBean { nullSafeList(peerCacheConfigurers).forEach(peerCacheConfigurer -> peerCacheConfigurer.configure(beanName, bean)); - private String cacheResolutionMessagePrefix; - private org.apache.geode.security.SecurityManager securityManager; /** @@ -151,8 +141,9 @@ public class CacheFactoryBean extends AbstractPdxConfigurableCacheFactoryBean { * @inheritDoc */ @Override - protected GemFireCache doGetObject() { - return init(); + @SuppressWarnings("unchecked") + protected T doFetchCache() { + return (T) CacheFactory.getAnyInstance(); } /** @@ -163,87 +154,6 @@ public class CacheFactoryBean extends AbstractPdxConfigurableCacheFactoryBean { return Cache.class; } - /** - * Initializes the {@link Cache}. - * - * @return a reference to the initialized {@link Cache}. - * @see org.apache.geode.cache.Cache - * @see #resolveCache() - * @see #postProcess(GemFireCache) - * @see #setCache(GemFireCache) - */ - @SuppressWarnings("deprecation") - GemFireCache init() { - - ClassLoader currentThreadContextClassLoader = Thread.currentThread().getContextClassLoader(); - - try { - // Use Spring Bean ClassLoader to load Spring configured Apache Geode classes - Thread.currentThread().setContextClassLoader(getBeanClassLoader()); - - setCache(postProcess(resolveCache())); - - Optional.ofNullable(getCache()).ifPresent(cache -> { - - Optional.ofNullable(cache.getDistributedSystem()) - .map(DistributedSystem::getDistributedMember) - .ifPresent(member -> - logInfo(() -> String.format("Connected to Distributed System [%1$s] as Member [%2$s]" - .concat(" in Group(s) [%3$s] with Role(s) [%4$s] on Host [%5$s] having PID [%6$d]"), - cache.getDistributedSystem().getName(), member.getId(), member.getGroups(), - member.getRoles(), member.getHost(), member.getProcessId()))); - - logInfo(() -> String.format("%1$s %2$s version [%3$s] Cache [%4$s]", this.cacheResolutionMessagePrefix, - apacheGeodeProductName(), apacheGeodeVersion(), cache.getName())); - - }); - - return getCache(); - } - catch (Exception cause) { - throw newRuntimeException(cause, "An error occurred while initializing the cache"); - } - finally { - Thread.currentThread().setContextClassLoader(currentThreadContextClassLoader); - } - } - - /** - * Resolves the {@link Cache} by first attempting to lookup an existing {@link Cache} instance in the JVM. - * If an existing {@link Cache} could not be found, then this method proceeds in attempting to create - * a new {@link Cache} instance. - * - * @param parameterized {@link Class} type extension of {@link GemFireCache}. - * @return the resolved {@link Cache} instance. - * @see org.apache.geode.cache.Cache - * @see #fetchCache() - * @see #resolveProperties() - * @see #createFactory(java.util.Properties) - * @see #configureFactory(Object) - * @see #createCache(Object) - */ - protected T resolveCache() { - - try { - - this.cacheResolutionMessagePrefix = "Found existing"; - - return fetchCache(); - } - catch (CacheClosedException cause) { - - this.cacheResolutionMessagePrefix = "Created new"; - - return createCache(postProcess(configureFactory(initializeFactory(createFactory(resolveProperties()))))); - } - } - - @Override - @SuppressWarnings("unchecked") - protected T doFetchCache() { - return (T) CacheFactory.getAnyInstance(); - } - /** * Constructs a new instance of {@link CacheFactory} initialized with the given Apache Geode {@link Properties} * used to construct, configure and initialize a new peer {@link Cache} instance. @@ -253,6 +163,7 @@ public class CacheFactoryBean extends AbstractPdxConfigurableCacheFactoryBean { * @see org.apache.geode.cache.CacheFactory * @see java.util.Properties */ + @Override protected @NonNull Object createFactory(@NonNull Properties gemfireProperties) { return new CacheFactory(gemfireProperties); } @@ -266,6 +177,7 @@ public class CacheFactoryBean extends AbstractPdxConfigurableCacheFactoryBean { * @see #configureSecurity(CacheFactory) * @see org.apache.geode.cache.CacheFactory */ + @Override protected @NonNull Object configureFactory(@NonNull Object factory) { return configureSecurity(configurePdx((CacheFactory) factory)); } @@ -299,17 +211,6 @@ public class CacheFactoryBean extends AbstractPdxConfigurableCacheFactoryBean { : cacheFactory; } - /** - * Post process the {@link CacheFactory} used to create the {@link Cache}. - * - * @param factory {@link CacheFactory} used to create the {@link Cache}. - * @return the post processed {@link CacheFactory}. - * @see org.apache.geode.cache.CacheFactory - */ - protected @NonNull Object postProcess(@NonNull Object factory) { - return factory; - } - /** * Creates a new {@link Cache} instance using the provided {@link Object factory}. * @@ -320,6 +221,7 @@ public class CacheFactoryBean extends AbstractPdxConfigurableCacheFactoryBean { * @see org.apache.geode.cache.GemFireCache */ @SuppressWarnings("unchecked") + @Override protected @NonNull T createCache(@NonNull Object factory) { return (T) ((CacheFactory) factory).create(); } @@ -340,12 +242,10 @@ public class CacheFactoryBean extends AbstractPdxConfigurableCacheFactoryBean { * @see #registerTransactionListeners(org.apache.geode.cache.GemFireCache) * @see #registerTransactionWriter(org.apache.geode.cache.GemFireCache) */ - // TODO: Refactor this garbage! + @Override protected @NonNull T postProcess(@NonNull T cache) { - loadCacheXml(cache); - - Optional.ofNullable(getCopyOnRead()).ifPresent(cache::setCopyOnRead); + super.postProcess(cache); if (cache instanceof Cache) { @@ -358,11 +258,7 @@ public class CacheFactoryBean extends AbstractPdxConfigurableCacheFactoryBean { Optional.ofNullable(getSearchTimeout()).ifPresent(peerCache::setSearchTimeout); } - configureHeapPercentages(cache); - configureOffHeapPercentages(cache); registerJndiDataSources(cache); - registerTransactionListeners(cache); - registerTransactionWriter(cache); return cache; } diff --git a/spring-data-geode/src/main/java/org/springframework/data/gemfire/client/ClientCacheFactoryBean.java b/spring-data-geode/src/main/java/org/springframework/data/gemfire/client/ClientCacheFactoryBean.java index 0162dd83..8ad14ca1 100644 --- a/spring-data-geode/src/main/java/org/springframework/data/gemfire/client/ClientCacheFactoryBean.java +++ b/spring-data-geode/src/main/java/org/springframework/data/gemfire/client/ClientCacheFactoryBean.java @@ -74,7 +74,6 @@ import org.springframework.util.StringUtils; * @see org.apache.geode.cache.client.SocketFactory * @see org.apache.geode.distributed.DistributedSystem * @see org.apache.geode.pdx.PdxSerializer - * @see org.springframework.beans.factory.BeanFactory * @see org.springframework.beans.factory.FactoryBean * @see org.springframework.context.ApplicationContext * @see org.springframework.context.ApplicationListener