手动实现服务熔断功能
·
在我们分布式系统中, 经常会用用到 Hystrix 、Resilience4j 等主流服务熔断器,平时使用得非常熟悉,但是否有想过,自己手动实现一个,熔断器呢 ?今天我们就以前来基于Spring AOP实现一个服务熔断器。
注意:下面的例子我在本地启用了一个 consul 的注册中心,用于注册生产者及消费者
一、熔断器概念
- 在断路器对象中封装受保护的⽅法调⽤;
- 该对象监控调⽤和断路情况;
- 调⽤失败触发阈值后,后续调⽤直接由断路器返回错误,不再执⾏实际调⽤;
下图为熔断器的工作原理图:

二、服务提供方
2.1 pom主要依赖
<properties>
<java.version>1.8</java.version>
<spring-cloud.version>Greenwich.SR1</spring-cloud.version>
</properties>
<dependencyManagement>
<dependencies>
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-dependencies</artifactId>
<version>${spring-cloud.version}</version>
<type>pom</type>
<scope>import</scope>
</dependency>
</dependencies>
</dependencyManagement>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-actuator</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-starter-consul-discovery</artifactId>
</dependency>
2.2 配置文件
application.properties
#0表示服务器随机端口
server.port=0
#consul 地址
spring.cloud.consul.host=localhost
#consul 端口
spring.cloud.consul.port=8500
spring.cloud.consul.discovery.prefer-ip-address=true
bootstrap.properties
#服务名称
spring.application.name=waiter-service
2.3 代码配置
@RestController
@RequestMapping("/waiter")
public class WaiterController {
@PostMapping("/successKey")
@ResponseBody
public String getSuccessKey(){
return "hello wrold!";
}
}
三、服务消费方
3.1 pom主要依赖
<properties>
<java.version>1.8</java.version>
<spring-cloud.version>Greenwich.SR1</spring-cloud.version>
</properties>
<dependencyManagement>
<dependencies>
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-dependencies</artifactId>
<version>${spring-cloud.version}</version>
<type>pom</type>
<scope>import</scope>
</dependency>
</dependencies>
</dependencyManagement>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-actuator</artifactId>
</dependency>
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-starter-consul-discovery</artifactId>
</dependency>
3.2 配置文件
application.properties
#表示服务器端口
server.port=8090
#consul 地址
spring.cloud.consul.host=localhost
#consul 端口
spring.cloud.consul.port=8500
spring.cloud.consul.discovery.prefer-ip-address=true
bootstrap.properties
#服务名称
spring.application.name=customer-service
3.3 代码配置
1)、启用切面 @EnableAspectJAutoProxy
@SpringBootApplication
@EnableDiscoveryClient
@EnableFeignClients
@EnableAspectJAutoProxy
public class CustomerServiceApplication {
public static void main(String[] args) {
SpringApplication.run(CustomerServiceApplication.class, args);
}
@Bean
public CloseableHttpClient httpClient() {
return HttpClients.custom()
.setConnectionTimeToLive(30, TimeUnit.SECONDS)
.evictIdleConnections(30, TimeUnit.SECONDS)
.setMaxConnTotal(200)
.setMaxConnPerRoute(20)
.disableAutomaticRetries()
.setKeepAliveStrategy(new CustomConnectionKeepAliveStrategy())
.build();
}
}
2)、定义一个切面
功能:
- 在geektime.spring.springbucks.customer.integration包下的所有类进行切面拦截;
- counter 用于保存服务调用识别次数;
- breakCounter 表示熔断后被调用次数;
- 服务异常3次后,当前服务开始熔断返回null ,连续返回 3次以上null后开始,尝试方式服务;
@Aspect
@Component
@Slf4j
public class CircuitBreakerAspect {
private static final Integer THRESHOLD = 3;
private Map<String, AtomicInteger> counter = new ConcurrentHashMap<>();
private Map<String, AtomicInteger> breakCounter = new ConcurrentHashMap<>();
@Around("execution(* geektime.spring.springbucks.customer.integration..*(..))")
public Object doWithCircuitBreaker(ProceedingJoinPoint pjp) throws Throwable {
String signature = pjp.getSignature().toLongString();
log.info("Invoke {}", signature);
Object retVal;
try {
if (counter.containsKey(signature)) {
if (counter.get(signature).get() > THRESHOLD &&
breakCounter.get(signature).get() < THRESHOLD) {
log.warn("Circuit breaker return null, break {} times.",
breakCounter.get(signature).incrementAndGet());
return null;
}
} else {
counter.put(signature, new AtomicInteger(0));
breakCounter.put(signature, new AtomicInteger(0));
}
retVal = pjp.proceed();
counter.get(signature).set(0);
breakCounter.get(signature).set(0);
} catch (Throwable t) {
log.warn("Circuit breaker counter: {}, Throwable {}",
counter.get(signature).incrementAndGet(), t.getMessage());
breakCounter.get(signature).set(0);
throw t;
}
return retVal;
}
}
3) 通过feign 方式消费服务
@FeignClient(name = "waiter-service", contextId = "circuitBreaker")
public interface CircuitBreakerService {
@PostMapping("/waiter/successKey")
String getSuccessKey();
}
4)映射一个http 接口便于测试
@RestController
@RequestMapping("/customer")
public class CustomerController {
@Autowired
private CircuitBreakerService circuitBreakerService ;
@GetMapping("/key")
public Object getKey(){
return circuitBreakerService.getSuccessKey();
}
}
- 测试验证
正常消费服务

生产者服务调用异常(我把生产者服务stop了)

3次请求生产者接口失败后熔断返回null

四、总结
在上面的手动实现服务熔断逻辑中,如果能加上全局的熔断策略及熔断后默认的返回值,就比较完美了。
更多推荐

所有评论(0)