本文目录导读:

我来详细解释Java分布式系统中的负载退避和负载均衡策略。
负载退避机制
指数退避算法
public class ExponentialBackoff {
private long baseDelay = 1000; // 基础延迟1秒
private long maxDelay = 60000; // 最大延迟60秒
private int maxRetries = 5;
public long getDelay(int attempt) {
// 2^attempt * baseDelay + random jitter
long delay = (long) Math.pow(2, attempt) * baseDelay;
delay = Math.min(delay, maxDelay);
// 添加随机抖动,避免惊群效应
long jitter = (long) (Math.random() * delay * 0.5);
return delay + jitter;
}
public void retryWithBackoff(Runnable task) {
for (int i = 0; i < maxRetries; i++) {
try {
task.run();
return; // 成功执行
} catch (Exception e) {
if (i == maxRetries - 1) {
throw new RuntimeException("Max retries exceeded", e);
}
try {
Thread.sleep(getDelay(i));
} catch (InterruptedException ie) {
Thread.currentThread().interrupt();
throw new RuntimeException("Retry interrupted", ie);
}
}
}
}
}
主流负载均衡算法
轮询算法
public class RoundRobinLoadBalancer {
private List<String> servers;
private AtomicInteger index = new AtomicInteger(0);
public RoundRobinLoadBalancer(List<String> servers) {
this.servers = servers;
}
public String getNextServer() {
int current = index.getAndUpdate(i -> (i + 1) % servers.size());
return servers.get(current);
}
}
加权轮询
public class WeightedRoundRobin {
private List<Server> servers;
private int totalWeight;
private int currentIndex = -1;
private int currentWeight = 0;
@Data
static class Server {
String address;
int weight; // 权重,越大表示处理能力越强
}
public synchronized String getServer() {
while (true) {
currentIndex = (currentIndex + 1) % servers.size();
if (currentIndex == 0) {
currentWeight--;
if (currentWeight <= 0) {
currentWeight = getMaxWeight();
if (currentWeight == 0) return null;
}
}
Server server = servers.get(currentIndex);
if (server.getWeight() >= currentWeight) {
return server.getAddress();
}
}
}
}
最少连接数
public class LeastConnectionBalancer {
private Map<String, AtomicInteger> connections = new ConcurrentHashMap<>();
public String getServer() {
return connections.entrySet().stream()
.min(Map.Entry.comparingByValue(
Comparator.comparingInt(AtomicInteger::get)))
.map(Map.Entry::getKey)
.orElseThrow(() -> new RuntimeException("No available servers"));
}
// 记录连接
public void acquireConnection(String server) {
connections.computeIfAbsent(server, k -> new AtomicInteger(0))
.incrementAndGet();
}
// 释放连接
public void releaseConnection(String server) {
connections.getOrDefault(server, new AtomicInteger(0))
.decrementAndGet();
}
}
一致性哈希
public class ConsistentHashLoadBalancer {
private TreeMap<Integer, String> hashRing = new TreeMap<>();
private int virtualNodes = 100; // 虚拟节点数
public ConsistentHashLoadBalancer(List<String> servers) {
for (String server : servers) {
addServer(server);
}
}
public void addServer(String server) {
// 为每个物理节点创建多个虚拟节点
for (int i = 0; i < virtualNodes; i++) {
int hash = hash(server + ":" + i);
hashRing.put(hash, server);
}
}
public String getServer(String key) {
if (hashRing.isEmpty()) return null;
int hash = hash(key);
// 找到大于等于hash值的最近节点
Map.Entry<Integer, String> entry = hashRing.ceilingEntry(hash);
if (entry == null) {
// 环的末尾,回到开头
entry = hashRing.firstEntry();
}
return entry.getValue();
}
private int hash(String key) {
// 使用MurmurHash或FNV哈希算法
return Math.abs(key.hashCode());
}
}
熔断器模式
public class CircuitBreaker {
private enum State { CLOSED, OPEN, HALF_OPEN }
private State state = State.CLOSED;
private int failureCount = 0;
private int threshold = 5; // 失败阈值
private long timeout = 30000; // 熔断时间
private long lastFailureTime;
private AtomicInteger requestCount = new AtomicInteger(0);
public <T> T execute(Supplier<T> operation) {
if (state == State.OPEN) {
if (System.currentTimeMillis() - lastFailureTime > timeout) {
state = State.HALF_OPEN;
} else {
throw new CircuitBreakerOpenException();
}
}
try {
T result = operation.get();
if (state == State.HALF_OPEN) {
state = State.CLOSED;
failureCount = 0;
}
return result;
} catch (Exception e) {
failureCount++;
lastFailureTime = System.currentTimeMillis();
if (failureCount >= threshold) {
state = State.OPEN;
}
throw e;
}
}
}
实际应用示例
Spring Cloud LoadBalancer
@Service
public class LoadBalancedService {
@Autowired
private LoadBalancerClient loadBalancer;
@Autowired
private RestTemplate restTemplate;
public String callService() {
ServiceInstance instance = loadBalancer.choose("service-name");
String url = instance.getUri() + "/api/endpoint";
ResponseEntity<String> response = restTemplate.exchange(
url,
HttpMethod.GET,
null,
String.class
);
return response.getBody();
}
}
自定义负载均衡配置
@Configuration
public class LoadBalanceConfig {
@Bean
public ReactorLoadBalancer<ServiceInstance> randomLoadBalancer(
Environment environment,
LoadBalancerClientFactory loadBalancerClientFactory) {
String name = environment.getProperty(LoadBalancerClientFactory.PROPERTY_NAME);
return new RandomLoadBalancer(
loadBalancerClientFactory.getLazyProvider(name, ServiceInstanceListSupplier.class),
name
);
}
}
性能监控与自适应
public class AdaptiveLoadBalancer {
private Map<String, ServerMetrics> metricsMap = new ConcurrentHashMap<>();
private static final double ALPHA = 0.3; // 平滑系数
@Data
static class ServerMetrics {
private double avgResponseTime;
private int activeConnections;
private double errorRate;
private double loadScore;
public void updateMetrics(long responseTime, boolean success) {
// 使用指数移动平均
avgResponseTime = ALPHA * responseTime + (1 - ALPHA) * avgResponseTime;
if (!success) {
errorRate = ALPHA * 1.0 + (1 - ALPHA) * errorRate;
} else {
errorRate = ALPHA * 0.0 + (1 - ALPHA) * errorRate;
}
// 综合负载评分
loadScore = avgResponseTime * 0.4 +
activeConnections * 0.3 +
errorRate * 1000 * 0.3;
}
}
public String selectOptimalServer() {
return metricsMap.entrySet().stream()
.min(Map.Entry.comparingByValue(
Comparator.comparingDouble(m -> m.getLoadScore())))
.map(Map.Entry::getKey)
.orElse(null);
}
}
建议的最佳实践
- 避免热点:使用一致性哈希配合虚拟节点
- 动态调整:基于实时性能数据动态调整权重
- 降级策略:实现熔断和降级机制
- 异步监控:使用独立线程收集性能指标
- 优雅关闭:实现平滑的节点摘除和加入
选择合适的负载均衡策略时,要考虑业务特性、服务器性能和网络环境等因素,必要时可以组合使用多种策略。