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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 22:45:45