GRPC自定义负载均衡器无法检测集群新增服务器,Java客户端临时方案失效
问题描述
我正在构建一个分布式工作流编排器,Worker通过GRPC与服务器集群通信。当集群中新增服务器时,GRPC客户端无法检测到这一变化。我尝试在服务器配置中添加MaxConnectionAge作为临时解决方案:
grpc.KeepaliveParams(keepalive.ServerParameters{ MaxConnectionAge: time.Minute * 1, })
我们有Golang和Java两种Worker实现,该方案在Golang客户端中运行正常——客户端每分钟新建连接,能检测到集群新服务器,但在Java客户端中完全无效。
Java客户端相关代码片段:
public CustomNameResolverFactory(String host, int port) { ManagedChannel managedChannel = NettyChannelBuilder .forAddress(host, port) .withOption( ChannelOption.CONNECT_TIMEOUT_MILLIS, 10000 ) .usePlaintext().build(); GetServersRequest request = GetServersRequest.newBuilder().build(); GetServersResponse servers = TaskServiceGrpc.newBlockingStub(managedChannel).getServers(request); List<Server> serversList = servers.getServersList(); System.out.println(servers); LOGGER.info("found servers {}", servers); for (Server server : serversList) { String rpcAddr = server.getRpcAddr(); String[] split = rpcAddr.split(":"); String hostName = split[0]; int portN = Integer.parseInt(split[1]); addresses.add(new EquivalentAddressGroup(new InetSocketAddress(hostName, portN))); } }
问题分析
从Java客户端代码来看,核心问题是服务器地址仅在CustomNameResolverFactory初始化时拉取一次,后续不会主动刷新地址列表。而Golang客户端的Resolver实现内置了地址刷新逻辑,配合服务器的MaxConnectionAge触发重连时会重新获取地址,因此能检测到新增服务器。
Java gRPC的自定义NameResolver如果没有实现动态刷新机制,即使服务器主动断开连接触发重连,客户端也不会去拉取最新的服务器列表,自然无法感知集群变化。
解决方案
1. 给自定义NameResolver添加定时刷新逻辑
让NameResolver定期调用getServers接口获取最新服务器地址,并通知gRPC通道更新。示例核心代码:
public class CustomNameResolver extends NameResolver { private final ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor(); private Listener listener; private final String host; private final int port; public CustomNameResolver(String host, int port) { this.host = host; this.port = port; // 每分钟刷新一次地址 scheduler.scheduleAtFixedRate(this::refreshAddresses, 0, 1, TimeUnit.MINUTES); } private void refreshAddresses() { try (ManagedChannel channel = NettyChannelBuilder .forAddress(host, port) .withOption(ChannelOption.CONNECT_TIMEOUT_MILLIS, 10000) .usePlaintext() .build()) { GetServersRequest request = GetServersRequest.newBuilder().build(); GetServersResponse servers = TaskServiceGrpc.newBlockingStub(channel).getServers(request); List<EquivalentAddressGroup> newAddresses = new ArrayList<>(); for (Server server : servers.getServersList()) { String[] split = server.getRpcAddr().split(":"); newAddresses.add(new EquivalentAddressGroup( new InetSocketAddress(split[0], Integer.parseInt(split[1])))); } // 通知通道地址更新 if (listener != null) { listener.onAddresses(newAddresses, Attributes.EMPTY); } } catch (Exception e) { LOGGER.error("刷新服务器地址失败", e); } } @Override public String getServiceAuthority() { return host + ":" + port; } @Override public void start(Listener listener) { this.listener = listener; refreshAddresses(); // 初始化时立即拉取一次 } @Override public void shutdown() { scheduler.shutdown(); } }
2. 配置Java客户端的连接存活参数
确保客户端能响应服务器的MaxConnectionAge配置,连接到期后主动重连,从而使用刷新后的地址列表。在NettyChannelBuilder中添加:
.keepAliveTime(30, TimeUnit.SECONDS) .keepAliveTimeout(5, TimeUnit.SECONDS) .keepAliveWithoutCalls(true)
3. 让通道正确使用自定义Resolver
创建ManagedChannel时指定自定义的NameResolverFactory,确保通道能感知地址变化:
ManagedChannel channel = NettyChannelBuilder .forTarget("custom:///" + host + ":" + port) .nameResolverFactory(new NameResolverFactory() { @Override public NameResolver newNameResolver(URI uri, NameResolver.Args args) { String authority = uri.getAuthority(); String[] split = authority.split(":"); return new CustomNameResolver(split[0], Integer.parseInt(split[1])); } @Override public String getDefaultScheme() { return "custom"; } }) .usePlaintext() .build();
关键要点
- Java gRPC不会自动刷新服务器地址,必须通过自定义
NameResolver主动实现地址更新逻辑 - 服务器的
MaxConnectionAge只是触发重连的条件,客户端需要在重连时能获取到最新地址才能生效 - 定时刷新地址时要做好异常处理,避免单次拉取失败导致服务中断
内容的提问来源于stack exchange,提问作者M Kumar
相关产品推荐
相关产品推荐

