-
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)
- Loading branch information
Mahesh Subramanian
committed
Jul 16, 2019
1 parent
1fe0610
commit 1892fcb
Showing
13 changed files
with
242 additions
and
74 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
77 changes: 77 additions & 0 deletions
77
...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,77 @@ | ||
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.eventsourcing.repository.jdbc.exception.OptimisticLockingRetryException; | ||
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.concurrent.TimeUnit; | ||
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 | ||
@Value(key = "maxRetryStreamOperation", defaultValue = "10") | ||
private long maxRetry; | ||
|
||
@Inject | ||
private EventSourceTransformation eventSourceTransformation; | ||
|
||
public void appendWithRetry(final UUID streamId, final EventStream eventStream, final Stream<JsonEnvelope> events) throws EventStreamException { | ||
executeStreamOperationWithRetry(streamId, () -> { | ||
eventStream.append(events); | ||
logger.info("Appended events to stream with ID - '{}'", streamId); | ||
}); | ||
|
||
} | ||
|
||
public void cloneWithRetry(final UUID streamId) throws EventStreamException { | ||
executeStreamOperationWithRetry(streamId, () -> { | ||
final UUID clonedStreamId = eventSourceTransformation.cloneStream(streamId); | ||
logger.info("Created backup stream '{}' from stream '{}'", clonedStreamId, streamId); | ||
}); | ||
} | ||
|
||
private void executeStreamOperationWithRetry(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; | ||
} | ||
sleep(); | ||
logger.warn(format("Encountered exception whilst completing operation on stream during attempt %s for stream with ID: %s", retryCount, streamId), e); | ||
} | ||
} | ||
|
||
} | ||
|
||
private void sleep() { | ||
try { | ||
TimeUnit.SECONDS.sleep(2); | ||
} catch (InterruptedException e) { | ||
logger.error("Thread sleep interrupted", e); | ||
Thread.currentThread().interrupt(); | ||
} | ||
} | ||
} |
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; | ||
} | ||
|
||
|
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
Oops, something went wrong.