如何从HttpAsyncRequestProducer读取请求体并支持多次读取?
解决方案
要解决从HttpAsyncRequestProducer读取请求体并支持多次读取的问题,核心思路是包装原请求生产者,缓存请求体内容并替换为可重复读取的HttpEntity,同时保证后续业务流程不受影响。以下是具体实现步骤:
1. 自定义可缓存的HttpAsyncRequestProducer包装类
创建包装类代理原HttpAsyncRequestProducer的所有方法,缓存请求体内容,并将原请求的HttpEntity替换为可重复读取的版本:
import org.apache.http.HttpEntity; import org.apache.http.HttpRequest; import org.apache.http.HttpEntityEnclosingRequest; import org.apache.http.nio.protocol.HttpAsyncRequestProducer; import org.apache.http.nio.ContentEncoder; import org.apache.http.nio.IOControl; import org.apache.http.protocol.HttpContext; import java.io.ByteArrayOutputStream; import java.io.IOException; import java.io.InputStream; public class CachedRequestProducer implements HttpAsyncRequestProducer { private final HttpAsyncRequestProducer delegate; private byte[] cachedRequestBody; private HttpRequest cachedRequest; public CachedRequestProducer(HttpAsyncRequestProducer delegate) { this.delegate = delegate; } @Override public HttpRequest generateRequest() throws IOException { if (cachedRequest == null) { HttpRequest request = delegate.generateRequest(); // 处理带实体的请求 if (request instanceof HttpEntityEnclosingRequest) { HttpEntityEnclosingRequest enclosingRequest = (HttpEntityEnclosingRequest) request; HttpEntity entity = enclosingRequest.getEntity(); if (entity != null) { // 读取实体内容并缓存 cachedRequestBody = readEntityContent(entity); // 替换为可重复读取的ByteArrayEntity enclosingRequest.setEntity(new org.apache.http.entity.ByteArrayEntity(cachedRequestBody)); } } cachedRequest = request; } return cachedRequest; } // 读取HttpEntity内容到byte数组 private byte[] readEntityContent(HttpEntity entity) throws IOException { try (InputStream is = entity.getContent(); ByteArrayOutputStream baos = new ByteArrayOutputStream()) { byte[] buffer = new byte[1024]; int len; while ((len = is.read(buffer)) != -1) { baos.write(buffer, 0, len); } return baos.toByteArray(); } } // 以下方法直接代理原producer的实现 @Override public void produceContent(ContentEncoder encoder, IOControl ioctrl) throws IOException { delegate.produceContent(encoder, ioctrl); } @Override public void requestCompleted(HttpContext context) { delegate.requestCompleted(context); } @Override public void failed(Exception ex) { delegate.failed(ex); } @Override public boolean isRepeatable() { return delegate.isRepeatable(); } @Override public void resetRequest() throws IOException { delegate.resetRequest(); // 重置缓存,适配原producer的重置逻辑 cachedRequest = null; cachedRequestBody = null; } @Override public void close() throws IOException { delegate.close(); } // 获取缓存的请求体 public byte[] getCachedRequestBody() { return cachedRequestBody; } }
2. 修改ByteBuddy拦截器代码
在拦截器中替换原HttpAsyncRequestProducer为自定义包装类,即可安全读取请求体,同时保证后续业务能重复读取:
import net.bytebuddy.implementation.bind.annotation.*; import org.apache.http.HttpRequest; import org.apache.http.HttpResponse; import org.apache.http.nio.protocol.HttpAsyncRequestProducer; import java.util.Arrays; import java.util.concurrent.Callable; import java.util.concurrent.Future; @RuntimeType public static Future<Object> doProceed(@Origin Method method, @SuperCall Callable<Future<Object>> callable, @AllArguments Object[] args) throws Exception { System.out.println("Inside ElasticSearchInterceptor"); System.out.println("ElasticSearchInterceptor: Method-> " + method); int i = 0; for (Object obj : args) { System.out.println("Argument[" + i + "]: " + obj); i++; } // 包装原requestProducer,替换入参 HttpAsyncRequestProducer originalProducer = (HttpAsyncRequestProducer) args[0]; CachedRequestProducer cachedProducer = new CachedRequestProducer(originalProducer); args[0] = cachedProducer; // 生成请求并缓存请求体 HttpRequest httpRequest = cachedProducer.generateRequest(); System.out.println("ElasticSearchInterceptor: Request Headers-> " + Arrays.toString(httpRequest.getAllHeaders())); // 读取缓存的请求体 byte[] requestBody = cachedProducer.getCachedRequestBody(); if (requestBody != null) { String bodyStr = new String(requestBody, "UTF-8"); // 根据实际编码调整 System.out.println("ElasticSearchInterceptor: Request Body-> " + bodyStr); } else { System.out.println("ElasticSearchInterceptor: No request body"); } System.out.println("ElasticSearchInterceptor: Before call"); Future<Object> call = callable.call(); HttpResponse res = (HttpResponse) call.get(); System.out.println("Response in elastic search " + res); System.out.println("ElasticSearchInterceptor: After call"); return call; }
关键说明
- 请求体缓存与复用:在
generateRequest()中读取原实体内容并缓存,同时替换为ByteArrayEntity(天然支持重复读取),确保后续流程多次读取请求体时不会出错。 - 无侵入代理:包装类完全代理原生产者的所有方法,除了请求体处理逻辑外,其他行为与原生产者一致,不会影响ElasticSearch客户端的正常运行。
- 编码适配:读取请求体时的UTF-8编码请根据实际业务场景调整,避免出现乱码问题。
内容的提问来源于stack exchange,提问作者Gourav Kumar
相关产品推荐
相关产品推荐

