* Uses finite state machine for three states: Closed, Open, Half-Open
* Closed state allows calls through, and on sequential failures exceeding the max# set - transitions to Open state. Intervening successes cause the failure count to reset to 0
* Open state throws a CircuitOpenException on every call until the reset timeout is reached which causes a transition to Half-Open state
* Half-Open state will allow the next single call through, if it succeeds - transition to Closed state, if it fails - transition back to Open state, starting the reset timer again
* Allow configuration for the call and reset timeouts, as well as the maximum number of sequential failures before opening
* Supports async or synchronous call protection
* Callbacks are supported for state entry into Closed, Open, Half-Open. These are run in the supplied execution context
* Both thrown exceptions and calls exceeding max call time are considered failures
* Uses akka scheduler for timer events
* Integrated into File-Based durable mailbox
* Sample documented for other durable mailboxes
83 lines
No EOL
2.7 KiB
Java
83 lines
No EOL
2.7 KiB
Java
/**
|
|
* Copyright (C) 2009-2012 Typesafe Inc. <http://www.typesafe.com>
|
|
*/
|
|
package docs.circuitbreaker;
|
|
|
|
//#imports1
|
|
|
|
import akka.actor.UntypedActor;
|
|
import akka.dispatch.Future;
|
|
import akka.event.LoggingAdapter;
|
|
import akka.util.Duration;
|
|
import akka.pattern.CircuitBreaker;
|
|
import akka.event.Logging;
|
|
|
|
import static akka.dispatch.Futures.future;
|
|
|
|
import java.util.concurrent.Callable;
|
|
|
|
//#imports1
|
|
|
|
//#circuit-breaker-initialization
|
|
public class DangerousJavaActor extends UntypedActor {
|
|
|
|
private final CircuitBreaker breaker;
|
|
private final LoggingAdapter log = Logging.getLogger(getContext().system(), this);
|
|
|
|
public DangerousJavaActor() {
|
|
this.breaker = new CircuitBreaker(
|
|
getContext().dispatcher(), getContext().system().scheduler(),
|
|
5, Duration.parse("10s"), Duration.parse("1m"))
|
|
.onOpen(new Callable<Object>() {
|
|
public Object call() throws Exception {
|
|
notifyMeOnOpen();
|
|
return null;
|
|
}
|
|
});
|
|
}
|
|
|
|
public void notifyMeOnOpen() {
|
|
log.warning("My CircuitBreaker is now open, and will not close for one minute");
|
|
}
|
|
//#circuit-breaker-initialization
|
|
|
|
//#circuit-breaker-usage
|
|
public String dangerousCall() {
|
|
return "This really isn't that dangerous of a call after all";
|
|
}
|
|
|
|
@Override
|
|
public void onReceive(Object message) {
|
|
if (message instanceof String) {
|
|
String m = (String) message;
|
|
if ("is my middle name".equals(m)) {
|
|
final Future<String> f = future(
|
|
new Callable<String>() {
|
|
public String call() {
|
|
return dangerousCall();
|
|
}
|
|
}, getContext().dispatcher());
|
|
|
|
getSender().tell(breaker
|
|
.callWithCircuitBreaker(
|
|
new Callable<Future<String>>() {
|
|
public Future<String> call() throws Exception {
|
|
return f;
|
|
}
|
|
}));
|
|
}
|
|
if ("block for me".equals(m)) {
|
|
getSender().tell(breaker
|
|
.callWithSyncCircuitBreaker(
|
|
new Callable<String>() {
|
|
@Override
|
|
public String call() throws Exception {
|
|
return dangerousCall();
|
|
}
|
|
}));
|
|
}
|
|
}
|
|
}
|
|
//#circuit-breaker-usage
|
|
|
|
} |