Class AbstractSingleFileGroupBacklog<T extends IReplicationOrderedPacket,CType extends AbstractSingleFileConfirmationHolder>
java.lang.Object
com.gigaspaces.internal.cluster.node.impl.backlog.AbstractSingleFileGroupBacklog<T,CType>
- All Implemented Interfaces:
IReplicationGroupBacklog,DynamicSourceGroupConfigHolder.IDynamicSourceGroupStateListener
- Direct Known Subclasses:
AbstractGlobalOrderGroupBacklog,AbstractMultiBucketSingleFileGroupBacklog,AbstractMultiSourceSingleFileGroupBacklog
public abstract class AbstractSingleFileGroupBacklog<T extends IReplicationOrderedPacket,CType extends AbstractSingleFileConfirmationHolder>
extends Object
implements IReplicationGroupBacklog, DynamicSourceGroupConfigHolder.IDynamicSourceGroupStateListener
A base class for
IReplicationGroupBacklog that contains a single IRedoLogFile and
treat the order of packets in that file for the global order of the backlog- Since:
- 8.0
- Author:
- eitany
-
Nested Class Summary
Nested ClassesModifier and TypeClassDescriptionclassclassclass -
Field Summary
FieldsModifier and TypeFieldDescriptionprotected final IPacketFilteredHandlerprotected final org.slf4j.Loggerprotected static final org.slf4j.Loggerprotected final org.slf4j.Loggerprotected final ReadWriteLock -
Constructor Summary
ConstructorsConstructorDescriptionAbstractSingleFileGroupBacklog(DynamicSourceGroupConfigHolder groupConfigHolder, String name, IReplicationPacketDataProducer<?> dataProducer) -
Method Summary
Modifier and TypeMethodDescriptionprotected voidvoidbeginSynchronizing(String memberName) voidbeginSynchronizing(String memberName, boolean isDirectPersistencySync) protected SynchronizingDatacheckSynchronizingDone(SynchronizingData synchronizingData, long currentKey, String memberName) protected voidcleanPendingErrorStateIfNeeded(String memberName, long packetKeykey, AbstractSingleFileConfirmationHolder confirmationHolder) protected voidvoidvoidclose()protected abstract TcreateBacklogOverflowPacket(long globalLastConfirmedKey, long firstKeyInBacklog, String memberName) createConfirmationMap(SourceGroupConfig groupConfig) protected abstract CTypevoiddecreaseMirrorDiscardedCount(long count) voiddecreaseWeight(String memberName, long lastConfirmedKey, long newlyConfirmedKey, AbstractSingleFileConfirmationHolder confirmationHolder) protected voiddecreaseWeightToAllMembersFromOldestPacket(long toKey) protected abstract voiddeleteBatchFromBacklog(long deletionBatchSize) protected voidensureLimit(IReplicationPacketData<?> data) protected TfilterPacketForSynchronizing(SynchronizingData synchronizingData, T packet, IPacketFilteredHandler filteredHandler, IReplicationPacketDataProducer dataProducer, org.slf4j.Logger logger, String memberName, PlatformLogicalVersion targetMemberVersion) intvoidprotected Collection<CType>getAllConfirmations(String... filterMembers) protected IRedoLogFile<T>protected CTypegetConfirmationHolderUnsafe(String memberName) getCurrentMarker(String memberName) Returns a marker of the current last position of the backlogprotected IPacketFilteredHandlerprotected longprotected longgetFirstRequiredKeyUnsafe(String memberName) protected <T extends SourceGroupConfig>
Tprotected longprotected abstract longgetLastConfirmedKeyUnsafe(String memberLookupName) protected longprotected StringgetMarker(IReplicationOrderedPacket packet, String membersGroupName) protected abstract longgetMemberUnconfirmedKey(CType value) protected longgetName()getNextHandshakeIteration(String memberName, IHandshakeContext handshakeContext) protected longgetPackets(String memberName, int maxSize, IReplicationChannelDataFilter filter, PlatformLogicalVersion targetMemberVersion, org.slf4j.Logger logger) getPacketsUnsafe(String memberName, int maxWeight, long upToKey, IReplicationChannelDataFilter dataFilter, IPacketFilteredHandler filteredHandler, PlatformLogicalVersion targetMemberVersion, org.slf4j.Logger logger) protected List<IReplicationOrderedPacket>getPacketsWithFullSerializedContent(String memberName, long fromKey, long upToKey, int maxWeight) getSpecificPacket(long packetKey) Gets a specific packet inside the backloggetSpecificPackets(long startPacketKey, long endPacketKey) getUnconfirmedMarker(String memberName) Returns a marker of the current unconfirmed packet of this given memberlonglonglonggetWeightUnsafe(String memberName) protected voidhandlePendingErrorBatchPackets(String memberName, List<IReplicationOrderedPacket> packets, Throwable error, long potentialLastUnprocessedKey) protected voidhandlePendingErrorSinglePacket(String memberName, IReplicationOrderedPacket packet, Throwable error) protected booleanbooleanprotected voidincreaseAllMembersWeight(long weight, long key) voidincreaseMirrorDiscardedCount(long count) voidincreaseWeight(String memberName, long weight, AbstractSingleFileConfirmationHolder confirmationHolder) protected voidinsertReplicationOrderedPacketToBacklog(T packet, ReplicationOutContext outContext) protected booleanbooleanisMarkerReached(String memberName, long markedKey) protected SynchronizingDataisSynchronizing(String memberName) protected voidlogPendingErrorResolved(String memberName, Throwable error) voidmakeMemberConfirmedOnAll(String memberName) Called once a new member addition is known to the available keepersvoidmemberAdded(MemberAddedEvent memberAddedParam, SourceGroupConfig newConfig) voidmemberRemoved(String memberName, SourceGroupConfig newConfig) voidmonitor(OperationWeightInfo info) protected voidnotifyDroppedIfListening(String member, BacklogConfig.LimitReachedPolicy policy) protected abstract voidonBeginSynchronization(String memberName) protected IReplicationOrderedPacketvoidvoidvoidprintRedoLog(String _name, String from) voidregisterWith(MetricRegistrator metricRegister) protected voidremoveSynchronizingState(long currentKey, String memberName) voidsetGroupHistory(IReplicationGroupHistory groupHistory) protected voidsetNextKeyUnsafe(long newNextKey) protected voidsetPacketWeight(IReplicationPacketData<?> data) voidsetPendingError(String memberName, Throwable error, IIdleStateData idleStateData) voidsetPendingError(String memberName, Throwable error, IReplicationOrderedPacket replicatedPacket) voidsetPendingError(String memberName, Throwable error, List<IReplicationOrderedPacket> replicatedPackets) voidsetStateListener(IReplicationBacklogStateListener stateListener) protected booleanshouldInsertPacket(IReplicationPacketData<?> data) longsize()longvoidstopSynchronization(String memberName) voidsynchronizationCopyStageDone(String memberName) booleansynchronizationDataGenerated(String memberName, String uid) voidsynchronizationDone(String memberName) protected longtakeNextKeyUnsafe(ReplicationOutContext replicationOutContext) toLogMessage(String memberName) protected StringtoString(List<IReplicationOrderedPacket> packets) protected voidupdateBacklogLimitations(SourceGroupConfig groupConfig) voidupdateMirrorWeightAfterCompaction(CompactionResult compactionResult) protected voidprotected voidvalidateReliableAsyncUpdateTargetsMatch(IReliableAsyncState reliableAsyncState, String sourceMemberName) voidMethods 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.backlog.IReplicationGroupBacklog
dumpState, fromWireForm, getConfirmed, getHandshakeRequest, getIdleStateData, getState, mergeWithDiscarded, processHandshakeResponse, processIdleStateDataResult, processResult, processResult, replaceWithDiscarded, supportDiscardMerge
-
Field Details
-
_loggerReplica
protected static final org.slf4j.Logger _loggerReplica -
_logger
protected final org.slf4j.Logger _logger -
_replicationLogger
protected final org.slf4j.Logger _replicationLogger -
_outOfSyncDueToDeletionTargets
-
_defaultFilteredHandler
-
_rwLock
-
-
Constructor Details
-
AbstractSingleFileGroupBacklog
public AbstractSingleFileGroupBacklog(DynamicSourceGroupConfigHolder groupConfigHolder, String name, IReplicationPacketDataProducer<?> dataProducer)
-
-
Method Details
-
updateBacklogLimitations
-
createConfirmationMap
-
hasExistingMember
protected boolean hasExistingMember() -
memberAdded
- Specified by:
memberAddedin interfaceDynamicSourceGroupConfigHolder.IDynamicSourceGroupStateListener
-
makeMemberConfirmedOnAll
Description copied from interface:IReplicationGroupBacklogCalled once a new member addition is known to the available keepers- Specified by:
makeMemberConfirmedOnAllin interfaceIReplicationGroupBacklog- Parameters:
memberName- new member name
-
memberRemoved
- Specified by:
memberRemovedin interfaceDynamicSourceGroupConfigHolder.IDynamicSourceGroupStateListener
-
getConfirmationHolderUnsafe
-
validateReliableAsyncUpdateTargetsMatch
protected void validateReliableAsyncUpdateTargetsMatch(IReliableAsyncState reliableAsyncState, String sourceMemberName) throws NoSuchReplicationMemberException, MissingReliableAsyncTargetStateException -
getMembersToValidateAgainst
-
getAllConfirmationHoldersUnsafe
-
getAllConfirmations
-
getAllConfirmations
-
createNewConfirmationHolder
-
getFirstKeyInBacklogInternal
protected long getFirstKeyInBacklogInternal() -
monitor
- Specified by:
monitorin interfaceIReplicationGroupBacklog- Throws:
RedoLogCapacityExceededException
-
isBacklogDroppedEntirely
protected boolean isBacklogDroppedEntirely() -
ensureLimit
-
notifyDroppedIfListening
-
peekOldestPacket
-
deleteBatchFromBacklog
protected abstract void deleteBatchFromBacklog(long deletionBatchSize) -
getInitialMaxAllowedDeleteUpTo
protected long getInitialMaxAllowedDeleteUpTo() -
getLastConfirmedKeyUnsafe
-
isSynchronizing
-
checkSynchronizingDone
protected SynchronizingData checkSynchronizingDone(SynchronizingData synchronizingData, long currentKey, String memberName) -
removeSynchronizingState
-
beginSynchronizing
- Specified by:
beginSynchronizingin interfaceIReplicationGroupBacklog
-
beginSynchronizing
- Specified by:
beginSynchronizingin interfaceIReplicationGroupBacklog
-
onBeginSynchronization
-
synchronizationDataGenerated
- Specified by:
synchronizationDataGeneratedin interfaceIReplicationGroupBacklog
-
synchronizationCopyStageDone
- Specified by:
synchronizationCopyStageDonein interfaceIReplicationGroupBacklog
-
synchronizationDone
- Specified by:
synchronizationDonein interfaceIReplicationGroupBacklog
-
stopSynchronization
- Specified by:
stopSynchronizationin interfaceIReplicationGroupBacklog
-
getPackets
public List<IReplicationOrderedPacket> getPackets(String memberName, int maxSize, IReplicationChannelDataFilter filter, PlatformLogicalVersion targetMemberVersion, org.slf4j.Logger logger) - Specified by:
getPacketsin interfaceIReplicationGroupBacklog
-
getPacketsUnsafe
public List<IReplicationOrderedPacket> getPacketsUnsafe(String memberName, int maxWeight, long upToKey, IReplicationChannelDataFilter dataFilter, IPacketFilteredHandler filteredHandler, PlatformLogicalVersion targetMemberVersion, org.slf4j.Logger logger) -
getFirstRequiredKeyUnsafe
-
createBacklogOverflowPacket
-
getFilteredHandler
-
filterPacketForSynchronizing
protected T filterPacketForSynchronizing(SynchronizingData synchronizingData, T packet, IPacketFilteredHandler filteredHandler, IReplicationPacketDataProducer dataProducer, org.slf4j.Logger logger, String memberName, PlatformLogicalVersion targetMemberVersion) -
clearReplicated
public void clearReplicated()- Specified by:
clearReplicatedin interfaceIReplicationGroupBacklog
-
clearConfirmedPackets
protected void clearConfirmedPackets() -
performCompaction
public void performCompaction()- Specified by:
performCompactionin interfaceIReplicationGroupBacklog
-
performCompactionUnsafe
public void performCompactionUnsafe() -
updateMirrorWeightAfterCompaction
-
decreaseMirrorDiscardedCount
public void decreaseMirrorDiscardedCount(long count) -
increaseMirrorDiscardedCount
public void increaseMirrorDiscardedCount(long count) -
hasMirror
public boolean hasMirror() -
getMinimumUnconfirmedKeyUnsafe
protected long getMinimumUnconfirmedKeyUnsafe() -
getMemberUnconfirmedKey
-
size
- Specified by:
sizein interfaceIReplicationGroupBacklog
-
size
public long size()- Specified by:
sizein interfaceIReplicationGroupBacklog
-
getCurrentMarker
Description copied from interface:IReplicationGroupBacklogReturns a marker of the current last position of the backlog- Specified by:
getCurrentMarkerin interfaceIReplicationGroupBacklog- Parameters:
memberName- the member we wish this marker will be attached to
-
getMarker
- Specified by:
getMarkerin interfaceIReplicationGroupBacklog
-
getNextKeyUnsafe
protected long getNextKeyUnsafe() -
getLastInsertedKeyToBacklogUnsafe
protected long getLastInsertedKeyToBacklogUnsafe() -
takeNextKeyUnsafe
-
setNextKeyUnsafe
protected void setNextKeyUnsafe(long newNextKey) -
getUnconfirmedMarker
Description copied from interface:IReplicationGroupBacklogReturns a marker of the current unconfirmed packet of this given member- Specified by:
getUnconfirmedMarkerin interfaceIReplicationGroupBacklog- Parameters:
memberName- the member we wish this marker will be attached to
-
isMarkerReached
-
toLogMessage
- Specified by:
toLogMessagein interfaceIReplicationGroupBacklog
-
validateIntegrity
protected void validateIntegrity() -
getNextHandshakeIteration
public IHandshakeIteration getNextHandshakeIteration(String memberName, IHandshakeContext handshakeContext) - Specified by:
getNextHandshakeIterationin interfaceIReplicationGroupBacklog
-
getStatistics
- Specified by:
getStatisticsin interfaceIReplicationGroupBacklog
-
flushRedoLogToStorage
public int flushRedoLogToStorage()- Specified by:
flushRedoLogToStoragein interfaceIReplicationGroupBacklog- Returns:
- number of flushed packets
-
getSwapStorageType
- Specified by:
getSwapStorageTypein interfaceIReplicationGroupBacklog
-
registerWith
- Specified by:
registerWithin interfaceIReplicationGroupBacklog
-
close
public void close()- Specified by:
closein interfaceIReplicationGroupBacklog
-
getName
-
getGroupName
-
getBacklogFile
-
getDataProducer
- Specified by:
getDataProducerin interfaceIReplicationGroupBacklog- Returns:
- the data producer this group backlog is using to generate replication data
-
setGroupHistory
- Specified by:
setGroupHistoryin interfaceIReplicationGroupBacklog
-
setStateListener
- Specified by:
setStateListenerin interfaceIReplicationGroupBacklog
-
setPendingError
- Specified by:
setPendingErrorin interfaceIReplicationGroupBacklog
-
setPendingError
public void setPendingError(String memberName, Throwable error, IReplicationOrderedPacket replicatedPacket) - Specified by:
setPendingErrorin interfaceIReplicationGroupBacklog
-
setPendingError
public void setPendingError(String memberName, Throwable error, List<IReplicationOrderedPacket> replicatedPackets) - Specified by:
setPendingErrorin interfaceIReplicationGroupBacklog
-
getLogPrefix
-
logPendingErrorResolved
-
handlePendingErrorBatchPackets
protected void handlePendingErrorBatchPackets(String memberName, List<IReplicationOrderedPacket> packets, Throwable error, long potentialLastUnprocessedKey) -
toString
-
handlePendingErrorSinglePacket
protected void handlePendingErrorSinglePacket(String memberName, IReplicationOrderedPacket packet, Throwable error) -
cleanPendingErrorStateIfNeeded
protected void cleanPendingErrorStateIfNeeded(String memberName, long packetKeykey, AbstractSingleFileConfirmationHolder confirmationHolder) -
setPacketWeight
-
shouldInsertPacket
-
getGroupConfigSnapshot
-
appendConfirmationStateString
-
getSpecificPacket
Description copied from interface:IReplicationGroupBacklogGets a specific packet inside the backlog- Specified by:
getSpecificPacketin interfaceIReplicationGroupBacklog
-
getSpecificPackets
-
getPacketsWithFullSerializedContent
protected List<IReplicationOrderedPacket> getPacketsWithFullSerializedContent(String memberName, long fromKey, long upToKey, int maxWeight) -
insertReplicationOrderedPacketToBacklog
-
writeLock
public void writeLock()- Specified by:
writeLockin interfaceIReplicationGroupBacklog
-
freeWriteLock
public void freeWriteLock()- Specified by:
freeWriteLockin interfaceIReplicationGroupBacklog
-
getWeight
public long getWeight()- Specified by:
getWeightin interfaceIReplicationGroupBacklog
-
getWeight
- Specified by:
getWeightin interfaceIReplicationGroupBacklog
-
getWeightUnsafe
-
increaseWeight
public void increaseWeight(String memberName, long weight, AbstractSingleFileConfirmationHolder confirmationHolder) - Specified by:
increaseWeightin interfaceIReplicationGroupBacklog
-
decreaseWeight
public void decreaseWeight(String memberName, long lastConfirmedKey, long newlyConfirmedKey, AbstractSingleFileConfirmationHolder confirmationHolder) - Specified by:
decreaseWeightin interfaceIReplicationGroupBacklog
-
decreaseWeightToAllMembersFromOldestPacket
protected void decreaseWeightToAllMembersFromOldestPacket(long toKey) -
increaseAllMembersWeight
protected void increaseAllMembersWeight(long weight, long key) -
printRedoLog
-