diff --git a/spring-web-reactive/src/main/java/org/springframework/core/codec/StringDecoder.java b/spring-web-reactive/src/main/java/org/springframework/core/codec/StringDecoder.java index 51dc6a5eb5..a96e14f536 100644 --- a/spring-web-reactive/src/main/java/org/springframework/core/codec/StringDecoder.java +++ b/spring-web-reactive/src/main/java/org/springframework/core/codec/StringDecoder.java @@ -91,16 +91,16 @@ public class StringDecoder extends AbstractDecoder { if (this.splitOnNewline) { inputFlux = Flux.from(inputStream).flatMap(StringDecoder::splitOnNewline); } - return decodeInternal(inputFlux, mimeType); + return inputFlux.map(buffer -> decodeDataBuffer(buffer, mimeType)); } @Override public Mono decodeToMono(Publisher inputStream, ResolvableType elementType, MimeType mimeType, Object... hints) { - return decodeInternal(Flux.from(inputStream), mimeType). - collect(StringBuilder::new, StringBuilder::append). - map(StringBuilder::toString); + return Flux.from(inputStream) + .reduce(DataBuffer::write) + .map(buffer -> decodeDataBuffer(buffer, mimeType)); } private static Flux splitOnNewline(DataBuffer dataBuffer) { @@ -120,13 +120,11 @@ public class StringDecoder extends AbstractDecoder { return Flux.fromIterable(results); } - private Flux decodeInternal(Flux inputFlux, MimeType mimeType) { + private String decodeDataBuffer(DataBuffer dataBuffer, MimeType mimeType) { Charset charset = getCharset(mimeType); - return inputFlux.map(dataBuffer -> { - CharBuffer charBuffer = charset.decode(dataBuffer.asByteBuffer()); - DataBufferUtils.release(dataBuffer); - return charBuffer.toString(); - }); + CharBuffer charBuffer = charset.decode(dataBuffer.asByteBuffer()); + DataBufferUtils.release(dataBuffer); + return charBuffer.toString(); } private Charset getCharset(MimeType mimeType) { diff --git a/spring-web-reactive/src/test/java/org/springframework/core/codec/StringDecoderTests.java b/spring-web-reactive/src/test/java/org/springframework/core/codec/StringDecoderTests.java index f3d4bb43b0..43492ae6a9 100644 --- a/spring-web-reactive/src/test/java/org/springframework/core/codec/StringDecoderTests.java +++ b/spring-web-reactive/src/test/java/org/springframework/core/codec/StringDecoderTests.java @@ -73,7 +73,18 @@ public class StringDecoderTests extends AbstractDataBufferAllocatingTestCase { } @Test - public void decodeEmpty() throws InterruptedException { + public void decodeEmptyFlux() throws InterruptedException { + Flux source = Flux.empty(); + Flux output = this.decoder.decode(source, ResolvableType.forClass(String.class), null); + + TestSubscriber.subscribe(output) + .assertNoError() + .assertComplete() + .assertNoValues(); + } + + @Test + public void decodeEmptyString() throws InterruptedException { Flux source = Flux.just(stringBuffer("")); Flux output = this.decoder.decode(source, ResolvableType.forClass(String.class), null); @@ -92,4 +103,15 @@ public class StringDecoderTests extends AbstractDataBufferAllocatingTestCase { .assertValues("foobarbaz"); } + @Test + public void decodeToMonoWithEmptyFlux() throws InterruptedException { + Flux source = Flux.empty(); + Mono output = this.decoder.decodeToMono(source, ResolvableType.forClass(String.class), null); + + TestSubscriber.subscribe(output) + .assertNoError() + .assertComplete() + .assertNoValues(); + } + }