如何将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
相关产品推荐
相关产品推荐

