Class AbstractMultiBucketSingleFileTargetProcessLog
java.lang.Object
com.gigaspaces.internal.cluster.node.impl.processlog.AbstractSingleFileTargetProcessLog
com.gigaspaces.internal.cluster.node.impl.processlog.multibucketsinglefile.AbstractMultiBucketSingleFileTargetProcessLog
- All Implemented Interfaces:
IReplicationTargetProcessLog,IMultiBucketSingleFileProcessLog
- Direct Known Subclasses:
MultiBucketSingleFileBatchConsumeTargetProcessLog,MultiBucketSingleFileSyncTargetProcessLog
public abstract class AbstractMultiBucketSingleFileTargetProcessLog
extends AbstractSingleFileTargetProcessLog
implements IMultiBucketSingleFileProcessLog
-
Field Summary
Fields inherited from class com.gigaspaces.internal.cluster.node.impl.processlog.AbstractSingleFileTargetProcessLog
_dataConsumer, _specificLogger -
Constructor Summary
ConstructorsModifierConstructorDescriptionprotectedAbstractMultiBucketSingleFileTargetProcessLog(ProcessLogConfig config, IReplicationPacketDataConsumer<?> dataConsumer, IReplicationProcessLogExceptionHandler exceptionHandler, IReplicationInFacade replicationInFacade, String name, String groupName, String sourceLookupName, long[] lastProcessedKeys, long[] lastGlobalProcessedKeys, boolean firstHandshakeForTarget, IReplicationGroupHistory groupHistory) AbstractMultiBucketSingleFileTargetProcessLog(ProcessLogConfig config, IReplicationPacketDataConsumer<?> dataConsumer, IReplicationProcessLogExceptionHandler exceptionHandler, IReplicationInFacade replicationInFacade, String name, String groupName, String sourceLookupName, IReplicationGroupHistory groupHistory) -
Method Summary
Modifier and TypeMethodDescriptionprotected voidafterSuccessfulConsumption(String sourceLookupName, IReplicationOrderedPacket packet) protected booleanbooleanvoidcreateBatchParallelProcessingContinuationTask(String sourceLookupName, IReplicationInFilterCallback inFilterCallback, ParallelBatchProcessingContext context, List<IReplicationOrderedPacket> batch, int segmentIndex) protected Stringlonglong[]long[]org.slf4j.LoggerbooleanperformHandshake(String memberName, IBacklogHandshakeRequest handshakeRequest) process(String sourceLookupName, DeletedMultiBucketOrderedPacket packet, IReplicationInFilterCallback inFilterCallback, ParallelBatchProcessingContext batchContext, int segmentIndex) process(String sourceLookupName, DiscardedMultiBucketOrderedPacket packet, IReplicationInFilterCallback inFilterCallback, ParallelBatchProcessingContext batchContext, int segmentIndex) process(String sourceLookupName, ISingleBucketReplicationOrderedPacket packet, IReplicationInFilterCallback inFilterCallback, ParallelBatchProcessingContext batchContext, int segmentIndex) process(String sourceLookupName, MultipleBucketOrderedPacket packet, IReplicationInFilterCallback inFilterCallback, ParallelBatchProcessingContext batchContext, int segmentIndex) process(String sourceLookupName, IReplicationOrderedPacket packet, IReplicationInFilterCallback inFilterCallback) process(String sourceLookupName, IReplicationOrderedPacket packet, IReplicationInFilterCallback inFilterCallback, ParallelBatchProcessingContext batchContext, int segmentIndex) processBatch(String sourceLookupName, List<IReplicationOrderedPacket> packets, IReplicationInFilterCallback inFilterCallback) voidprocessHandshakeIteration(String sourceMemberName, IHandshakeIteration handshakeIteration) processIdleStateData(String string, IIdleStateData idleStateData, IReplicationInFilterCallback inFilterCallback) resync(IBacklogHandshakeRequest handshakeRequest) protected booleanprotected voidtoWireForm(IProcessResult processResult) voidvoidMethods inherited from class com.gigaspaces.internal.cluster.node.impl.processlog.AbstractSingleFileTargetProcessLog
contentRequiredWhileProcessing, createReplicationInContext, getDataConsumer, getExceptionHandler, getGroupHistory, getGroupName, getReplicationInContext, getReplicationInFacade, getSourceLookupName, throwIfRepetitiveError, toLogMessageMethods inherited from class java.lang.Object
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, waitMethods inherited from interface com.gigaspaces.internal.cluster.node.impl.processlog.IReplicationTargetProcessLog
getDataConsumer, toLogMessage
-
Constructor Details
-
AbstractMultiBucketSingleFileTargetProcessLog
public AbstractMultiBucketSingleFileTargetProcessLog(ProcessLogConfig config, IReplicationPacketDataConsumer<?> dataConsumer, IReplicationProcessLogExceptionHandler exceptionHandler, IReplicationInFacade replicationInFacade, String name, String groupName, String sourceLookupName, IReplicationGroupHistory groupHistory) -
AbstractMultiBucketSingleFileTargetProcessLog
protected AbstractMultiBucketSingleFileTargetProcessLog(ProcessLogConfig config, IReplicationPacketDataConsumer<?> dataConsumer, IReplicationProcessLogExceptionHandler exceptionHandler, IReplicationInFacade replicationInFacade, String name, String groupName, String sourceLookupName, long[] lastProcessedKeys, long[] lastGlobalProcessedKeys, boolean firstHandshakeForTarget, IReplicationGroupHistory groupHistory)
-
-
Method Details
-
getExecutorService
-
getLastProcessedKeys
public long[] getLastProcessedKeys() -
getLastGlobalProcessedKeys
public long[] getLastGlobalProcessedKeys() -
isFirstHandshakeForTarget
public boolean isFirstHandshakeForTarget() -
performHandshake
public MultiBucketSingleFileHandshakeResponse performHandshake(String memberName, IBacklogHandshakeRequest handshakeRequest) throws IncomingReplicationOutOfSyncException - Specified by:
performHandshakein interfaceIReplicationTargetProcessLog- Throws:
IncomingReplicationOutOfSyncException
-
canResetState
protected boolean canResetState() -
validateOpen
public void validateOpen() -
validateNotClosed
public void validateNotClosed() -
throwClosedException
protected void throwClosedException() -
processBatch
public MultiBucketSingleFileProcessResult processBatch(String sourceLookupName, List<IReplicationOrderedPacket> packets, IReplicationInFilterCallback inFilterCallback) - Specified by:
processBatchin interfaceIReplicationTargetProcessLog
-
createBatchParallelProcessingContinuationTask
public void createBatchParallelProcessingContinuationTask(String sourceLookupName, IReplicationInFilterCallback inFilterCallback, ParallelBatchProcessingContext context, List<IReplicationOrderedPacket> batch, int segmentIndex) -
process
public MultiBucketSingleFileProcessResult process(String sourceLookupName, IReplicationOrderedPacket packet, IReplicationInFilterCallback inFilterCallback, ParallelBatchProcessingContext batchContext, int segmentIndex) -
process
public MultiBucketSingleFileProcessResult process(String sourceLookupName, IReplicationOrderedPacket packet, IReplicationInFilterCallback inFilterCallback) - Specified by:
processin interfaceIReplicationTargetProcessLog
-
process
public MultiBucketSingleFileProcessResult process(String sourceLookupName, ISingleBucketReplicationOrderedPacket packet, IReplicationInFilterCallback inFilterCallback, ParallelBatchProcessingContext batchContext, int segmentIndex) - Specified by:
processin interfaceIMultiBucketSingleFileProcessLog
-
process
public MultiBucketSingleFileProcessResult process(String sourceLookupName, MultipleBucketOrderedPacket packet, IReplicationInFilterCallback inFilterCallback, ParallelBatchProcessingContext batchContext, int segmentIndex) - Specified by:
processin interfaceIMultiBucketSingleFileProcessLog
-
process
public MultiBucketSingleFileProcessResult process(String sourceLookupName, DiscardedMultiBucketOrderedPacket packet, IReplicationInFilterCallback inFilterCallback, ParallelBatchProcessingContext batchContext, int segmentIndex) - Specified by:
processin interfaceIMultiBucketSingleFileProcessLog
-
process
public MultiBucketSingleFileProcessResult process(String sourceLookupName, DeletedMultiBucketOrderedPacket packet, IReplicationInFilterCallback inFilterCallback, ParallelBatchProcessingContext batchContext, int segmentIndex) - Specified by:
processin interfaceIMultiBucketSingleFileProcessLog
-
close
- Specified by:
closein interfaceIReplicationTargetProcessLog- Throws:
InterruptedException
-
getConsumeTimeout
public long getConsumeTimeout() -
processHandshakeIteration
public void processHandshakeIteration(String sourceMemberName, IHandshakeIteration handshakeIteration) - Specified by:
processHandshakeIterationin interfaceIReplicationTargetProcessLog
-
shouldCloneOnFilter
protected boolean shouldCloneOnFilter() -
afterSuccessfulConsumption
protected void afterSuccessfulConsumption(String sourceLookupName, IReplicationOrderedPacket packet) -
resync
- Specified by:
resyncin interfaceIReplicationTargetProcessLog
-
toWireForm
- Specified by:
toWireFormin interfaceIReplicationTargetProcessLog
-
dumpState
- Specified by:
dumpStatein interfaceIReplicationTargetProcessLog- Specified by:
dumpStatein classAbstractSingleFileTargetProcessLog
-
dumpStateExtra
-
getSpecificLogger
public org.slf4j.Logger getSpecificLogger() -
processIdleStateData
public IProcessResult processIdleStateData(String string, IIdleStateData idleStateData, IReplicationInFilterCallback inFilterCallback) - Specified by:
processIdleStateDatain interfaceIReplicationTargetProcessLog
-