如何用Java(Apache HttpClient5)实时获取SSE响应单个Payload
实时处理SSE响应Payload的解决方案
使用Apache HttpClient 5实现实时处理
现有代码的核心问题是EntityUtils.toByteArray会一次性读取完整响应实体,无法做到实时接收Payload。要实现逐段处理,需要直接读取响应输入流,按SSE格式规则解析每一行内容:
实现逻辑
- 直接获取响应实体的输入流,避免一次性加载全部内容
- 逐行读取流数据,识别
data:前缀的行,累积多行data内容(SSE允许单个事件拆分到多个data行) - 遇到空行时,代表当前事件结束,触发Payload处理逻辑
- 处理流读取的异常与资源自动关闭
修改后的代码示例
// Set the socket timeout final ConnectionConfig connConfig = ConnectionConfig.custom() .setSocketTimeout(socketTimeout, TimeUnit.MILLISECONDS) .build(); // Custom config final BasicHttpClientConnectionManager cm = new BasicHttpClientConnectionManager(); cm.setConnectionConfig(connConfig); // Build the client try (final CloseableHttpClient client = HttpClientBuilder.create().setConnectionManager(cm).build()) { // Execute the request and process response stream return client.execute(request.getRequest(), response -> { // 校验响应状态码 if (response.getCode() != HttpStatus.SC_OK) { return HttpResponse.builder() .withCode(response.getCode()) .build(); } StringBuilder currentData = new StringBuilder(); try (InputStream inputStream = response.getEntity().getContent(); BufferedReader reader = new BufferedReader(new InputStreamReader(inputStream, StandardCharsets.UTF_8))) { String line; while ((line = reader.readLine()) != null) { // 处理SSE格式行 if (line.startsWith("data: ")) { // 移除前缀,累积多段data内容 String dataContent = line.substring(6); if (currentData.length() > 0) { currentData.append("\n"); } currentData.append(dataContent); } else if (line.isEmpty()) { // 空行代表当前事件结束,触发处理 if (currentData.length() > 0) { processPayload(currentData.toString()); currentData.setLength(0); // 重置缓冲区 } } // 可按需处理event:、id:等其他SSE字段 } // 处理未以空行结尾的最后一个事件 if (currentData.length() > 0) { processPayload(currentData.toString()); } } catch (IOException e) { throw new RuntimeException("读取SSE流失败", e); } return HttpResponse.builder() .withCode(response.getCode()) .build(); }); }
实时处理Payload的辅助方法
private void processPayload(String payload) { // 这里替换为你的业务逻辑,比如解析JSON、触发回调等 System.out.println("实时收到Payload: " + payload); }
替代方案:使用OkHttp处理SSE
如果觉得HttpClient的流处理不够简洁,可使用OkHttp,它对HTTP流的原生支持更友好:
OkHttp代码示例
OkHttpClient client = new OkHttpClient.Builder() .connectTimeout(socketTimeout, TimeUnit.MILLISECONDS) .readTimeout(socketTimeout, TimeUnit.MILLISECONDS) .build(); // 构造POST请求体(示例为JSON格式,可替换为FormBody等) RequestBody requestBody = RequestBody.create( MediaType.parse("application/json"), "{\"param\":\"value\"}" ); Request request = new Request.Builder() .url("你的服务器地址") .post(requestBody) .build(); client.newCall(request).enqueue(new Callback() { @Override public void onFailure(Call call, IOException e) { e.printStackTrace(); } @Override public void onResponse(Call call, Response response) throws IOException { if (!response.isSuccessful()) { throw new IOException("响应异常: " + response); } StringBuilder currentData = new StringBuilder(); BufferedReader reader = new BufferedReader( new InputStreamReader(response.body().byteStream(), StandardCharsets.UTF_8) ); String line; while ((line = reader.readLine()) != null) { if (line.startsWith("data: ")) { String dataContent = line.substring(6); if (currentData.length() > 0) { currentData.append("\n"); } currentData.append(dataContent); } else if (line.isEmpty()) { if (currentData.length() > 0) { processPayload(currentData.toString()); currentData.setLength(0); } } } // 处理最后一个事件 if (currentData.length() > 0) { processPayload(currentData.toString()); } response.close(); } });
内容的提问来源于stack exchange,提问作者tarka
相关产品推荐
相关产品推荐

