Netty HttpClient如何在发送请求后接收FullHttpResponse
Netty HttpClient 聚合HttpContent为FullHttpResponse并在main函数输出的问题
我正在使用Netty编写HttpClient,现有代码中的HttpClientHandler会打印响应状态、请求头及内容。我需要将HttpContent聚合为FullHttpResponse,并在HttpHelloWorldClientApp的main函数中输出该响应,但目前cf.get()返回null,请问该如何实现?
原代码实现
HttpHelloWorldClientApp类
import io.netty.bootstrap.Bootstrap; import io.netty.buffer.Unpooled; import io.netty.channel.Channel; import io.netty.channel.ChannelFuture; import io.netty.channel.EventLoopGroup; import io.netty.channel.nio.NioEventLoopGroup; import io.netty.channel.socket.nio.NioSocketChannel; import io.netty.handler.codec.http.*; import io.netty.handler.ssl.SslContext; import io.netty.handler.ssl.SslContextBuilder; import io.netty.handler.ssl.util.InsecureTrustManagerFactory; import java.net.URI; public class HttpHelloWorldClientApp { static final String URL = System.getProperty("url", "https://www.google.com/"); public static void main(String[] args) throws Exception { URI uri = new URI(URL); String scheme = uri.getScheme() == null? "http" : uri.getScheme(); String host = uri.getHost() == null? "127.0.0.1" : uri.getHost(); int port = uri.getPort(); if (port == -1) { if ("http".equalsIgnoreCase(scheme)) { port = 80; } else if ("https".equalsIgnoreCase(scheme)) { port = 443; } } if (!"http".equalsIgnoreCase(scheme) && !"https".equalsIgnoreCase(scheme)) { System.err.println("Only HTTP(S) is supported."); return; } // Configure SSL context if necessary. final boolean ssl = "https".equalsIgnoreCase(scheme); final SslContext sslCtx; if (ssl) { sslCtx = SslContextBuilder.forClient() .trustManager(InsecureTrustManagerFactory.INSTANCE).build(); } else { sslCtx = null; } // Configure the client. EventLoopGroup group = new NioEventLoopGroup(); try { Bootstrap b = new Bootstrap(); b.group(group) .channel(NioSocketChannel.class) .handler(new HttpClientInitializer(sslCtx)); // Make the connection attempt. Channel ch = b.connect(host, port).sync().channel(); // Prepare the HTTP request. FullHttpRequest request = new DefaultFullHttpRequest( HttpVersion.HTTP_1_1, HttpMethod.GET, uri.getRawPath(), Unpooled.EMPTY_BUFFER); request.headers().set(HttpHeaderNames.HOST, host); request.headers().set(HttpHeaderNames.CONNECTION, HttpHeaderValues.CLOSE); request.headers().set(HttpHeaderNames.ACCEPT_ENCODING, HttpHeaderValues.GZIP); // Send the HTTP request. ChannelFuture cf = ch.writeAndFlush(request); // Wait for the server to close the connection. ch.closeFuture().sync(); System.out.println(cf.get()); } finally { // Shut down executor threads to exit. group.shutdownGracefully(); } } }
HttpClientInitializer类
import io.netty.channel.ChannelInitializer; import io.netty.channel.ChannelPipeline; import io.netty.channel.socket.SocketChannel; import io.netty.handler.codec.http.HttpClientCodec; import io.netty.handler.codec.http.HttpContentDecompressor; import io.netty.handler.codec.http.HttpObjectAggregator; import io.netty.handler.ssl.SslContext; public class HttpClientInitializer extends ChannelInitializer<SocketChannel> { private final SslContext sslCtx; public HttpClientInitializer(SslContext sslCtx) { this.sslCtx = sslCtx; } @Override public void initChannel(SocketChannel ch) { ChannelPipeline p = ch.pipeline(); // Enable HTTPS if necessary. if (sslCtx != null) { p.addLast(sslCtx.newHandler(ch.alloc())); } p.addLast(new HttpClientCodec()); // Remove the following line if you don't want automatic content decompression. p.addLast(new HttpContentDecompressor()); // Uncomment the following line if you don't want to handle HttpContents. p.addLast(new HttpObjectAggregator(1048576)); p.addLast(new HttpClientHandler()); } }
HttpClientHandler类
import io.netty.channel.ChannelHandlerContext; import io.netty.channel.SimpleChannelInboundHandler; import io.netty.handler.codec.http.*; import io.netty.util.CharsetUtil; public class HttpClientHandler extends SimpleChannelInboundHandler<HttpObject> { @Override public void channelRead0(ChannelHandlerContext ctx, HttpObject msg) { if (msg instanceof HttpResponse response) { System.err.println("STATUS: " + response.status()); System.err.println("VERSION: " + response.protocolVersion()); System.err.println(); if (!response.headers().isEmpty()) { for (CharSequence name: response.headers().names()) { for (CharSequence value: response.headers().getAll(name)) { System.err.println("HEADER: " + name + " = " + value); } } System.err.println(); } if (HttpUtil.isTransferEncodingChunked(response)) { System.err.println("CHUNKED CONTENT {"); } else { System.err.println("CONTENT {"); } } if (msg instanceof HttpContent content) { System.err.print(content.content().toString(CharsetUtil.UTF_8)); System.err.flush(); if (content instanceof LastHttpContent) { System.err.println("} END OF CONTENT"); ctx.close(); } } } @Override public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) { cause.printStackTrace(); ctx.close(); } }
解决方案
问题核心原因
cf.get()返回null是因为ChannelFuture仅代表请求发送操作完成的结果,和服务器返回的响应完全无关。要获取聚合后的FullHttpResponse,需要在Handler中把响应传递回主线程,不能直接通过writeAndFlush的Future获取。
另外你已经在Pipeline中添加了HttpObjectAggregator,它会自动将HttpResponse和后续的HttpContent片段聚合为完整的FullHttpResponse,因此Handler可以直接处理FullHttpResponse类型,无需再拆分处理不同的HttpObject。
具体修改步骤
1. 改造HttpClientHandler,通过Promise传递响应
import io.netty.channel.ChannelHandlerContext; import io.netty.channel.SimpleChannelInboundHandler; import io.netty.handler.codec.http.FullHttpResponse; import io.netty.util.concurrent.Promise; public class HttpClientHandler extends SimpleChannelInboundHandler<FullHttpResponse> { private final Promise<FullHttpResponse> responsePromise; public HttpClientHandler(Promise<FullHttpResponse> responsePromise) { this.responsePromise = responsePromise; } @Override protected void channelRead0(ChannelHandlerContext ctx, FullHttpResponse msg) { // 将聚合后的响应存入Promise,传递给主线程 responsePromise.setSuccess(msg.retain()); ctx.close(); } @Override public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) { // 异常时通知主线程 responsePromise.setFailure(cause); ctx.close(); } }
2. 修改HttpClientInitializer,支持传入Promise
import io.netty.channel.ChannelInitializer; import io.netty.channel.ChannelPipeline; import io.netty.channel.socket.SocketChannel; import io.netty.handler.codec.http.HttpClientCodec; import io.netty.handler.codec.http.HttpContentDecompressor; import io.netty.handler.codec.http.HttpObjectAggregator; import io.netty.handler.ssl.SslContext; import io.netty.util.concurrent.Promise; public class HttpClientInitializer extends ChannelInitializer<SocketChannel> { private final SslContext sslCtx; private final Promise<FullHttpResponse> responsePromise; public HttpClientInitializer(SslContext sslCtx, Promise<FullHttpResponse> responsePromise) { this.sslCtx = sslCtx; this.responsePromise = responsePromise; } @Override public void initChannel(SocketChannel ch) { ChannelPipeline p = ch.pipeline(); if (sslCtx != null) { p.addLast(sslCtx.newHandler(ch.alloc())); } p.addLast(new HttpClientCodec()); p.addLast(new HttpContentDecompressor()); p.addLast(new HttpObjectAggregator(1048576)); // 将Promise传入Handler p.addLast(new HttpClientHandler(responsePromise)); } }
3. 改造main方法,从Promise中获取并输出响应
import io.netty.bootstrap.Bootstrap; import io.netty.buffer.Unpooled; import io.netty.channel.Channel; import io.netty.channel.EventLoopGroup; import io.netty.channel.nio.NioEventLoopGroup; import io.netty.channel.socket.nio.NioSocketChannel; import io.netty.handler.codec.http.*; import io.netty.handler.ssl.SslContext; import io.netty.handler.ssl.SslContextBuilder; import io.netty.handler.ssl.util.InsecureTrustManagerFactory; import io.netty.util.concurrent.DefaultPromise; import io.netty.util.concurrent.Promise; import java.net.URI; public class HttpHelloWorldClientApp { static final String URL = System.getProperty("url", "https://www.google.com/"); public static void main(String[] args) throws Exception { URI uri = new URI(URL); String scheme = uri.getScheme() == null ? "http" : uri.getScheme(); String host = uri.getHost() == null ? "127.0.0.1" : uri.getHost(); int port = uri.getPort(); if (port == -1) { if ("http".equalsIgnoreCase(scheme)) { port = 80; } else if ("https".equalsIgnoreCase(scheme)) { port = 443; } } if (!"http".equalsIgnoreCase(scheme) && !"https".equalsIgnoreCase(scheme)) { System.err.println("Only HTTP(S) is supported."); return; } final boolean ssl = "https".equalsIgnoreCase(scheme); final SslContext sslCtx; if (ssl) { sslCtx = SslContextBuilder.forClient() .trustManager(InsecureTrustManagerFactory.INSTANCE).build(); } else { sslCtx = null; } EventLoopGroup group = new NioEventLoopGroup(); try { Bootstrap b = new Bootstrap(); // 创建Promise用于接收异步响应 Promise<FullHttpResponse> responsePromise = new DefaultPromise<>(group.next()); b.group(group) .channel(NioSocketChannel.class) .handler(new HttpClientInitializer(sslCtx, responsePromise)); Channel ch = b.connect(host, port).sync().channel(); FullHttpRequest request = new DefaultFullHttpRequest( HttpVersion.HTTP_1_1, HttpMethod.GET, uri.getRawPath(), Unpooled.EMPTY_BUFFER); request.headers().set(HttpHeaderNames.HOST, host); request.headers().set(HttpHeaderNames.CONNECTION, HttpHeaderValues.CLOSE); request.headers().set(HttpHeaderNames.ACCEPT_ENCODING, HttpHeaderValues.GZIP); ch.writeAndFlush(request); // 等待响应返回,同步获取结果 FullHttpResponse response = responsePromise.sync().get(); try { // 在main函数中输出完整响应 System.out.println("STATUS: " + response.status()); System.out.println("VERSION: " + response.protocolVersion()); System.out.println(); if (!response.headers().isEmpty()) { for (CharSequence name : response.headers().names()) { for (CharSequence value : response.headers().getAll(name)) { System.out.println("HEADER: " + name + " = " + value); } } System.out.println(); } System.out.println("CONTENT:"); System.out.println(response.content().toString(io.netty.util.CharsetUtil.UTF_8)); } finally { // 释放响应内存,避免内存泄漏 response.release(); } ch.closeFuture().sync(); } finally { group.shutdownGracefully(); } } }
关键说明
HttpObjectAggregator的作用是自动聚合分块响应,所以Handler可以直接处理FullHttpResponse,无需手动拼接内容。- 使用Netty的
Promise实现异步结果传递,主线程通过responsePromise.sync()等待响应完成,这是Netty中跨线程传递结果的标准方式。 - 必须调用
response.release()释放FullHttpResponse的内存,因为Netty的ByteBuf是池化的,不释放会导致内存泄漏。
内容的提问来源于stack exchange,提问作者user51
相关产品推荐
相关产品推荐

