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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 10:13:15