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

如何将Netty HTTP响应从Channel Handler返回至execute方法

实现Netty HTTP客户端用CompletableFuture返回请求结果

核心思路是将请求与对应的CompletableFuture绑定,在HttpResponseHandler中拿到关联的Future并完成它,同时处理异常场景避免Future一直处于未完成状态。以下分两种场景给出具体实现:

场景1:短连接(一次请求对应一个Channel)

这种场景下每个请求会创建新Channel,直接用Channel的属性存储Future即可。

1. 定义全局AttributeKey

在工具类或Handler中定义存储Future的键:

private static final AttributeKey<CompletableFuture<FullHttpResponse>> RESPONSE_FUTURE = 
    AttributeKey.valueOf("RESPONSE_FUTURE");

2. 修改NettyHttpClient的executeRequest方法

创建CompletableFuture并绑定到Channel属性,发送请求后返回Future:

public CompletableFuture<FullHttpResponse> executeRequest(String host, int port, HttpRequest request) {
    CompletableFuture<FullHttpResponse> future = new CompletableFuture<>();
    
    // 初始化Bootstrap并建立连接
    Bootstrap bootstrap = new Bootstrap()
        .group(new NioEventLoopGroup())
        .channel(NioSocketChannel.class)
        .handler(new HttpChannelInitializer());

    try {
        Channel channel = bootstrap.connect(host, port).sync().channel();
        // 将Future绑定到Channel属性
        channel.attr(RESPONSE_FUTURE).set(future);
        // 发送HTTP请求
        channel.writeAndFlush(request);
    } catch (InterruptedException e) {
        future.completeExceptionally(e);
        Thread.currentThread().interrupt();
    }

    // 可选:添加超时处理,避免Future永久挂起
    future.orTimeout(10, TimeUnit.SECONDS);
    return future;
}

3. 修改HttpResponseHandler处理响应与异常

在channelRead0中完成Future,同时在异常/通道关闭时标记Future异常:

public class HttpResponseHandler extends SimpleChannelInboundHandler<FullHttpResponse> {

    @Override
    protected void channelRead0(ChannelHandlerContext ctx, FullHttpResponse response) throws Exception {
        // 从Channel属性中获取关联的Future
        CompletableFuture<FullHttpResponse> future = ctx.channel().attr(RESPONSE_FUTURE).get();
        if (future != null && !future.isDone()) {
            // 调用retain防止Netty自动释放响应的ByteBuf
            future.complete(response.retain());
            // 完成后清空属性,避免内存泄漏
            ctx.channel().attr(RESPONSE_FUTURE).set(null);
        }

        // 短连接场景下,根据响应头关闭通道
        if (!HttpUtil.isKeepAlive(response)) {
            ctx.channel().close();
        }
    }

    @Override
    public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {
        // 异常时标记Future为失败
        CompletableFuture<FullHttpResponse> future = ctx.channel().attr(RESPONSE_FUTURE).get();
        if (future != null && !future.isDone()) {
            future.completeExceptionally(cause);
            ctx.channel().attr(RESPONSE_FUTURE).set(null);
        }
        ctx.close();
    }

    @Override
    public void channelInactive(ChannelHandlerContext ctx) throws Exception {
        // 通道关闭但未收到响应时,标记Future失败
        CompletableFuture<FullHttpResponse> future = ctx.channel().attr(RESPONSE_FUTURE).get();
        if (future != null && !future.isDone()) {
            future.completeExceptionally(new IOException("Channel closed before receiving response"));
            ctx.channel().attr(RESPONSE_FUTURE).set(null);
        }
        super.channelInactive(ctx);
    }
}

场景2:长连接(多个请求复用同一个Channel)

长连接下需要用请求唯一ID映射请求与Future,避免Channel属性被覆盖。

1. 在NettyHttpClient中维护请求-Future映射

添加全局映射表和ID生成器:

public class NettyHttpClient {
    private final ConcurrentHashMap<String, CompletableFuture<FullHttpResponse>> requestFutureMap = new ConcurrentHashMap<>();
    private final AtomicInteger requestIdGenerator = new AtomicInteger(0);
    private Channel persistentChannel;

    // 初始化长连接逻辑...
}

2. 修改executeRequest方法

生成唯一请求ID,添加到请求头,并存入映射表:

public CompletableFuture<FullHttpResponse> executeRequest(HttpRequest request) {
    String requestId = String.valueOf(requestIdGenerator.incrementAndGet());
    // 自定义Header存储请求ID
    request.headers().set("X-Request-ID", requestId);

    CompletableFuture<FullHttpResponse> future = new CompletableFuture<>();
    requestFutureMap.put(requestId, future);

    // 复用长连接发送请求
    persistentChannel.writeAndFlush(request);

    // 超时处理+完成后自动清理映射
    future.orTimeout(10, TimeUnit.SECONDS)
          .whenComplete((resp, err) -> requestFutureMap.remove(requestId));
    return future;
}

3. 修改HttpResponseHandler匹配请求ID

从响应头中获取请求ID,找到对应Future并完成:

public class HttpResponseHandler extends SimpleChannelInboundHandler<FullHttpResponse> {
    private final NettyHttpClient client;

    // 通过构造器注入HttpClient实例
    public HttpResponseHandler(NettyHttpClient client) {
        this.client = client;
    }

    @Override
    protected void channelRead0(ChannelHandlerContext ctx, FullHttpResponse response) throws Exception {
        String requestId = response.headers().get("X-Request-ID");
        if (requestId != null) {
            // 从映射表中移除并获取对应的Future
            CompletableFuture<FullHttpResponse> future = client.requestFutureMap.remove(requestId);
            if (future != null && !future.isDone()) {
                future.complete(response.retain());
            }
        }

        // 长连接保持通道打开,除非响应头明确要求关闭
        if (!HttpUtil.isKeepAlive(response)) {
            ctx.channel().close();
        }
    }

    @Override
    public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {
        // 异常时将所有未完成的Future标记为失败
        client.requestFutureMap.values().forEach(future -> {
            if (!future.isDone()) {
                future.completeExceptionally(cause);
            }
        });
        client.requestFutureMap.clear();
        ctx.close();
    }
}

关键注意事项

  • 内存泄漏预防:Future完成后必须从映射表/Channel属性中移除,避免无用对象堆积。
  • ByteBuf引用管理:调用future.complete(response.retain()),防止Netty自动释放响应的ByteBuf导致后续读取失败。
  • 超时处理:必须给Future添加超时逻辑,避免请求无响应时Future永久阻塞。

内容的提问来源于stack exchange,提问作者user51

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 09:48:11