ProtobufJsonEncoder actually supports streaming

Closes gh-37158
This commit is contained in:
rstoyanchev
2026-08-24 16:48:48 +01:00
parent 37c8f41633
commit 82cf15c60f
2 changed files with 56 additions and 17 deletions
@@ -54,6 +54,8 @@ import org.springframework.util.MimeType;
*/
public class ProtobufJsonEncoder implements HttpMessageEncoder<Message> {
private static final byte[] NEWLINE_SEPARATOR = {'\n'};
private static final byte[] EMPTY_BYTES = new byte[0];
private static final ResolvableType MESSAGE_TYPE = ResolvableType.forClass(Message.class);
@@ -117,22 +119,33 @@ public class ProtobufJsonEncoder implements HttpMessageEncoder<Message> {
.map(value -> encodeValue(value, bufferFactory, elementType, mimeType, hints))
.flux();
}
JsonArrayJoinHelper helper = new JsonArrayJoinHelper();
// Do not prepend JSON array prefix until first signal is known, onNext vs onError
// Keeps response not committed for error handling
return Flux.from(inputStream)
.map(value -> {
byte[] prefix = helper.getPrefix();
byte[] delimiter = helper.getDelimiter();
DataBuffer delimiterBuffer = bufferFactory.wrap(delimiter);
DataBuffer dataBuffer = encodeValue(value, bufferFactory, MESSAGE_TYPE, mimeType, hints);
return (prefix.length > 0 ?
bufferFactory.join(List.of(bufferFactory.wrap(prefix), delimiterBuffer, dataBuffer)) :
bufferFactory.join(List.of(delimiterBuffer, dataBuffer)));
})
.switchIfEmpty(Mono.fromCallable(() -> bufferFactory.wrap(helper.getPrefix())))
.concatWith(Mono.fromCallable(() -> bufferFactory.wrap(helper.getSuffix())));
byte[] separator = getStreamingMediaTypeSeparator(mimeType);
if (separator != null) {
return Flux.from(inputStream)
.map(value -> {
DataBuffer dataBuffer = encodeValue(value, bufferFactory, MESSAGE_TYPE, mimeType, hints);
return bufferFactory.join(List.of(dataBuffer, bufferFactory.wrap(separator)));
});
}
else {
JsonArrayJoinHelper helper = new JsonArrayJoinHelper();
// Do not prepend JSON array prefix until first signal is known, onNext vs onError
// Keeps response not committed for error handling
return Flux.from(inputStream)
.map(value -> {
byte[] prefix = helper.getPrefix();
byte[] delimiter = helper.getDelimiter();
DataBuffer delimiterBuffer = bufferFactory.wrap(delimiter);
DataBuffer dataBuffer = encodeValue(value, bufferFactory, MESSAGE_TYPE, mimeType, hints);
return (prefix.length > 0 ?
bufferFactory.join(List.of(bufferFactory.wrap(prefix), delimiterBuffer, dataBuffer)) :
bufferFactory.join(List.of(delimiterBuffer, dataBuffer)));
})
.switchIfEmpty(Mono.fromCallable(() -> bufferFactory.wrap(helper.getPrefix())))
.concatWith(Mono.fromCallable(() -> bufferFactory.wrap(helper.getSuffix())));
}
}
@Override
@@ -153,6 +166,22 @@ public class ProtobufJsonEncoder implements HttpMessageEncoder<Message> {
}
}
/**
* Return the separator to use for the given mime type.
* <p>By default, this method returns new line {@code "\n"} if the given
* mime type is one of the configured {@link #getStreamingMediaTypes
* streaming} mime types.
* @since 7.0.10
*/
protected byte @Nullable [] getStreamingMediaTypeSeparator(@Nullable MimeType mimeType) {
for (MediaType streamingMediaType : getStreamingMediaTypes()) {
if (streamingMediaType.isCompatibleWith(mimeType)) {
return NEWLINE_SEPARATOR;
}
}
return null;
}
private static class JsonArrayJoinHelper {
@@ -92,7 +92,7 @@ class ProtobufJsonEncoderTests extends AbstractEncoderTests<ProtobufJsonEncoder>
}
@Test
void encodeStream() {
void encodeNonStream() {
Flux<Message> input = Flux.just(this.msg1, this.msg2);
ResolvableType inputType = forClass(Msg.class);
@@ -104,7 +104,7 @@ class ProtobufJsonEncoderTests extends AbstractEncoderTests<ProtobufJsonEncoder>
}
@Test
void encodeEmptyFlux() {
void encodeNonStreamEmpty() {
Flux<Message> input = Flux.empty();
ResolvableType inputType = forClass(Msg.class);
Flux<DataBuffer> result = this.encoder.encode(input, this.bufferFactory, inputType,
@@ -115,6 +115,16 @@ class ProtobufJsonEncoderTests extends AbstractEncoderTests<ProtobufJsonEncoder>
.verifyComplete();
}
@Test
void encodeStream() {
Flux<Message> input = Flux.just(this.msg1, this.msg2);
ResolvableType inputType = forClass(Msg.class);
testEncode(input, inputType, MediaType.APPLICATION_NDJSON, null, step -> step
.assertNext(buffer -> assertBufferEqualsJson(buffer, "{\"foo\":\"Foo\",\"blah\":{\"blah\":123}}\n"))
.assertNext(buffer -> assertBufferEqualsJson(buffer, "{\"foo\":\"Bar\",\"blah\":{\"blah\":456}}\n"))
.verifyComplete());
}
private void assertBufferEqualsJson(DataBuffer actual, String expected) {
byte[] bytes = DataBufferTestUtils.dumpBytes(actual);