Java实现入站HTTP请求体到出站HTTP请求体的非阻塞转发
非阻塞HTTP POST请求体管道传输示例
以下示例基于Java 11+内置的HttpClient和Servlet 3.1+异步API实现,全程采用非阻塞流式处理,避免将请求体完整加载到内存,适合大文件或大流量场景。
import jakarta.servlet.AsyncContext; import jakarta.servlet.ServletException; import jakarta.servlet.annotation.WebServlet; import jakarta.servlet.http.HttpServlet; import jakarta.servlet.http.HttpServletRequest; import jakarta.servlet.http.HttpServletResponse; import java.io.IOException; import java.net.URI; import java.net.http.HttpClient; import java.net.http.HttpRequest; import java.net.http.HttpResponse; import java.nio.ByteBuffer; import java.util.concurrent.Flow; @WebServlet(urlPatterns = "/proxy", asyncSupported = true) public class ProxyServlet extends HttpServlet { private static final HttpClient HTTP_CLIENT = HttpClient.newHttpClient(); // 替换为实际的目标服务地址 private static final String TARGET_URL = "https://your-target-service.com/api/endpoint"; @Override protected void doPost(HttpServletRequest req, HttpServletResponse resp) throws ServletException, IOException { // 启动异步上下文,释放容器线程 AsyncContext asyncContext = req.startAsync(); // 将入站请求体包装为非阻塞的Flow.Publisher Flow.Publisher<ByteBuffer> requestBodyPublisher = subscriber -> { req.getInputStream().transferTo(new Flow.Subscriber<>() { @Override public void onSubscribe(Flow.Subscription subscription) { subscriber.onSubscribe(subscription); } @Override public void onNext(ByteBuffer item) { subscriber.onNext(item); } @Override public void onError(Throwable throwable) { subscriber.onError(throwable); asyncContext.complete(); } @Override public void onComplete() { subscriber.onComplete(); } }); }; // 构建出站POST请求,流式传输请求体 HttpRequest outboundRequest = HttpRequest.newBuilder() .uri(URI.create(TARGET_URL)) .POST(HttpRequest.BodyPublishers.fromPublisher(requestBodyPublisher)) // 透传Content-Type等必要头部 .headers("Content-Type", req.getHeader("Content-Type")) .build(); // 异步发送请求并处理响应 HTTP_CLIENT.sendAsync(outboundRequest, HttpResponse.BodyHandlers.ofInputStream()) .thenAccept(outboundResponse -> { try { // 透传响应状态码和头部 resp.setStatus(outboundResponse.statusCode()); outboundResponse.headers().map().forEach((key, values) -> { values.forEach(value -> resp.addHeader(key, value)); }); // 流式回写响应体到客户端 outboundResponse.body().transferTo(resp.getOutputStream()); } catch (IOException e) { resp.setStatus(HttpServletResponse.SC_INTERNAL_SERVER_ERROR); } finally { // 完成异步处理 asyncContext.complete(); } }) .exceptionally(e -> { resp.setStatus(HttpServletResponse.SC_INTERNAL_SERVER_ERROR); asyncContext.complete(); return null; }); } }
关键说明
- 异步Servlet:通过
asyncSupported = true标记Servlet支持异步处理,调用startAsync()释放容器线程,避免阻塞线程池。 - 非阻塞请求体传输:利用
Flow.Publisher将Servlet的输入流转为HttpClient可消费的流式数据源,数据边读边发,无需缓存整个请求体。 - 异步HTTP客户端:
HttpClient.sendAsync()以非阻塞方式发送请求,响应处理在回调线程中完成,不占用容器线程。 - 流式响应回写:将目标服务的响应输入流直接传输到客户端的输出流,同样采用非阻塞流式处理。
内容的提问来源于stack exchange,提问作者Mark_Sagecy
相关产品推荐
相关产品推荐

