本文目录导读:

我来详细介绍Java服务容错案例的开发方法,涵盖常见的容错模式和实战示例。
项目框架搭建
Maven依赖配置
<dependencies>
<!-- Spring Boot -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<!-- Resilience4j 容错库 -->
<dependency>
<groupId>io.github.resilience4j</groupId>
<artifactId>resilience4j-spring-boot2</artifactId>
<version>2.0.2</version>
</dependency>
<dependency>
<groupId>io.github.resilience4j</groupId>
<artifactId>resilience4j-circuitbreaker</artifactId>
</dependency>
<dependency>
<groupId>io.github.resilience4j</groupId>
<artifactId>resilience4j-retry</artifactId>
</dependency>
<dependency>
<groupId>io.github.resilience4j</groupId>
<artifactId>resilience4j-bulkhead</artifactId>
</dependency>
<dependency>
<groupId>io.github.resilience4j</groupId>
<artifactId>resilience4j-ratelimiter</artifactId>
</dependency>
<!-- Hystrix(备用方案) -->
<dependency>
<groupId>org.springframework.cloud</groupId>
<artifactId>spring-cloud-starter-netflix-hystrix</artifactId>
<version>2.2.10.RELEASE</version>
</dependency>
</dependencies>
核心容错模式实现
1 熔断器模式 (Circuit Breaker)
import io.github.resilience4j.circuitbreaker.annotation.CircuitBreaker;
import org.springframework.stereotype.Service;
@Service
public class PaymentService {
private static final String PAYMENT_SERVICE = "paymentService";
@CircuitBreaker(name = PAYMENT_SERVICE, fallbackMethod = "fallbackPayment")
public String processPayment(String orderId, double amount) {
// 模拟调用远程支付服务
if (amount > 10000) {
throw new RuntimeException("支付服务异常:金额过大");
}
// 调用支付网关
return callPaymentGateway(orderId, amount);
}
// 熔断降级方法
public String fallbackPayment(String orderId, double amount, Throwable t) {
// 记录错误日志
log.error("支付服务熔断,订单:{},金额:{},错误:{}", orderId, amount, t.getMessage());
// 返回降级响应
return "支付服务暂时不可用,订单:" + orderId + " 已加入重试队列";
}
private String callPaymentGateway(String orderId, double amount) {
// 模拟支付网关调用
return "支付成功,订单:" + orderId;
}
}
2 重试机制 (Retry)
import io.github.resilience4j.retry.annotation.Retry;
import org.springframework.stereotype.Service;
@Service
public class OrderService {
private static final String ORDER_SERVICE = "orderService";
private int retryCount = 0;
@Retry(name = ORDER_SERVICE, fallbackMethod = "fallbackCreateOrder")
public Order createOrder(String userId, List<Product> products) {
retryCount++;
System.out.println("第 " + retryCount + " 次尝试创建订单");
// 模拟数据库操作
if (retryCount < 3) {
throw new DatabaseTimeoutException("数据库连接超时");
}
// 正常创建订单
return new Order(userId, products);
}
public Order fallbackCreateOrder(String userId, List<Product> products, Throwable t) {
System.out.println("创建订单失败,已重试 " + retryCount + " 次");
retryCount = 0;
// 返回一个默认的失败订单
return new Order(userId, products, OrderStatus.FAILED);
}
}
3 限流器 (Rate Limiter)
import io.github.resilience4j.ratelimiter.annotation.RateLimiter;
import org.springframework.stereotype.Service;
@Service
public class SmsService {
private static final String SMS_SERVICE = "smsService";
@RateLimiter(name = SMS_SERVICE, fallbackMethod = "fallbackSendSms")
public String sendSms(String phoneNumber, String message) {
// 调用短信服务商API
return smsProvider.sendSms(phoneNumber, message);
}
public String fallbackSendSms(String phoneNumber, String message, Throwable t) {
// 限流降级处理
return "短信发送频率过高,请稍后再试";
}
}
4 舱壁隔离 (Bulkhead)
import io.github.resilience4j.bulkhead.annotation.Bulkhead;
import org.springframework.stereotype.Service;
@Service
public class DatabaseService {
private static final String DB_SERVICE = "databaseService";
@Bulkhead(name = DB_SERVICE, type = Bulkhead.Type.THREADPOOL,
fallbackMethod = "fallbackQuery")
public List<User> queryUsers(String department) {
// 数据库查询操作
return userRepository.findByDepartment(department);
}
public List<User> fallbackQuery(String department, Throwable t) {
// 返回缓存数据或空列表
return cacheService.getUserCache(department);
}
}
配置文件
application.yml
resilience4j:
circuitbreaker:
instances:
paymentService:
registerHealthIndicator: true
slidingWindowSize: 10
minimumNumberOfCalls: 5
permittedNumberOfCallsInHalfOpenState: 3
automaticTransitionFromOpenToHalfOpenEnabled: true
waitDurationInOpenState: 5s
failureRateThreshold: 50
eventConsumerBufferSize: 10
retry:
instances:
orderService:
maxRetryAttempts: 3
waitDuration: 2s
retryExceptions:
- com.example.exception.DatabaseTimeoutException
- java.net.ConnectException
ratelimiter:
instances:
smsService:
limitForPeriod: 10
limitRefreshPeriod: 1s
timeoutDuration: 0s
bulkhead:
instances:
databaseService:
maxConcurrentCalls: 10
maxWaitDuration: 100ms
spring:
application:
name: fault-tolerance-demo
高级容错场景
1 缓存降级
@Component
public class CacheFallbackService {
@Autowired
private RedisTemplate<String, Object> redisTemplate;
@Cacheable(value = "product", key = "#productId")
@CircuitBreaker(name = "productService", fallbackMethod = "fallbackGetProduct")
public Product getProduct(String productId) {
// 从数据库加载
return productRepository.findById(productId)
.orElseThrow(() -> new ProductNotFoundException(productId));
}
public Product fallbackGetProduct(String productId, Throwable t) {
// 从Redis缓存获取
return (Product) redisTemplate.opsForValue()
.get("product:" + productId);
}
}
2 服务降级策略器
@Component
public class SmartDegradationStrategy {
private final Map<String, DegradationLevel> serviceStatus = new ConcurrentHashMap<>();
public enum DegradationLevel {
NORMAL(0), // 正常
DEGRADED(1), // 降级
LIMITED(2), // 限制
CLOSED(3); // 关闭
private final int level;
DegradationLevel(int level) { this.level = level; }
}
public Response<?> executeWithDegradation(String serviceName, Supplier<?> supplier) {
DegradationLevel level = serviceStatus.getOrDefault(serviceName, DegradationLevel.NORMAL);
return switch (level) {
case NORMAL -> executeNormal(supplier);
case DEGRADED -> executeDegraded(supplier);
case LIMITED -> executeLimited(supplier);
case CLOSED -> executeClosed(serviceName);
};
}
private Response<?> executeNormal(Supplier<?> supplier) {
try {
Object result = supplier.get();
return Response.success(result);
} catch (Exception e) {
return Response.failure("服务异常", e);
}
}
private Response<?> executeDegraded(Supplier<?> supplier) {
// 降级处理:只返回部分功能
return Response.success("服务降级中,提供简化服务");
}
private Response<?> executeClosed(String serviceName) {
// 服务关闭:返回友好提示
return Response.failure(serviceName + " 服务正在维护中");
}
}
监控和告警
1 健康检查端点
@RestController
@RequestMapping("/health")
public class HealthController {
@Autowired
private CircuitBreakerRegistry circuitBreakerRegistry;
@GetMapping("/circuitbreakers")
public Map<String, Object> getCircuitBreakers() {
Map<String, Object> result = new HashMap<>();
circuitBreakerRegistry.getAllCircuitBreakers().forEach(cb -> {
CircuitBreaker.Metrics metrics = cb.getMetrics();
Map<String, Object> breakerInfo = new HashMap<>();
breakerInfo.put("state", cb.getState());
breakerInfo.put("failureRate", metrics.getFailureRate());
breakerInfo.put("successfulCalls", metrics.getNumberOfSuccessfulCalls());
breakerInfo.put("failedCalls", metrics.getNumberOfFailedCalls());
result.put(cb.getName(), breakerInfo);
});
return result;
}
}
2 事件监听器
@Component
public class ResilienceEventListener {
@EventListener
public void handleCircuitBreakerEvent(CircuitBreakerEvent event) {
switch (event.getEventType()) {
case SUCCESS:
log.info("熔断器调用成功: {}", event.getCircuitBreakerName());
break;
case ERROR:
log.error("熔断器调用失败: {}", event.getCircuitBreakerName());
break;
case STATE_TRANSITION:
log.warn("熔断器状态变更: {} -> {}",
event.getCircuitBreakerName(),
event.getNewState());
// 发送告警通知
sendAlert(event);
break;
}
}
private void sendAlert(CircuitBreakerEvent event) {
// 实现告警通知逻辑(邮件、短信、钉钉等)
alertService.sendAlert(
"熔断器状态变更告警",
String.format("服务 %s 状态变更为 %s",
event.getCircuitBreakerName(),
event.getNewState())
);
}
}
测试用例
@SpringBootTest
class FaultToleranceTest {
@Autowired
private PaymentService paymentService;
@Test
void testCircuitBreaker() {
// 模拟大量失败请求
for (int i = 0; i < 20; i++) {
try {
if (i % 2 == 0) {
paymentService.processPayment("order" + i, 999);
} else {
paymentService.processPayment("order" + i, 10001);
}
} catch (Exception e) {
// 预期会有异常
}
}
// 验证熔断器状态
CircuitBreaker circuitBreaker = circuitBreakerRegistry
.circuitBreaker("paymentService");
assertThat(circuitBreaker.getState())
.isEqualTo(CircuitBreaker.State.OPEN);
}
@Test
void testRetryMechanism() {
Order order = orderService.createOrder("user1", products);
assertThat(order.getStatus()).isEqualTo(OrderStatus.FAILED);
// 验证重试次数
assertThat(retryCount).isEqualTo(3);
}
}
最佳实践建议
- 合理配置参数:根据业务特性调整熔断阈值、超时时间等
- 优雅的降级方案:提供有意义的降级响应,而非简单返回错误
- 渐进式容错:先限流、再重试、最后熔断
- 监控告警体系:建立完善的监控和告警机制
- 缓存配合:使用本地缓存或Redis缓存提升容错能力
这个案例覆盖了主要的服务容错模式,可以根据实际业务需求进行调整和扩展。