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

使用Apache Http5客户端读取分块压缩HTTP响应Trailer时崩溃

解决Apache Http5客户端处理分块传输+Trailer时的压缩兼容性问题

问题背景

基于JAX-RS的Java REST服务通过分块传输编码发送大数据流,客户端使用Apache Http5。当服务端发送部分数据后抛出异常时,客户端无法正确检测问题,抛出MalformedChunkCodingException且无法获取Trailer字段。仅在禁用压缩(设置Accept-Encoding: identity)或不发送Trailer时,流程才能正常运行。

服务端核心代码(JAX-RS MessageBodyWriter)

@Override
public void writeTo(
    Stream<?> stream,
    Class<?> type,
    Type genericType,
    Annotation[] annotations,
    MediaType mediaType,
    MultivaluedMap<String, Object> httpHeaders,
    OutputStream entityStream)
    throws IOException {
     httpServletResponse.setHeader("Transfer-Encoding","chunked");
  final Map<String,String> trailers = new HashMap<>();
  httpServletResponse.setTrailerFields(()->trailers);
  try (stream;
       PrintWriter w = new PrintWriter(new OutputStreamWriter(entityStream,StandardCharsets.UTF_8))
  ) {   

    try {
        
          IntStream.rangeClosed(1, 100)
                    .forEachOrdered(
                        record -> {
                          try {
                             w.write(STRING_DATA);

                          } catch (Exception e) {
                            throw new ProcessingException(new IOException(e));
                          }
                        })
     
    } catch (Exception e) {
      trailers.put("X-Stream-Error","TT xxxxxxxxxxxxxxxxxxxxxx");

    }
  }
}

客户端核心代码(Apache Http5)

public static void main(String[] args) throws Exception{
  TrustStrategy acceptingTrustStrategy = (cert, authType) -> true;
  SSLContext sslContext = SSLContexts.custom().loadTrustMaterial(null, acceptingTrustStrategy).build();
  SSLConnectionSocketFactory sslsf = new SSLConnectionSocketFactory(sslContext,
          NoopHostnameVerifier.INSTANCE);
Registry<ConnectionSocketFactory> socketFactoryRegistry =
  RegistryBuilder.<ConnectionSocketFactory>create()
      .register("https", sslsf)
      .register("http", new PlainConnectionSocketFactory())
      .build();

  BasicHttpClientConnectionManager connectionManager =
          new BasicHttpClientConnectionManager(socketFactoryRegistry);

try (CloseableHttpClient client =
  HttpClients.custom() /*.setSSLSocketFactory(sslsf)*/
      .setConnectionManager(connectionManager)
     
      .build()) {

      HttpPost request = new HttpPost(URL);
      request.setHeader("Content-Type", "application/x-www-form-urlencoded");
      request.setHeader("TE","trailers");
      request.setHeader("Trailer","X-Stream-Error");


      List<NameValuePair> params = new ArrayList<>();

     

      request.setEntity(new UrlEncodedFormEntity(params));
      Supplier<List<? extends Header>> trailersSupplier = null;
      try{CloseableHttpResponse response = client.execute(request);
          HttpEntity entity = response.getEntity();
 
          int row=0;
          if (entity != null) {
              try
             
              {
                  InputStream inputStream = entity.getContent();

                  BufferedReader reader = new BufferedReader(new InputStreamReader(inputStream));
                  String line;
                  while ((line = reader.readLine()) != null) {
                      System.out.println("data => " + line);
                  }
                  trailersSupplier=entity.getTrailers();
                  List<? extends Header> trailers = trailersSupplier.get();
                  System.out.println("trailers=>"+trailers);
                  for (Header h :trailers){
                      System.out.println(h.getName()+":"+h.getValue());
                  }                     
                
                  inputStream.close();
                  entity.close();
                  response.close();

              }catch (Exception ex){
                  ex.printStackTrace();
              }
              System.out.println("done ="+row);
          }

          // Check for trailers

      } catch (IOException e) {
          e.printStackTrace();
      }
  } catch (IOException e) {
      e.printStackTrace();
  }
}

抛出异常

org.apache.hc.core5.http.MalformedChunkCodingException: CRLF expected at end of chunk
    at org.apache.hc.core5.http.impl.io.ChunkedInputStream.getChunkSize(ChunkedInputStream.java:250)
    at org.apache.hc.core5.http.impl.io.ChunkedInputStream.nextChunk(ChunkedInputStream.java:222)
    at org.apache.hc.core5.http.impl.io.ChunkedInputStream.read(ChunkedInputStream.java:147)
    at org.apache.hc.core5.http.impl.io.ChunkedInputStream.close(ChunkedInputStream.java:314)

