Spring WebFlux Client如何正确处理SSE中的[DONE]标记
解决Spring WebFlux Client处理OpenAI SSE中[DONE]的反序列化问题
OpenAI的SSE响应末尾会发送一行[DONE]标记流结束,但该字符串不是有效的JSON结构,直接用Jackson反序列化会抛出异常。以下是两种可行的解决方式:
方案1:直接使用WebClient手动处理响应流
绕过声明式接口的自动反序列化,先过滤无效行再手动解析JSON:
import com.fasterxml.jackson.databind.ObjectMapper; import org.springframework.http.MediaType; import org.springframework.web.reactive.function.client.WebClient; import reactor.core.publisher.Flux; public class OpenAIClient { private final WebClient webClient; private final ObjectMapper objectMapper; public OpenAIClient(WebClient.Builder webClientBuilder, ObjectMapper objectMapper) { this.webClient = webClientBuilder.baseUrl("https://api.openai.com/v1").build(); this.objectMapper = objectMapper; } public Flux<ChatCompletionResult> requestChatCompletion(AbstractChatCompletionRequest request) { return webClient.post() .uri("/chat/completions") .accept(MediaType.TEXT_EVENT_STREAM_VALUE) .bodyValue(request) .retrieve() .bodyToFlux(String.class) // 过滤掉[DONE]标记行 .filter(line -> !"[DONE]".equals(line.trim())) // 手动反序列化JSON到目标对象 .map(line -> { try { return objectMapper.readValue(line, ChatCompletionResult.class); } catch (Exception e) { // 可根据需求调整异常处理逻辑,比如跳过无效行 throw new RuntimeException("解析ChatCompletion chunk失败", e); } }); } }
方案2:为声明式客户端配置自定义解码器
如果想保留OpenAIAsyncClient声明式接口,可以通过自定义消息阅读器过滤无效行:
步骤1:实现自定义SSE消息阅读器
import com.fasterxml.jackson.databind.ObjectMapper; import org.springframework.core.ResolvableType; import org.springframework.http.MediaType; import org.springframework.http.codec.HttpMessageReader; import reactor.core.publisher.Flux; import java.util.List; public class OpenAISseMessageReader implements HttpMessageReader<Object> { private final ObjectMapper objectMapper; public OpenAISseMessageReader(ObjectMapper objectMapper) { this.objectMapper = objectMapper; } @Override public List<MediaType> getReadableMediaTypes() { return List.of(MediaType.TEXT_EVENT_STREAM); } @Override public boolean canRead(ResolvableType elementType, MediaType mediaType) { return mediaType.equals(MediaType.TEXT_EVENT_STREAM) && objectMapper.canDeserialize(elementType.getType()); } @Override public Flux<Object> read(ResolvableType elementType, org.springframework.http.ReactiveHttpInputMessage message, List<Object> hints) { return message.getBody() .map(buffer -> buffer.toString()) .filter(line -> !"[DONE]".equals(line.trim())) .map(line -> { try { return objectMapper.readValue(line, elementType.getType()); } catch (Exception e) { throw new RuntimeException("解析SSE chunk失败", e); } }); } }
步骤2:配置WebClient使用自定义解码器
import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.web.reactive.function.client.WebClient; import org.springframework.web.reactive.function.client.support.WebClientAdapter; import org.springframework.web.service.invoker.HttpServiceProxyFactory; @Configuration public class OpenAIClientConfig { @Bean public OpenAIAsyncClient openAIAsyncClient(WebClient.Builder webClientBuilder, ObjectMapper objectMapper) { WebClient webClient = webClientBuilder .baseUrl("https://api.openai.com/v1") .codecs(configurer -> { // 注册自定义SSE解码器 configurer.customCodecs().register(new OpenAISseMessageReader(objectMapper)); }) .build(); HttpServiceProxyFactory factory = HttpServiceProxyFactory.builderFor(WebClientAdapter.create(webClient)).build(); return factory.createClient(OpenAIAsyncClient.class); } }
注意事项
- 确保
ChatCompletionResult及其嵌套类(如Choice、Delta)的字段与OpenAI返回的JSON结构完全匹配。 - 可根据业务需求调整异常处理逻辑,例如跳过无效行而非直接抛出异常。
内容的提问来源于stack exchange,提问作者Oreste Viron
相关产品推荐
相关产品推荐

