SGF-137, SGF-138

This commit is contained in:
David Turanski
2012-11-29 10:44:08 -05:00
parent 3df0e93175
commit 7653b88c9b
10 changed files with 448 additions and 35 deletions

View File

@@ -21,13 +21,14 @@ import java.lang.reflect.Field;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.springframework.beans.factory.DisposableBean;
import org.springframework.context.SmartLifecycle;
import org.springframework.core.io.Resource;
import org.springframework.data.gemfire.client.ClientRegionFactoryBean;
import org.springframework.data.gemfire.wan.SmartLifecycleGatewaySender;
import org.springframework.util.Assert;
import org.springframework.util.ObjectUtils;
import org.springframework.util.ReflectionUtils;
import com.gemstone.gemfire.cache.asyncqueue.AsyncEventQueue;
import com.gemstone.gemfire.cache.AttributesFactory;
import com.gemstone.gemfire.cache.Cache;
import com.gemstone.gemfire.cache.CacheClosedException;
@@ -40,6 +41,7 @@ import com.gemstone.gemfire.cache.Region;
import com.gemstone.gemfire.cache.RegionAttributes;
import com.gemstone.gemfire.cache.RegionFactory;
import com.gemstone.gemfire.cache.Scope;
import com.gemstone.gemfire.cache.asyncqueue.AsyncEventQueue;
import com.gemstone.gemfire.cache.wan.GatewaySender;
/**
@@ -55,7 +57,11 @@ import com.gemstone.gemfire.cache.wan.GatewaySender;
* @author Costin Leau
* @author David Turanski
*/
public class RegionFactoryBean<K, V> extends RegionLookupFactoryBean<K, V> implements DisposableBean {
public class RegionFactoryBean<K, V> extends RegionLookupFactoryBean<K, V> implements DisposableBean, SmartLifecycle {
private boolean autoStartup = true;
private boolean running;
protected final Log log = LogFactory.getLog(getClass());
@@ -150,13 +156,13 @@ public class RegionFactoryBean<K, V> extends RegionLookupFactoryBean<K, V> imple
if (cacheWriter != null) {
regionFactory.setCacheWriter(cacheWriter);
}
if (diskStoreName != null) {
regionFactory.setDiskStoreName(diskStoreName);
Assert.isTrue(!isNotPersistent(),"it is invalid to specify a disk store if 'persistent' is set to false.");
Assert.isTrue(!isNotPersistent(), "it is invalid to specify a disk store if 'persistent' is set to false.");
persistent = true;
}
resolveDataPolicy(regionFactory, persistent, dataPolicy);
if (scope != null) {
@@ -194,8 +200,7 @@ public class RegionFactoryBean<K, V> extends RegionLookupFactoryBean<K, V> imple
if (dataPolicy == null) {
if (isPersistent()) {
regionFactory.setDataPolicy(DataPolicy.PERSISTENT_REPLICATE);
}
else {
} else {
regionFactory.setDataPolicy(DataPolicy.DEFAULT);
}
return;
@@ -238,23 +243,21 @@ public class RegionFactoryBean<K, V> extends RegionLookupFactoryBean<K, V> imple
* @param region
*/
protected void postProcess(Region<K, V> region) {
// do nothing
}
@Override
public void destroy() throws Exception {
if (region != null) {
if (close) {
if (!region.getCache().isClosed()) {
if (!region.getRegionService().isClosed()) {
try {
region.close();
}
catch (CacheClosedException cce) {
} catch (CacheClosedException cce) {
// nothing to see folks, move on.
}
}
}
else if (destroy) {
} else if (destroy) {
region.destroyRegion();
}
}
@@ -414,4 +417,73 @@ public class RegionFactoryBean<K, V> extends RegionLookupFactoryBean<K, V> imple
protected boolean isNotPersistent() {
return persistent != null && !persistent;
}
/* (non-Javadoc)
* @see org.springframework.context.Lifecycle#start()
*/
@Override
public void start() {
if (!ObjectUtils.isEmpty(gatewaySenders)) {
synchronized (gatewaySenders) {
for (Object obj : gatewaySenders) {
SmartLifecycleGatewaySender gws = (SmartLifecycleGatewaySender) obj;
if (gws.isAutoStartup() && !gws.isRunning()) {
gws.start();
}
}
}
}
this.running = true;
}
/* (non-Javadoc)
* @see org.springframework.context.Lifecycle#stop()
*/
@Override
public void stop() {
if (!ObjectUtils.isEmpty(gatewaySenders)) {
synchronized (gatewaySenders) {
for (Object obj : gatewaySenders) {
SmartLifecycleGatewaySender gws = (SmartLifecycleGatewaySender) obj;
gws.stop();
}
}
}
this.running = false;
}
/* (non-Javadoc)
* @see org.springframework.context.Lifecycle#isRunning()
*/
@Override
public boolean isRunning() {
return this.running;
}
/* (non-Javadoc)
* @see org.springframework.context.Phased#getPhase()
*/
@Override
public int getPhase() {
return Integer.MAX_VALUE;
}
/* (non-Javadoc)
* @see org.springframework.context.SmartLifecycle#isAutoStartup()
*/
@Override
public boolean isAutoStartup() {
return this.autoStartup;
}
/* (non-Javadoc)
* @see org.springframework.context.SmartLifecycle#stop(java.lang.Runnable)
*/
@Override
public void stop(Runnable callback) {
stop();
callback.run();
}
}

View File

@@ -19,11 +19,10 @@ import org.springframework.beans.factory.support.BeanDefinitionBuilder;
import org.springframework.beans.factory.xml.AbstractSingleBeanDefinitionParser;
import org.springframework.beans.factory.xml.ParserContext;
import org.springframework.data.gemfire.wan.AsyncEventQueueFactoryBean;
import org.springframework.util.StringUtils;
import org.springframework.util.xml.DomUtils;
import org.w3c.dom.Element;
import com.gemstone.gemfire.internal.lang.StringUtils;
/**
* @author David Turanski
*
@@ -41,8 +40,10 @@ public class AsyncEventQueueParser extends AbstractSingleBeanDefinitionParser {
Element asyncEventListenerElement = DomUtils.getChildElementByTagName(element, "async-event-listener");
Object asyncEventListener = ParsingUtils.parseRefOrSingleNestedBeanDeclaration(parserContext,
asyncEventListenerElement, builder);
String cacheName = StringUtils.isEmpty(element.getAttribute("cache-ref")) ? "gemfireCache" : element
.getAttribute("cache-ref");
String cacheName = !StringUtils.hasText(element.getAttribute("cache-ref")) ? GemfireConstants.DEFAULT_GEMFIRE_CACHE_NAME
: element.getAttribute("cache-ref");
builder.addConstructorArgReference(cacheName);
builder.addConstructorArgValue(asyncEventListener);
ParsingUtils.setPropertyValue(element, builder, "batch-size");
@@ -50,5 +51,23 @@ public class AsyncEventQueueParser extends AbstractSingleBeanDefinitionParser {
ParsingUtils.setPropertyValue(element, builder, "disk-store-ref");
ParsingUtils.setPropertyValue(element, builder, "persistent");
ParsingUtils.setPropertyValue(element, builder, "parallel");
ParsingUtils.setPropertyValue(element, builder, NAME_ATTRIBUTE);
if (!StringUtils.hasText(element.getAttribute(NAME_ATTRIBUTE))) {
if (element.getParentNode().getNodeName().endsWith("region")) {
Element region = (Element) element.getParentNode();
String regionName = StringUtils.hasText(region.getAttribute("name")) ? region.getAttribute("name")
: region.getAttribute("id");
int i = 0;
String name = regionName + ".asyncEventQueue#" + i;
while (parserContext.getRegistry().isBeanNameInUse(name)) {
i++;
name = regionName + ".asyncEventQueue#" + i;
}
builder.addPropertyValue("name", name);
}
}
}
}

View File

@@ -36,10 +36,12 @@ class GatewaySenderParser extends AbstractSimpleBeanDefinitionParser {
@Override
protected void doParse(Element element, ParserContext parserContext, BeanDefinitionBuilder builder) {
builder.setLazyInit(false);
String cacheRef = element.getAttribute("cache-ref");
// add cache reference (fallback to default if nothing is specified)
builder.addConstructorArgReference((StringUtils.hasText(cacheRef) ? cacheRef : "gemfireCache"));
builder.addConstructorArgReference((StringUtils.hasText(cacheRef) ? cacheRef
: GemfireConstants.DEFAULT_GEMFIRE_CACHE_NAME));
ParsingUtils.setPropertyValue(element, builder, "alert-threshold");
ParsingUtils.setPropertyValue(element, builder, "batch-size");
ParsingUtils.setPropertyValue(element, builder, "batch-time-interval");
@@ -55,6 +57,7 @@ class GatewaySenderParser extends AbstractSimpleBeanDefinitionParser {
ParsingUtils.setPropertyValue(element, builder, "socket-read-timeout");
ParsingUtils.setPropertyValue(element, builder, "persistent");
ParsingUtils.setPropertyValue(element, builder, "parallel");
ParsingUtils.setPropertyValue(element, builder, NAME_ATTRIBUTE);
Element eventFilterElement = DomUtils.getChildElementByTagName(element, "event-filter");
if (eventFilterElement != null) {
@@ -63,5 +66,25 @@ class GatewaySenderParser extends AbstractSimpleBeanDefinitionParser {
}
ParsingUtils.parseTransportFilters(element, parserContext, builder);
/**
* set the name for an inner bean
*/
if (!StringUtils.hasText(element.getAttribute(NAME_ATTRIBUTE))) {
if (element.getParentNode().getNodeName().endsWith("region")) {
Element region = (Element) element.getParentNode();
String regionName = StringUtils.hasText(region.getAttribute("name")) ? region.getAttribute("name")
: region.getAttribute("id");
int i = 0;
String name = regionName + ".gatewaySender#" + i;
while (parserContext.getRegistry().isBeanNameInUse(name)) {
i++;
name = regionName + ".gatewaySender#" + i;
}
builder.addPropertyValue("name", name);
}
}
}
}

View File

@@ -34,15 +34,25 @@ public abstract class AbstractWANComponentFactoryBean<T> implements FactoryBean<
DisposableBean {
protected Log log = LogFactory.getLog(this.getClass());
protected String name;
private String name;
protected final Cache cache;
protected Object factory;
private String beanName;
protected AbstractWANComponentFactoryBean(Cache cache) {
this.cache = cache;
}
public void setName(String name) {
this.name = name;
}
public String getName() {
return name!=null ? name: beanName;
}
@Override
public void destroy() throws Exception {
@@ -50,13 +60,13 @@ public abstract class AbstractWANComponentFactoryBean<T> implements FactoryBean<
}
@Override
public final void setBeanName(String name) {
this.name = name;
public final void setBeanName(String beanName) {
this.beanName = beanName;
}
@Override
public final void afterPropertiesSet() throws Exception {
Assert.notNull(name, "Name cannot be null");
Assert.notNull(getName(), "Name cannot be null");
Assert.notNull(cache, "Cache cannot be null");
doInit();
}

View File

@@ -92,7 +92,7 @@ public class AsyncEventQueueFactoryBean extends AbstractWANComponentFactoryBean<
asyncEventQueueFactory.setMaximumQueueMemory(maximumQueueMemory);
}
asyncEventQueue = asyncEventQueueFactory.create(name, asyncEventListener);
asyncEventQueue = asyncEventQueueFactory.create(getName(), asyncEventListener);
}
@Override

View File

@@ -76,6 +76,7 @@ public class GatewayHubFactoryBean extends AbstractWANComponentFactoryBean<Gatew
@Override
protected void doInit() {
String name = getName();
gatewayHub = cache.addGatewayHub(name, port == null ? GatewayHub.DEFAULT_PORT : port);
if (log.isDebugEnabled()) {

View File

@@ -52,8 +52,6 @@ public class GatewayReceiverFactoryBean extends AbstractWANComponentFactoryBean<
*/
public GatewayReceiverFactoryBean(Cache cache) {
super(cache);
// Bean name not required.
this.name = "gateway-receiver";
}
@Override

View File

@@ -33,7 +33,7 @@ import com.gemstone.gemfire.cache.wan.GatewayTransportFilter;
* @author David Turanski
*
*/
public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<GatewaySender> {
public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<GatewaySender> {
private static List<String> validOrderPolicyValues = Arrays.asList("KEY", "PARTITION", "THREAD");
private GatewaySender gatewaySender;
@@ -58,7 +58,7 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
private Integer dispatcherThreads;
private Boolean manualStart;
private boolean manualStart = false;
private Integer maximumQueueMemory;
@@ -82,12 +82,12 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
@Override
public GatewaySender getObject() throws Exception {
return gatewaySender;
return new SmartLifecycleGatewaySender(gatewaySender, !manualStart);
}
@Override
public Class<?> getObjectType() {
return GatewaySender.class;
return SmartLifecycleGatewaySender.class;
}
@Override
@@ -151,9 +151,9 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
if (dispatcherThreads != null) {
gatewaySenderFactory.setDispatcherThreads(dispatcherThreads);
}
if (manualStart != null) {
gatewaySenderFactory.setManualStart(manualStart);
}
gatewaySenderFactory.setManualStart(true);
if (maximumQueueMemory != null) {
gatewaySenderFactory.setMaximumQueueMemory(maximumQueueMemory);
}
@@ -163,7 +163,7 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
if (socketReadTimeout != null) {
gatewaySenderFactory.setSocketReadTimeout(socketReadTimeout);
}
gatewaySender = gatewaySenderFactory.create(name, remoteDistributedSystemId);
gatewaySender = gatewaySenderFactory.create(getName(), remoteDistributedSystemId);
}
public void setRemoteDistributedSystemId(int remoteDistributedSystemId) {
@@ -233,5 +233,4 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
public void setSocketReadTimeout(Integer socketReadTimeout) {
this.socketReadTimeout = socketReadTimeout;
}
}

View File

@@ -0,0 +1,277 @@
/*
* Copyright 2002-2012 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
*
* http://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.wan;
import java.util.List;
import org.springframework.context.SmartLifecycle;
import org.springframework.util.Assert;
import com.gemstone.gemfire.cache.util.Gateway.OrderPolicy;
import com.gemstone.gemfire.cache.wan.GatewayEventFilter;
import com.gemstone.gemfire.cache.wan.GatewaySender;
import com.gemstone.gemfire.cache.wan.GatewayTransportFilter;
/**
* A {@link GatewaySender} controlled by {@link SmartLifecycle}
* @author David Turanski
*
*/
public class SmartLifecycleGatewaySender implements GatewaySender, SmartLifecycle {
private final GatewaySender delegate;
private final boolean autoStartup;
public SmartLifecycleGatewaySender(GatewaySender delegate, boolean autoStartup) {
Assert.notNull(delegate);
this.delegate = delegate;
this.autoStartup = autoStartup;
}
/* (non-Javadoc)
* @see com.gemstone.gemfire.cache.wan.GatewaySender#addGatewayEventFilter(com.gemstone.gemfire.cache.wan.GatewayEventFilter)
*/
@Override
public void addGatewayEventFilter(GatewayEventFilter eventFilter) {
this.delegate.addGatewayEventFilter(eventFilter);
}
/* (non-Javadoc)
* @see com.gemstone.gemfire.cache.wan.GatewaySender#getAlertThreshold()
*/
@Override
public int getAlertThreshold() {
return this.delegate.getAlertThreshold();
}
/* (non-Javadoc)
* @see com.gemstone.gemfire.cache.wan.GatewaySender#getBatchSize()
*/
@Override
public int getBatchSize() {
return this.delegate.getBatchSize();
}
/* (non-Javadoc)
* @see com.gemstone.gemfire.cache.wan.GatewaySender#getBatchTimeInterval()
*/
@Override
public int getBatchTimeInterval() {
return this.delegate.getBatchTimeInterval();
}
/* (non-Javadoc)
* @see com.gemstone.gemfire.cache.wan.GatewaySender#getDiskStoreName()
*/
@Override
public String getDiskStoreName() {
return this.delegate.getDiskStoreName();
}
/* (non-Javadoc)
* @see com.gemstone.gemfire.cache.wan.GatewaySender#getDispatcherThreads()
*/
@Override
public int getDispatcherThreads() {
return this.delegate.getDispatcherThreads();
}
/* (non-Javadoc)
* @see com.gemstone.gemfire.cache.wan.GatewaySender#getGatewayEventFilters()
*/
@Override
public List<GatewayEventFilter> getGatewayEventFilters() {
return this.delegate.getGatewayEventFilters();
}
/* (non-Javadoc)
* @see com.gemstone.gemfire.cache.wan.GatewaySender#getGatewayTransportFilters()
*/
@Override
public List<GatewayTransportFilter> getGatewayTransportFilters() {
return this.delegate.getGatewayTransportFilters();
}
/* (non-Javadoc)
* @see com.gemstone.gemfire.cache.wan.GatewaySender#getId()
*/
@Override
public String getId() {
return this.delegate.getId();
}
/* (non-Javadoc)
* @see com.gemstone.gemfire.cache.wan.GatewaySender#getMaximumQueueMemory()
*/
@Override
public int getMaximumQueueMemory() {
return this.delegate.getMaximumQueueMemory();
}
/* (non-Javadoc)
* @see com.gemstone.gemfire.cache.wan.GatewaySender#getOrderPolicy()
*/
@Override
public OrderPolicy getOrderPolicy() {
return this.delegate.getOrderPolicy();
}
/* (non-Javadoc)
* @see com.gemstone.gemfire.cache.wan.GatewaySender#getRemoteDSId()
*/
@Override
public int getRemoteDSId() {
return this.delegate.getRemoteDSId();
}
/* (non-Javadoc)
* @see com.gemstone.gemfire.cache.wan.GatewaySender#getSocketBufferSize()
*/
@Override
public int getSocketBufferSize() {
return this.delegate.getSocketBufferSize();
}
/* (non-Javadoc)
* @see com.gemstone.gemfire.cache.wan.GatewaySender#getSocketReadTimeout()
*/
@Override
public int getSocketReadTimeout() {
return this.delegate.getSocketReadTimeout();
}
/* (non-Javadoc)
* @see com.gemstone.gemfire.cache.wan.GatewaySender#isBatchConflationEnabled()
*/
@Override
public boolean isBatchConflationEnabled() {
return this.delegate.isBatchConflationEnabled();
}
/* (non-Javadoc)
* @see com.gemstone.gemfire.cache.wan.GatewaySender#isDiskSynchronous()
*/
@Override
public boolean isDiskSynchronous() {
return this.delegate.isDiskSynchronous();
}
/* (non-Javadoc)
* @see com.gemstone.gemfire.cache.wan.GatewaySender#isManualStart()
*/
@Override
public boolean isManualStart() {
return true;
}
/* (non-Javadoc)
* @see com.gemstone.gemfire.cache.wan.GatewaySender#isParallel()
*/
@Override
public boolean isParallel() {
return this.delegate.isParallel();
}
/* (non-Javadoc)
* @see com.gemstone.gemfire.cache.wan.GatewaySender#isPaused()
*/
@Override
public boolean isPaused() {
return this.delegate.isPaused();
}
/* (non-Javadoc)
* @see com.gemstone.gemfire.cache.wan.GatewaySender#isPersistenceEnabled()
*/
@Override
public boolean isPersistenceEnabled() {
return this.delegate.isPersistenceEnabled();
}
/* (non-Javadoc)
* @see com.gemstone.gemfire.cache.wan.GatewaySender#isRunning()
*/
@Override
public boolean isRunning() {
return this.delegate.isRunning();
}
/* (non-Javadoc)
* @see com.gemstone.gemfire.cache.wan.GatewaySender#pause()
*/
@Override
public void pause() {
this.delegate.pause();
}
/* (non-Javadoc)
* @see com.gemstone.gemfire.cache.wan.GatewaySender#removeGatewayEventFilter(com.gemstone.gemfire.cache.wan.GatewayEventFilter)
*/
@Override
public void removeGatewayEventFilter(GatewayEventFilter eventFilter) {
this.delegate.removeGatewayEventFilter(eventFilter);
}
/* (non-Javadoc)
* @see com.gemstone.gemfire.cache.wan.GatewaySender#resume()
*/
@Override
public void resume() {
this.delegate.resume();
}
/* (non-Javadoc)
* @see com.gemstone.gemfire.cache.wan.GatewaySender#start()
*/
@Override
public void start() {
this.delegate.start();
}
/* (non-Javadoc)
* @see com.gemstone.gemfire.cache.wan.GatewaySender#stop()
*/
@Override
public void stop() {
this.delegate.stop();
}
/* (non-Javadoc)
* @see org.springframework.context.Phased#getPhase()
*/
@Override
public int getPhase() {
return Integer.MAX_VALUE;
}
/* (non-Javadoc)
* @see org.springframework.context.SmartLifecycle#isAutoStartup()
*/
@Override
public boolean isAutoStartup() {
return this.autoStartup ;
}
/* (non-Javadoc)
* @see org.springframework.context.SmartLifecycle#stop(java.lang.Runnable)
*/
@Override
public void stop(Runnable callback) {
stop();
callback.run();
}
}

View File

@@ -2258,6 +2258,13 @@ Inner bean definition of the event filter
minOccurs="0" maxOccurs="1" />
</xsd:sequence>
<xsd:attributeGroup ref="commonWANQueueAttributes" />
<xsd:attribute name="name" type="xsd:string" use="optional">
<xsd:annotation>
<xsd:documentation><![CDATA[
Optionally specifies the GemFire gateway sender id. By default this value is the bean id or a generated value if an inner bean.
]]></xsd:documentation>
</xsd:annotation>
</xsd:attribute>
<xsd:attribute name="remote-distributed-system-id"
type="xsd:string" use="required">
<xsd:annotation>
@@ -2506,6 +2513,13 @@ use inner bean declarations.
</xsd:element>
</xsd:sequence>
<xsd:attributeGroup ref="commonWANQueueAttributes" />
<xsd:attribute name="name" type="xsd:string" use="optional">
<xsd:annotation>
<xsd:documentation><![CDATA[
Optionally specifies the GemFire gateway sender id. By default this value is the bean id or a generated value if an inner bean.
]]></xsd:documentation>
</xsd:annotation>
</xsd:attribute>
</xsd:complexType>
<!-- -->
<xsd:element name="gateway-sender">