SGF-570 - Respect manual-start on Gateway Senders/Receivers but no longer couple the start/stop lifecycle to the Spring container.

(cherry picked from commit 1b9e65cdc1)
Signed-off-by: John Blum <jblum@pivotal.io>
This commit is contained in:
John Blum
2016-11-18 19:43:09 -08:00
parent 49b86a0ed0
commit 6d8bee2b34
6 changed files with 73 additions and 156 deletions

View File

@@ -15,20 +15,19 @@
*/
package org.springframework.data.gemfire.wan;
import java.io.IOException;
import java.util.List;
import org.apache.geode.cache.Cache;
import org.apache.geode.cache.wan.GatewayReceiver;
import org.apache.geode.cache.wan.GatewayReceiverFactory;
import org.apache.geode.cache.wan.GatewayTransportFilter;
import org.springframework.context.SmartLifecycle;
import org.springframework.beans.factory.FactoryBean;
import org.springframework.data.gemfire.util.CollectionUtils;
import org.springframework.util.Assert;
import org.springframework.util.CollectionUtils;
import org.springframework.util.StringUtils;
/**
* Spring FactoryBean for creating a GemFire {@link GatewayReceiver}.
* Spring {@link FactoryBean} for creating a GemFire {@link GatewayReceiver}.
*
* @author David Turanski
* @author John Blum
@@ -40,8 +39,7 @@ import org.springframework.util.StringUtils;
* @since 1.2.2
*/
@SuppressWarnings("unused")
public class GatewayReceiverFactoryBean extends AbstractWANComponentFactoryBean<GatewayReceiver>
implements SmartLifecycle {
public class GatewayReceiverFactoryBean extends AbstractWANComponentFactoryBean<GatewayReceiver> {
private boolean manualStart = false;
@@ -58,35 +56,39 @@ public class GatewayReceiverFactoryBean extends AbstractWANComponentFactoryBean<
private String hostnameForSenders;
/**
* Constructs an instance of the GatewayReceiverFactoryBean class for configuring an initializing
* a GemFire Gateway Receiver.
* Constructs an instance of the {@link GatewayReceiverFactoryBean} class initialized with a reference to
* the GemFire {@link Cache} used to configure and initialize a GemFire {@link GatewayReceiver}.
*
* @param cache a reference to the GemFire Cache used to setup the Gateway Receiver.
* @param cache reference to the GemFire {@link Cache} used to create the {@link GatewayReceiver}.
* @see org.apache.geode.cache.Cache
*/
public GatewayReceiverFactoryBean(Cache cache) {
super(cache);
}
/**
* @inheritDoc
*/
@Override
public GatewayReceiver getObject() throws Exception {
return gatewayReceiver;
}
/**
* @inheritDoc
*/
@Override
public Class<?> getObjectType() {
return (gatewayReceiver != null ? gatewayReceiver.getClass() : GatewayReceiver.class);
}
/**
* @inheritDoc
*/
@Override
protected void doInit() throws Exception {
GatewayReceiverFactory gatewayReceiverFactory = cache.createGatewayReceiverFactory();
if (!CollectionUtils.isEmpty(transportFilters)) {
for (GatewayTransportFilter transportFilter : transportFilters) {
gatewayReceiverFactory.addGatewayTransportFilter(transportFilter);
}
}
if (StringUtils.hasText(bindAddress)) {
gatewayReceiverFactory.setBindAddress(bindAddress);
}
@@ -95,28 +97,36 @@ public class GatewayReceiverFactoryBean extends AbstractWANComponentFactoryBean<
gatewayReceiverFactory.setHostnameForSenders(hostnameForSenders);
}
if (maximumTimeBetweenPings != null) {
gatewayReceiverFactory.setMaximumTimeBetweenPings(maximumTimeBetweenPings);
}
int localStartPort = defaultPort(startPort, GatewayReceiver.DEFAULT_START_PORT);
int localEndPort = defaultPort(endPort, GatewayReceiver.DEFAULT_END_PORT);
int localStartPort = (startPort != null ? startPort : GatewayReceiver.DEFAULT_START_PORT);
int localEndPort = (endPort != null ? endPort : GatewayReceiver.DEFAULT_END_PORT);
Assert.isTrue(localStartPort <= localEndPort, String.format("'startPort' must be less than or equal to %1$d.",
localEndPort));
Assert.isTrue(localStartPort <= localEndPort,
String.format("'startPort' must be less than or equal to %d.", localEndPort));
gatewayReceiverFactory.setStartPort(localStartPort);
gatewayReceiverFactory.setEndPort(localEndPort);
gatewayReceiverFactory.setManualStart(true);
gatewayReceiverFactory.setManualStart(manualStart);
if (maximumTimeBetweenPings != null) {
gatewayReceiverFactory.setMaximumTimeBetweenPings(maximumTimeBetweenPings);
}
if (socketBufferSize != null) {
gatewayReceiverFactory.setSocketBufferSize(socketBufferSize);
}
for (GatewayTransportFilter transportFilter : CollectionUtils.nullSafeList(transportFilters)) {
gatewayReceiverFactory.addGatewayTransportFilter(transportFilter);
}
gatewayReceiver = gatewayReceiverFactory.create();
}
public void setGatewayReceiver(final GatewayReceiver gatewayReceiver) {
protected int defaultPort(Integer port, int defaultPort) {
return (port != null ? port : defaultPort);
}
public void setGatewayReceiver(GatewayReceiver gatewayReceiver) {
this.gatewayReceiver = gatewayReceiver;
}
@@ -151,45 +161,4 @@ public class GatewayReceiverFactoryBean extends AbstractWANComponentFactoryBean<
public void setTransportFilters(List<GatewayTransportFilter> transportFilters) {
this.transportFilters = transportFilters;
}
@Override
public boolean isAutoStartup() {
return !manualStart;
}
@Override
public int getPhase() {
return Integer.MAX_VALUE;
}
@Override
public boolean isRunning() {
return gatewayReceiver.isRunning();
}
@Override
public void start() {
Assert.state(gatewayReceiver != null, "The GatewayReceiver was not properly configured and initialized!");
if (!isRunning()) {
try {
gatewayReceiver.start();
}
catch (IOException e) {
throw new RuntimeException("Failed to start Gateway Receiver due to I/O error.", e);
}
}
}
@Override
public void stop() {
gatewayReceiver.stop();
}
@Override
public void stop(final Runnable callback) {
stop();
callback.run();
}
}

View File

@@ -23,16 +23,15 @@ import org.apache.geode.cache.wan.GatewayEventSubstitutionFilter;
import org.apache.geode.cache.wan.GatewaySender;
import org.apache.geode.cache.wan.GatewaySenderFactory;
import org.apache.geode.cache.wan.GatewayTransportFilter;
import org.springframework.context.SmartLifecycle;
import org.springframework.beans.factory.FactoryBean;
import org.springframework.data.gemfire.util.CollectionUtils;
import org.springframework.util.Assert;
import org.springframework.util.CollectionUtils;
/**
* FactoryBean for creating a parallel or serial GemFire {@link GatewaySender}.
* Spring {@link FactoryBean} for creating a parallel or serial GemFire {@link GatewaySender}.
*
* @author David Turanski
* @author John Blum
* @see org.springframework.context.SmartLifecycle
* @see org.springframework.data.gemfire.wan.AbstractWANComponentFactoryBean
* @see org.apache.geode.cache.Cache
* @see org.apache.geode.cache.util.Gateway
@@ -41,8 +40,7 @@ import org.springframework.util.CollectionUtils;
* @since 1.2.2
*/
@SuppressWarnings("unused")
public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<GatewaySender>
implements SmartLifecycle {
public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<GatewaySender> {
private boolean manualStart = false;
@@ -73,25 +71,19 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
private String orderPolicy;
/**
* Constructs an instance of the GatewaySenderFactoryBean class initialized with a reference to the GemFire cache.
* Constructs an instance of the {@link GatewaySenderFactoryBean} class initialized with a reference to
* the GemFire {@link Cache} used to configured and initialized a GemFire {@link GatewaySender}.
*
* @param cache the Gemfire cache reference.
* @param cache reference to the GemFire {@link Cache} used to create the GemFire {@link GatewaySender}.
* @see org.apache.geode.cache.Cache
*/
public GatewaySenderFactoryBean(final Cache cache) {
public GatewaySenderFactoryBean(Cache cache) {
super(cache);
}
@Override
public GatewaySender getObject() throws Exception {
return gatewaySender;
}
@Override
public Class<?> getObjectType() {
return (gatewaySender != null ? gatewaySender.getClass() : GatewaySender.class);
}
/**
* @inheritDoc
*/
@Override
protected void doInit() {
GatewaySenderFactory gatewaySenderFactory = (this.factory != null ? (GatewaySenderFactory) factory
@@ -101,6 +93,10 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
gatewaySenderFactory.setAlertThreshold(alertThreshold);
}
if (batchConflationEnabled != null) {
gatewaySenderFactory.setBatchConflationEnabled(batchConflationEnabled);
}
if (batchSize != null) {
gatewaySenderFactory.setBatchSize(batchSize);
}
@@ -121,21 +117,15 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
gatewaySenderFactory.setDispatcherThreads(dispatcherThreads);
}
if (batchConflationEnabled != null) {
gatewaySenderFactory.setBatchConflationEnabled(batchConflationEnabled);
}
if (!CollectionUtils.isEmpty(eventFilters)) {
for (GatewayEventFilter eventFilter : eventFilters) {
gatewaySenderFactory.addGatewayEventFilter(eventFilter);
}
for (GatewayEventFilter eventFilter : CollectionUtils.nullSafeList(eventFilters)) {
gatewaySenderFactory.addGatewayEventFilter(eventFilter);
}
if (eventSubstitutionFilter != null) {
gatewaySenderFactory.setGatewayEventSubstitutionFilter(eventSubstitutionFilter);
}
gatewaySenderFactory.setManualStart(true);
gatewaySenderFactory.setManualStart(manualStart);
if (maximumQueueMemory != null) {
gatewaySenderFactory.setMaximumQueueMemory(maximumQueueMemory);
@@ -145,7 +135,7 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
Assert.isTrue(isSerialGatewaySender(), "Order Policy cannot be used with a Parallel Gateway Sender Queue.");
Assert.isTrue(VALID_ORDER_POLICIES.contains(orderPolicy.toUpperCase()),
String.format("The value for Order Policy '%1$s' is invalid.", orderPolicy));
String.format("The value for Order Policy '%s' is invalid.", orderPolicy));
gatewaySenderFactory.setOrderPolicy(GatewaySender.OrderPolicy.valueOf(orderPolicy.toUpperCase()));
}
@@ -161,10 +151,8 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
gatewaySenderFactory.setSocketReadTimeout(socketReadTimeout);
}
if (!CollectionUtils.isEmpty(transportFilters)) {
for (GatewayTransportFilter transportFilter : transportFilters) {
gatewaySenderFactory.addGatewayTransportFilter(transportFilter);
}
for (GatewayTransportFilter transportFilter : CollectionUtils.nullSafeList(transportFilters)) {
gatewaySenderFactory.addGatewayTransportFilter(transportFilter);
}
GatewaySenderWrapper wrapper = new GatewaySenderWrapper(gatewaySenderFactory.create(getName(),
@@ -174,6 +162,22 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
gatewaySender = wrapper;
}
/**
* @inheritDoc
*/
@Override
public GatewaySender getObject() throws Exception {
return gatewaySender;
}
/**
* @inheritDoc
*/
@Override
public Class<?> getObjectType() {
return (gatewaySender != null ? gatewaySender.getClass() : GatewaySender.class);
}
public void setAlertThreshold(Integer alertThreshold) {
this.alertThreshold = alertThreshold;
}
@@ -273,57 +277,4 @@ public class GatewaySenderFactoryBean extends AbstractWANComponentFactoryBean<Ga
public void setTransportFilters(List<GatewayTransportFilter> gatewayTransportFilters) {
this.transportFilters = gatewayTransportFilters;
}
/* (non-Javadoc)
* @see org.springframework.context.SmartLifecycle#isAutoStartup()
*/
@Override
public boolean isAutoStartup() {
return !manualStart;
}
/* (non-Javadoc)
* @see org.springframework.context.Phased#getPhase()
*/
@Override
public int getPhase() {
return Integer.MAX_VALUE;
}
/* (non-Javadoc)
* @see org.springframework.context.Lifecycle#isRunning()
*/
@Override
public boolean isRunning() {
return gatewaySender.isRunning();
}
/* (non-Javadoc)
* @see org.springframework.context.Lifecycle#start()
*/
@Override
public synchronized void start() {
Assert.notNull(gatewaySender, "The GatewaySender was not properly configured and initialized!");
if (!isRunning()){
gatewaySender.start();
}
}
/* (non-Javadoc)
* @see org.springframework.context.Lifecycle#stop()
*/
@Override
public void stop() {
gatewaySender.stop();
}
/* (non-Javadoc)
* @see org.springframework.context.SmartLifecycle#stop(java.lang.Runnable)
*/
@Override
public void stop(Runnable callback) {
stop();
callback.run();
}
}