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

