The WAN replication feature allows 2 remote data centers, or 2 availability zones, to maintain data consistency. In the case where one data center cannot process incoming events for any reason, the other data center should retain the failed events so that no data is lost. Currently if data center 1 (DC1) is able to connect to data center 2 (DC2) and send it events, those events are removed from the queue on DC1 when the ack from DC2 is received, regardless of what happens to them on DC2. This behavior is controlled by the internal system property REMOVE_FROM_QUEUE_ON_EXCEPTION which defaults to true. Most common exceptions thrown from a receiving site include:
We will provide a mechanism for users to preserve events on the gateway sender that do not get successfully processed on the receiving data center. Our example implementation will store these events on disk at the sending data center and notify the user what events did not get transmitted.
Our current design approach is as follows:
Create new Java API
Define callback API for senders to set callback to dispatchers
If sender is configured with a callback, invoke the callback if batch exception occurs prior to batch removal
Implement a default callback API (see item 5 below)
Add properties on gateway receiver factory to specify # retries for a failed event and wait time between retries.
Modify Gfsh commands
Add option to gfsh ‘create gateway sender’ command to specify custom callback
Add options to gfsh ‘create gateway receiver’ command to set # retries and wait time between retries
Store new options in cluster config
Sender: callback implementation
Receiver: # of retries and wait time between retries
Add support in cache.xml for specifying new callback for gateway sender and setting new options for gateway receiver
Create example implementation of Sender callback that writes event(s) and associated exceptions to a file
Security features
Define privileges needed to deploy and configure sender callback
With security, callback should only write eventIds and exceptions, i.e. no entry values should be written to disk.
Add logging and statistics for callback
Log messages for gateway receiver for start time and results of retries
Add statistics and MBean for callbacks in-progress, completed, # and duration
New workflow for setting up WAN gateway using gfsh:
The new GatewayEventFailureListener interface is defined like:
public interface GatewayEventFailureListener extends CacheCallback {
/**
* Callback invoked on the GatewaySender when an event fails to be processed by the
* GatewayReceiver
*
* @param event The event that failed
*
* @param exception The exception that occurred
*/
void onFailure(GatewayQueueEvent event, Throwable exception);
} |
Example:
public class LoggingGatewayEventFailureListener implements GatewayEventFailureListener, Declarable {
private Cache cache;
public void onFailure(GatewayQueueEvent event, Throwable exception) {
this.cache.getLogger().warning("LoggingGatewayEventFailureListener onFailure: region=" + event.getRegion().getName() + "; operation=" + event.getOperation() + "; key=" + event.getKey() + "; value=" + event.getDeserializedValue() + "; exception=" + exception);
}
public void initialize(Cache cache, Properties properties) {
this.cache = cache;
}
} |
This LoggingGatewayEventFailureListener will log warnings like:
[warning 2018/11/05 17:30:41.613 PST ln-1 <AckReaderThread for : Event Processor for GatewaySender_ny_3> tid=0x75] LoggingGatewayEventFailureListener onFailure: region=data; operation=CREATE; key=8360; value=Trade[id=8360; cusip=PVTL; shares=100; price=18]; exception=org.apache.geode.cache.persistence.PartitionOfflineException: Region /data bucket 73 has persistent data that is no longer online stored at these locations: [...] |
The GatewaySenderFactory adds the ability to add a GatewayEventFailureListener:
/** * Sets the provided <code>GatewayEventFailureListener</code> in this GatewaySenderFactory. * * @param listener The <code>GatewayEventFailureListener</code> */ GatewaySenderFactory setGatewayEventFailureListener(GatewayEventFailureListener listener); |
The GatewaySender adds the ability to get a GatewayEventFailureListener:
/** * Returns this <code>GatewaySender's</code> <code>GatewayEventFailureListener</code>. * * @return this <code>GatewaySender's</code> <code>GatewayEventFailureListener</code> */ GatewayEventFailureListener getGatewayEventFailureListener(); |
Example:
GatewaySender sender = cache.createGatewaySenderFactory()
.setParallel(true)
.setGatewayEventFailureListener(new FileGatewayEventFailureListener(new File(...)))
.create("ln", 2); |
The GatewayReceiverFactory adds the ability to set retry attempts and wait time between retry attempts:
/** * Sets the number of retry attempts to apply failing events from remote GatewaySenders * * @param retryAttempts The retry attempts */ GatewayReceiverFactory setRetryAttempts(int retryAttempts); /** * Sets the wait time between retry attempts to apply failing events from remote GatewaySenders * * @param waitTimeBetweenRetryAttempts The wait time in milliseconds */ GatewayReceiverFactory setWaitTimeBetweenRetryAttempts(long waitTimeBetweenRetryAttempts); |
The GatewayReceiver adds the ability to get retry attempts and wait time between retry attempts:
/** * Returns the number of times to retry a failing event before throwing an exception. * * @return the number of times to retry a failing event before throwing an exception */ int getRetryAttempts(); /** * Returns the amount of time in milliseconds to wait between attempts to apply a failing event. * * @return the amount of time in milliseconds to wait between attempts to apply a failing event */ long getWaitTimeBetweenRetryAttempts(); |
Example:
GatewayReceiver receiver = cache.createGatewayReceiverFactory() .setRetryAttempts(10) .setWaitTimeBetweenRetryAttempts(100) .create(); |
The create gateway-sender command defines this new parameter:
| Name | Description |
|---|---|
| gateway-event-failure-listener | The fully qualified class name of GatewayEventFailureListener to be set in the GatewaySender |
Example:
Cluster-1 gfsh>create gateway-sender --id=ln --parallel=true --remote-distributed-system-id=2 --gateway-event-failure-listener=LoggingGatewayEventFailureListener Member | Status ------ | ------------------------------------ ny-1 | GatewaySender "ln" created on "ny-1" |
The create gateway-receiver command defines these new parameters:
| Name | Description |
|---|---|
| retry-attempts | The number of retry attempts for failed events processed by the GatewayReceiver |
| wait-time-between-retry-attempts | The amount of time to wait between retry attempts for failed events processed by the GatewayReceiver |
Example:
Cluster-2 gfsh>create gateway-receiver --retry-attempts=10 --wait-time-between-retry-attempts=100 Member | Status | Message ------ | ------ | --------------------------------------------------------------------------- ln-1 | OK | GatewayReceiver created on member "ln-1" and will listen on the port "5296" |
The <gateway-sender> element defines the <gateway-event-failure-listener> sub-element. The <gateway-event-failure-listener> sub-element is like any other Declarable.
Example:
<gateway-sender id="..."> <gateway-event-failure-listener> <class-name>FileGatewayEventFailureListener</class-name> </gateway-event-failure-listener> </gateway-sender> |
The <gateway-receiver> element defines the retry-attempts and wait-time-between-retry-attempts attributes.
Example:
<gateway-receiver retry-attempts="5" wait-time-between-retry-attempts="100"/> |
Risks and Unknowns
How to handle class not found exception for sender callback
entity EventProcessor as A
entity RemoteDispatcher as B
entity ServerConnection as C
entity ReceiverCommand as D
box "Site 1" #LightBlue
participant A
participant B
endbox
box "Site 2" #LightBlue
participant C
participant D
endbox
A -> A: peekBatchFromQueue
A -> B: dispatchBatch
B -> B: getConnection
B -> C: sendBatch
C -> C: readRequest
C -> C: createCommand
C -> D: execute
D -> D: readBatchEvents
loop For Each Batch Event
loop Retry
D -> D: determineOperation (create, update, destroy)
D -> D: executeOperation
alt Successful executeOperation:
D -> D: breakRetry
else Failed executeOperation:
alt Remove from queue on exception:
D -> D: storeException
D -> D: breakRetry
else Keep in queue on exception:
D -> D: sleep N milliseconds
D -> D: continueRetry
end
end
end
end
D -> B: sendAcknowledgement
B -> B: readAcknowledgement
B -> B: logExceptions (if necessary)
A -> A: removeBatchFromQueue |
entity EventProcessor as A
entity RemoteDispatcher as B
entity ServerConnection as C
entity ReceiverCommand as D
entity FailedEventHandler as E
box "Site 1" #LightBlue
participant A
participant B
participant E
endbox
box "Site 2" #LightBlue
participant C
participant D
endbox
A -> A: peekBatchFromQueue
A -> B: dispatchBatch
B -> B: getConnection
B -> C: sendBatch
C -> C: readRequest
C -> C: createCommand
C -> D: execute
D -> D: readBatchEvents
loop For Each Batch Event
loop Retry numberOfRetries
D -> D: determineOperation (create, update, destroy)
D -> D: executeOperation
alt Successful executeOperation:
D -> D: breakRetry
else Failed executeOperation:
D -> D: storeException
D -> D: sleep waitTimeBetweenRetries milliseconds
D -> D: continueRetry
end
end
end
D -> B: sendAcknowledgement
B -> B: readAcknowledgement
loop For Each Failed Batch Event
B -> E: onException
end
A -> A: removeBatchFromQueue |