You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.12 07:01:10