-
Notifications
You must be signed in to change notification settings - Fork 3
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
SCTP-341: Addressing retry mechanism when performing operations on a …
…stream (append or clone) (#33)
- Loading branch information
Mahesh Subramanian
committed
Jul 17, 2019
1 parent
1fe0610
commit 1fff5e4
Showing
17 changed files
with
316 additions
and
76 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
43 changes: 43 additions & 0 deletions
43
...n/java/uk/gov/justice/tools/eventsourcing/transformation/service/RetryStreamOperator.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,43 @@ | ||
package uk.gov.justice.tools.eventsourcing.transformation.service; | ||
|
||
import uk.gov.justice.services.eventsourcing.source.core.EventSourceTransformation; | ||
import uk.gov.justice.services.eventsourcing.source.core.EventStream; | ||
import uk.gov.justice.services.eventsourcing.source.core.exception.EventStreamException; | ||
import uk.gov.justice.services.messaging.JsonEnvelope; | ||
|
||
import java.util.UUID; | ||
import java.util.stream.Stream; | ||
|
||
import javax.enterprise.context.ApplicationScoped; | ||
import javax.inject.Inject; | ||
|
||
import org.slf4j.Logger; | ||
|
||
@ApplicationScoped | ||
public class RetryStreamOperator { | ||
|
||
@Inject | ||
private Logger logger; | ||
|
||
@Inject | ||
private StreamOperationRetryableExecutor streamOperationRetryableExecutor; | ||
|
||
@Inject | ||
private EventSourceTransformation eventSourceTransformation; | ||
|
||
public void appendWithRetry(final UUID streamId, final EventStream eventStream, final Stream<JsonEnvelope> events) throws EventStreamException { | ||
streamOperationRetryableExecutor.execute(streamId, () -> { | ||
eventStream.append(events); | ||
logger.info("Appended events to stream with ID - '{}'", streamId); | ||
}); | ||
|
||
} | ||
|
||
public void cloneWithRetry(final UUID streamId) throws EventStreamException { | ||
streamOperationRetryableExecutor.execute(streamId, () -> { | ||
final UUID clonedStreamId = eventSourceTransformation.cloneStream(streamId); | ||
logger.info("Created backup stream '{}' from stream '{}'", clonedStreamId, streamId); | ||
}); | ||
} | ||
|
||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
11 changes: 11 additions & 0 deletions
11
...a/uk/gov/justice/tools/eventsourcing/transformation/service/StreamOperationRetryable.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,11 @@ | ||
package uk.gov.justice.tools.eventsourcing.transformation.service; | ||
|
||
import uk.gov.justice.services.eventsourcing.source.core.exception.EventStreamException; | ||
|
||
@FunctionalInterface | ||
public interface StreamOperationRetryable { | ||
|
||
void execute() throws EventStreamException; | ||
} | ||
|
||
|
50 changes: 50 additions & 0 deletions
50
.../justice/tools/eventsourcing/transformation/service/StreamOperationRetryableExecutor.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,50 @@ | ||
package uk.gov.justice.tools.eventsourcing.transformation.service; | ||
|
||
import static java.lang.String.format; | ||
|
||
import uk.gov.justice.services.common.configuration.Value; | ||
import uk.gov.justice.services.common.util.Sleeper; | ||
import uk.gov.justice.services.eventsourcing.repository.jdbc.exception.OptimisticLockingRetryException; | ||
import uk.gov.justice.services.eventsourcing.source.core.exception.EventStreamException; | ||
|
||
import java.util.UUID; | ||
|
||
import javax.enterprise.context.ApplicationScoped; | ||
import javax.inject.Inject; | ||
|
||
import org.slf4j.Logger; | ||
|
||
@ApplicationScoped | ||
public class StreamOperationRetryableExecutor { | ||
|
||
private static final int SLEEP_TIME_IN_MILLISECONDS = 2000; | ||
|
||
@Inject | ||
@Value(key = "maxRetryStreamOperation", defaultValue = "10") | ||
private long maxRetry; | ||
|
||
@Inject | ||
private Logger logger; | ||
|
||
@Inject | ||
private Sleeper sleeper; | ||
|
||
public void execute(final UUID streamId, final StreamOperationRetryable streamOperationRetryable) throws EventStreamException { | ||
boolean operationCompletedSuccesfully = false; | ||
long retryCount = 0L; | ||
while (!operationCompletedSuccesfully) { | ||
try { | ||
streamOperationRetryable.execute(); | ||
operationCompletedSuccesfully = true; | ||
} catch (OptimisticLockingRetryException e) { | ||
retryCount++; | ||
if (retryCount >= maxRetry) { | ||
logger.error("Failed to complete operation on stream '{}' due to concurrency issues. Exhausted all retries.", streamId); | ||
throw e; | ||
} | ||
sleeper.sleepFor(SLEEP_TIME_IN_MILLISECONDS); | ||
logger.warn(format("Encountered exception whilst completing operation on stream during attempt %s for stream with ID: %s", retryCount, streamId), e); | ||
} | ||
} | ||
} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.