使用WebFlux消费Server Sent Events无有效响应问题排查
问题描述
我在Java 11中尝试用WebFlux消费Server Sent Events(SSE),以下是最小可运行代码:
public static void main(String[] args) { consumeServerSentEvent(); } public static void consumeServerSentEvent() { WebClient client = WebClient.create("https://chat.poemhub.top"); ParameterizedTypeReference<ServerSentEvent<String>> type = new ParameterizedTypeReference<>() {}; Flux<ServerSentEvent<String>> eventStream = client.get() .uri("/v1/completions?q=hello") .headers(httpHeaders -> { httpHeaders.set("Authorization", "Bearer sk-DRjSqzDl3xHL1234663AfET3BlbkFJfQXGId4540GNmZqr7ss"); httpHeaders.set("Content-Type",MediaType.TEXT_EVENT_STREAM.toString()); }) .retrieve() .bodyToFlux(type); eventStream.subscribe( content -> { log.info("Time: {} - event: name[{}], id [{}], content[{}] ", LocalTime.now(), content.event(), content.id(), content.data()); }, error -> { log.error("Error receiving SSE: {}", error); }, () -> { log.info("Completed!!!"); } ); }
用curl测试该API可正常返回数据:
➜ ~ curl -N -H "Authorization: Bearer sk-DRjSqzDlwieHL12UvxAfET3BlbkFJfQioueNmZqr7ss" https://chat.poemhub.top/v1/completions\?q\=hello [{"id":"cmpl-6xEWYsb3cce1Xni4Xsqv8OdkU2EpQ","object":"text_completion","created":1679575202,"choices":[{"text":"\n","index":0,"logprobs":null,"finish_reason":null}],"model":"text-davinci-003"},{"id":"cmpl-6xEWYsb3cce1Xni4Xsqv8OdkU2EpQ","object":"text_completion","created":1679575202,"choices":[{"text":"\n","index":0,"logprobs":null,"finish_reason":null}],"model":"text-davinci-003"},{"id":"cmpl-6xEWYsb3cce1Xni4Xsqv8OdkU2EpQ","object":"text_completion","created":1679575202,"choices":[{"text":"Hello","index":0,"logprobs":null,"finish_reason":null}],"model":"text-davinci-003"},{"id":"cmpl-6xEWYsb3cce1Xni4Xsqv8OdkU2EpQ","object":"text_completion","created":1679575202,"choices":[{"text":" there","index":0,"logprobs":null,"finish_reason":null}],"model":"text-davinci-003"},{"id":"cmpl-6xEWYsb3cce1Xni4Xsqv8OdkU2EpQ","object":"text_completion","created":1679575202,"choices":[{"text":"!","index":0,"logprobs":null,"finish_reason":null}],"model":"text-davinci-003"},{"id":"cmpl-6xEWYsb3cce1Xni4Xsqv8OdkU2EpQ","object":"text_completion","created":1679575202,"choices":[{"text":" How","index":0,"logprobs":null,"finish_reason":null}],"model":"text-davinci-003"},{"id":"cmpl-6xEWYsb3cce1Xni4Xsqv8OdkU2EpQ","object":"text_completion","created":1679575202,"choices":[{"text":" can","index":0,"logprobs":null,"finish_reason":null}],"model":"text-davinci-003"},{"id":"cmpl-6xEWYsb3cce1Xni4Xsqv8OdkU2EpQ","object":"text_completion","created":1679575202,"choices":[{"text":" I","index":0,"logprobs":null,"finish_reason":null}],"model":"text-davinci-003"},{"id":"cmpl-6xEWYsb3cce1Xni4Xsqv8OdkU2EpQ","object":"text_completion","created":1679575202,"choices":[{"text":" help","index":0,"logprobs":null,"finish_reason":null}],"model":"text-davinci-003"},{"id":"cmpl-6xEWYsb3cce1Xni4Xsqv8OdkU2EpQ","object":"text_completion","created":1679575202,"choices":[{"text":" you","index":0,"logprobs":null,"finish_reason":null}],"model":"text-davinci-003"},{"id":"cmpl-6xEWYsb3cce1Xni4Xsqv8OdkU2EpQ","object":"text_completion","created":1679575202,"choices":[{"text":"?","index":0,"logprobs":null,"finish_reason":null}],"model":"text-davinci-003"}]%
但运行WebFlux代码时无任何日志输出,请问遗漏了什么配置或逻辑?该如何解决?
解决方案
1. 解决JVM提前退出问题
WebFlux基于Reactor实现异步非阻塞逻辑,subscribe()方法是异步执行的,main方法调用后会直接结束,导致JVM终止,订阅的事件还没处理就退出了。需要通过阻塞等待来让JVM保持运行,直到流处理完成。
方法一:使用CountDownLatch
import java.util.concurrent.CountDownLatch; import java.time.LocalTime; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.web.reactive.function.client.WebClient; import reactor.core.publisher.Flux; public class SseConsumer { private static final Logger log = LoggerFactory.getLogger(SseConsumer.class); public static void main(String[] args) throws InterruptedException { CountDownLatch latch = new CountDownLatch(1); consumeServerSentEvent(latch); latch.await(); // 等待流处理完成 } public static void consumeServerSentEvent(CountDownLatch latch) { WebClient client = WebClient.create("https://chat.poemhub.top"); Flux<String> eventStream = client.get() .uri("/v1/completions?q=hello") .headers(httpHeaders -> { httpHeaders.set("Authorization", "Bearer sk-DRjSqzDl3xHL1234663AfET3BlbkFJfQXGId4540GNmZqr7ss"); httpHeaders.set("Accept", "application/x-ndjson"); }) .retrieve() .bodyToFlux(String.class); eventStream.subscribe( content -> log.info("Time: {} - content[{}] ", LocalTime.now(), content), error -> { log.error("Error receiving SSE: {}", error); latch.countDown(); }, () -> { log.info("Completed!!!"); latch.countDown(); // 流结束时释放 latch } ); } }
方法二:使用blockLast()(适合测试场景)
import java.time.LocalTime; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.web.reactive.function.client.WebClient; import reactor.core.publisher.Flux; public class SseConsumer { private static final Logger log = LoggerFactory.getLogger(SseConsumer.class); public static void main(String[] args) { consumeServerSentEvent() .doOnNext(content -> log.info("Time: {} - content[{}] ", LocalTime.now(), content)) .doOnError(error -> log.error("Error receiving SSE: {}", error)) .doOnComplete(() -> log.info("Completed!!!")) .blockLast(); // 阻塞直到流结束 } public static Flux<String> consumeServerSentEvent() { WebClient client = WebClient.create("https://chat.poemhub.top"); return client.get() .uri("/v1/completions?q=hello") .headers(httpHeaders -> { httpHeaders.set("Authorization", "Bearer sk-DRjSqzDl3xHL1234663AfET3BlbkFJfQXGId4540GNmZqr7ss"); httpHeaders.set("Accept", "application/x-ndjson"); }) .retrieve() .bodyToFlux(String.class); } }
2. 解决SSE解析不匹配问题
从curl输出可以看出,该API返回的是换行分隔JSON(NDJSON)流,而非标准SSE格式(标准SSE需包含data: 前缀和换行分隔符)。用ServerSentEvent<String>解析会因格式不匹配导致无法识别事件,需修改解析逻辑:
步骤1:定义对应POJO类(可选,如需结构化解析)
import java.util.List; public class CompletionResponse { private String id; private String object; private long created; private List<Choice> choices; private String model; // Getters and Setters public static class Choice { private String text; private int index; private Object logprobs; private String finish_reason; // Getters and Setters } }
步骤2:修改为结构化解析
将bodyToFlux(String.class)替换为bodyToFlux(CompletionResponse.class),即可直接得到结构化对象:
public static Flux<CompletionResponse> consumeServerSentEvent() { WebClient client = WebClient.create("https://chat.poemhub.top"); return client.get() .uri("/v1/completions?q=hello") .headers(httpHeaders -> { httpHeaders.set("Authorization", "Bearer sk-DRjSqzDl3xHL1234663AfET3BlbkFJfQXGId4540GNmZqr7ss"); httpHeaders.set("Accept", "application/x-ndjson"); }) .retrieve() .bodyToFlux(CompletionResponse.class); }
步骤3:调整日志输出逻辑
.doOnNext(response -> { log.info("Time: {} - id[{}], text[{}] ", LocalTime.now(), response.getId(), response.getChoices().get(0).getText()); })
额外说明
- 生产环境中应避免使用
block()或blockLast(),这类方法会阻塞线程,违背Reactor异步非阻塞设计,建议整合到Spring应用上下文,由框架管理生命周期。 - 若API后续改为标准SSE格式,可换回
ServerSentEvent解析,并确保响应的Content-Type为text/event-stream。
内容的提问来源于stack exchange,提问作者Dolphin
相关产品推荐
相关产品推荐

