在我们分布式系统中, 经常会用用到 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();
        }
    }
  1. 测试验证
    正常消费服务
    在这里插入图片描述
    生产者服务调用异常(我把生产者服务stop了)
    在这里插入图片描述
    3次请求生产者接口失败后熔断返回null
    在这里插入图片描述

四、总结

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

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