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

如何拦截Apache HttpClient5请求体?OpenSearchClient场景求助

解决OpenSearchClient基于Apache HttpClient5拦截搜索请求体的问题

问题分析

你在使用基于Apache HttpClient5的OpenSearchClient时,尝试通过HttpRequestInterceptor拦截请求记录查询语句,但发现拿到的EntityDetails是InternalAbstractHttpAsyncClient的内部类实例,无法直接提取请求体。这是因为异步客户端的请求实体是通过AsyncEntityProducer异步生成的,常规拦截器无法直接获取内容,需要用装饰器模式包装实体生产者来实现。

解决方案

1. 实现日志装饰类包装AsyncEntityProducer

创建一个装饰类,拦截AsyncEntityProducer的内容输出,记录请求体:

import org.apache.hc.core5.http.EntityDetails;
import org.apache.hc.core5.http.nio.AsyncEntityProducer;
import org.apache.hc.core5.http.nio.DataStreamChannel;
import org.apache.hc.core5.http.nio.SessionOutputBuffer;
import org.apache.hc.core5.util.Timeout;
import java.io.ByteArrayOutputStream;
import java.io.IOException;
import java.nio.ByteBuffer;
import java.nio.charset.StandardCharsets;

public class LoggingAsyncEntityProducer implements AsyncEntityProducer {

    private final AsyncEntityProducer delegate;
    private final ByteArrayOutputStream contentBuffer = new ByteArrayOutputStream();

    public LoggingAsyncEntityProducer(AsyncEntityProducer delegate) {
        this.delegate = delegate;
    }

    @Override
    public void produce(SessionOutputBuffer sessionOutputBuffer, DataStreamChannel dataStreamChannel) throws IOException {
        // 装饰输出缓冲区,复制请求体内容到本地buffer
        delegate.produce(new SessionOutputBuffer() {
            @Override
            public int write(ByteBuffer src) throws IOException {
                int written = sessionOutputBuffer.write(src);
                // 复制当前ByteBuffer中的内容到记录buffer
                if (src.position() > 0) {
                    src.flip();
                    byte[] bytes = new byte[src.remaining()];
                    src.get(bytes);
                    contentBuffer.write(bytes);
                    src.flip(); // 恢复缓冲区状态,不影响原输出
                }
                return written;
            }

            @Override
            public void flush() throws IOException {
                sessionOutputBuffer.flush();
            }

            @Override
            public int capacity() {
                return sessionOutputBuffer.capacity();
            }

            @Override
            public int length() {
                return sessionOutputBuffer.length();
            }

            @Override
            public boolean hasData() {
                return sessionOutputBuffer.hasData();
            }
        }, dataStreamChannel);

        // 实体生产完成后,输出记录的请求体
        if (isDone()) {
            String requestBody = contentBuffer.toString(StandardCharsets.UTF_8);
            log.info("OpenSearch 查询语句: {}", requestBody); // 替换为你的日志组件
        }
    }

    @Override
    public void releaseResources() {
        delegate.releaseResources();
    }

    @Override
    public boolean isRepeatable() {
        return delegate.isRepeatable();
    }

    @Override
    public boolean isDone() {
        return delegate.isDone();
    }

    @Override
    public void failed(Exception cause) {
        delegate.failed(cause);
    }

    @Override
    public Timeout getTimeout() {
        return delegate.getTimeout();
    }

    @Override
    public void setTimeout(Timeout timeout) {
        delegate.setTimeout(timeout);
    }

    @Override
    public EntityDetails getEntityDetails() {
        return delegate.getEntityDetails();
    }
}

2. 修改OpenSearchClient配置,添加拦截器替换实体生产者

在客户端配置中添加自定义拦截器,将原有的AsyncEntityProducer替换为装饰后的实例:

import org.apache.hc.client5.http.async.methods.SimpleHttpRequest;
import org.apache.hc.core5.http.HttpRequest;
import org.apache.hc.core5.http.HttpRequestInterceptor;
import org.apache.hc.core5.http.nio.AsyncEntityProducer;
import java.lang.reflect.Field;

public OpenSearchClient opensearchClient() {
    val builder = ApacheHttpClient5TransportBuilder.builder(new HttpHost(protocol, host, port));
    builder.setHttpClientConfigCallback(
        httpClientBuilder -> {
          val connectionManager =
              PoolingAsyncClientConnectionManagerBuilder.create()
                  .setDefaultConnectionConfig(
                      ConnectionConfig.custom()
                          .setConnectTimeout(timeout, TimeUnit.MILLISECONDS)
                          .setSocketTimeout(3000, TimeUnit.MILLISECONDS)
                          .build())
                  .build();

          // 自定义请求拦截器,替换实体生产者为日志装饰类
          HttpRequestInterceptor loggingInterceptor = (request, entityDetails, context) -> {
              if (request instanceof SimpleHttpRequest) {
                  // 处理SimpleHttpRequest类型的请求,直接替换实体生产者
                  SimpleHttpRequest simpleRequest = (SimpleHttpRequest) request;
                  AsyncEntityProducer originalProducer = simpleRequest.getEntityProducer();
                  if (originalProducer != null) {
                      simpleRequest.setEntityProducer(new LoggingAsyncEntityProducer(originalProducer));
                  }
              } else {
                  // 处理InternalAbstractHttpAsyncClient内部类请求,通过反射获取实体生产者
                  try {
                      Field entityProducerField = request.getClass().getDeclaredField("entityProducer");
                      entityProducerField.setAccessible(true);
                      AsyncEntityProducer originalProducer = (AsyncEntityProducer) entityProducerField.get(request);
                      if (originalProducer != null) {
                          entityProducerField.set(request, new LoggingAsyncEntityProducer(originalProducer));
                      }
                  } catch (NoSuchFieldException | IllegalAccessException e) {
                      log.error("获取请求实体生产者失败", e);
                  }
              }
          };

          return httpClientBuilder
              .setConnectionManager(connectionManager)
              .addRequestInterceptorFirst(loggingInterceptor);
        });

    return new OpenSearchClient(builder.build());
}

注意事项

  • 反射操作依赖HttpClient内部类结构,若后续版本更新内部字段名,需同步调整反射代码。
  • 若请求体过大,建议将ByteArrayOutputStream改为流式记录(比如直接写入日志文件),避免内存占用过高。
  • 确保日志内容符合隐私规范,不要记录敏感数据。

内容的提问来源于stack exchange,提问作者Shuvradeb Saha

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 13:25:02