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

如何从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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 22:20:32