问题根源

  1. 手动设置Transfer-Encoding干扰压缩逻辑:启用HTTP压缩时,Servlet容器会自动包装输出流做压缩处理并管理分块传输,手动设置Transfer-Encoding: chunked会破坏容器的压缩分块逻辑,导致分块格式错误。
  2. PrintWriter未及时flush:循环写入数据时,PrintWriter默认缓存数据,若未主动flush,异常发生时缓存数据可能未发送,导致分块结构不完整。
  3. 异常处理不当:抛出ProcessingException会中断容器正常响应流程,导致Trailer字段无法被正确附加到响应末尾。
  4. 客户端资源管理错误:手动关闭流的顺序可能导致Trailer未被正确读取,应使用自动资源管理机制。

修复方案

1. 服务端代码修正

@Override
public void writeTo(
    Stream<?> stream,
    Class<?> type,
    Type genericType,
    Annotation[] annotations,
    MediaType mediaType,
    MultivaluedMap<String, Object> httpHeaders,
    OutputStream entityStream)
    throws IOException {
    // 移除手动设置Transfer-Encoding,交由容器自动处理
    final Map<String, String> trailers = new HashMap<>();
    httpServletResponse.setTrailerFields(() -> trailers);
    
    try (stream;
         PrintWriter w = new PrintWriter(new BufferedWriter(new OutputStreamWriter(entityStream, StandardCharsets.UTF_8)))) {
        try {
            IntStream.rangeClosed(1, 100)
                    .forEachOrdered(record -> {
                        try {
                            w.write(STRING_DATA);
                            w.flush(); // 每次写入后强制flush,确保数据及时发送
                        } catch (Exception e) {
                            trailers.put("X-Stream-Error", "TT xxxxxxxxxxxxxxxxxxxxxx");
                            // 用RuntimeException终止循环,不中断容器流程
                            throw new RuntimeException(e);
                        }
                    });
        } catch (RuntimeException e) {
            // 补全异常处理,确保Trailer被正确设置
            if (!trailers.containsKey("X-Stream-Error")) {
                trailers.put("X-Stream-Error", "Unexpected error: " + e.getMessage());
            }
        }
        w.flush(); // 最终flush,确保所有数据发送完毕
    }
}

2. 客户端代码修正

public static void main(String[] args) throws Exception {
    TrustStrategy acceptingTrustStrategy = (cert, authType) -> true;
    SSLContext sslContext = SSLContexts.custom().loadTrustMaterial(null, acceptingTrustStrategy).build();
    SSLConnectionSocketFactory sslsf = new SSLConnectionSocketFactory(sslContext, NoopHostnameVerifier.INSTANCE);
    
    Registry<ConnectionSocketFactory> socketFactoryRegistry = RegistryBuilder.<ConnectionSocketFactory>create()
            .register("https", sslsf)
            .register("http", new PlainConnectionSocketFactory())
            .build();
    
    BasicHttpClientConnectionManager connectionManager = new BasicHttpClientConnectionManager(socketFactoryRegistry);
    
    try (CloseableHttpClient client = HttpClients.custom()
            .setConnectionManager(connectionManager)
            .build()) {
        
        HttpPost request = new HttpPost(URL);
        request.setHeader("Content-Type", "application/x-www-form-urlencoded");
        request.setHeader("TE", "trailers");
        request.setHeader("Trailer", "X-Stream-Error");
        // 如需保留压缩,无需设置Accept-Encoding: identity
        // request.setHeader("Accept-Encoding", "gzip, deflate");
        
        List<NameValuePair> params = new ArrayList<>();
        request.setEntity(new UrlEncodedFormEntity(params));
        
        // 使用try-with-resources自动管理响应资源
        try (CloseableHttpResponse response = client.execute(request)) {
            HttpEntity entity = response.getEntity();
            if (entity != null) {
                // 自动管理输入流
                try (InputStream inputStream = entity.getContent();
                     BufferedReader reader = new BufferedReader(new InputStreamReader(inputStream))) {
                    
                    String line;
                    while ((line = reader.readLine()) != null) {
                        System.out.println("data => " + line);
                    }
                    
                    // 读取完所有数据后再获取Trailer
                    Supplier<List<? extends Header>> trailersSupplier = entity.getTrailers();
                    if (trailersSupplier != null) {
                        List<? extends Header> trailers = trailersSupplier.get();
                        System.out.println("trailers=>" + trailers);
                        for (Header h : trailers) {
                            System.out.println(h.getName() + ":" + h.getValue());
                        }
                    }
                } catch (Exception ex) {
                    ex.printStackTrace();
                }
            }
        } catch (IOException e) {
            e.printStackTrace();
        }
    } catch (IOException e) {
        e.printStackTrace();
    }
}

关键说明

  • 移除手动分块设置:Servlet容器会根据响应大小和压缩配置自动处理分块传输,手动设置会导致压缩流与分块逻辑冲突。
  • 及时flush输出流:确保每次写入的数据都被发送到客户端,避免缓存导致的分块不完整。
  • 避免中断容器流程:异常时使用RuntimeException终止循环而非抛出ProcessingException,让容器能正常附加Trailer并结束响应。
  • 客户端自动资源管理:用try-with-resources确保流和响应被正确关闭,避免手动关闭顺序错误导致的Trailer读取失败。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 02:53:11