Batch SSE events writes when possible
Prior to this commit, the `SseEventBuilder` would be used to create SSE events and write them to the connection using the `ResponseBodyEmitter`. This would send each data item one by one, effectively writing and flushing to the network for each. Since multiple data lines are prepared by the `SseEventBuilder`, a typical write of an SSE event performs multiple flushes operations. This commit adds a method on `ResponseBodyEmitter` to perform batch writes (given a `Set<DataWithMediaType>`) and only flush once all elements of the set have been written. This also applies in case of early writes, where now all buffered elements are written then flushed altogether. Fixes gh-30912
This commit is contained in:
@@ -128,9 +128,7 @@ public class ResponseBodyEmitter {
|
||||
this.handler = handler;
|
||||
|
||||
try {
|
||||
for (DataWithMediaType sendAttempt : this.earlySendAttempts) {
|
||||
sendInternal(sendAttempt.getData(), sendAttempt.getMediaType());
|
||||
}
|
||||
sendInternal(this.earlySendAttempts);
|
||||
}
|
||||
finally {
|
||||
this.earlySendAttempts.clear();
|
||||
@@ -194,11 +192,7 @@ public class ResponseBodyEmitter {
|
||||
*/
|
||||
public synchronized void send(Object object, @Nullable MediaType mediaType) throws IOException {
|
||||
Assert.state(!this.complete, () -> "ResponseBodyEmitter has already completed" +
|
||||
(this.failure != null ? " with error: " + this.failure : ""));
|
||||
sendInternal(object, mediaType);
|
||||
}
|
||||
|
||||
private void sendInternal(Object object, @Nullable MediaType mediaType) throws IOException {
|
||||
(this.failure != null ? " with error: " + this.failure : ""));
|
||||
if (this.handler != null) {
|
||||
try {
|
||||
this.handler.send(object, mediaType);
|
||||
@@ -217,6 +211,43 @@ public class ResponseBodyEmitter {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Write a set of data and MediaType pairs in a batch.
|
||||
* <p>Compared to {@link #send(Object, MediaType)}, this batches the write operations
|
||||
* and flushes to the network at the end.
|
||||
* @param items the object and media type pairs to write
|
||||
* @throws IOException raised when an I/O error occurs
|
||||
* @throws java.lang.IllegalStateException wraps any other errors
|
||||
* @since 6.0.12
|
||||
*/
|
||||
public synchronized void send(Set<DataWithMediaType> items) throws IOException {
|
||||
Assert.state(!this.complete, () -> "ResponseBodyEmitter has already completed" +
|
||||
(this.failure != null ? " with error: " + this.failure : ""));
|
||||
sendInternal(items);
|
||||
}
|
||||
|
||||
private void sendInternal(Set<DataWithMediaType> items) throws IOException {
|
||||
if (items.isEmpty()) {
|
||||
return;
|
||||
}
|
||||
if (this.handler != null) {
|
||||
try {
|
||||
this.handler.send(items);
|
||||
}
|
||||
catch (IOException ex) {
|
||||
this.sendFailed = true;
|
||||
throw ex;
|
||||
}
|
||||
catch (Throwable ex) {
|
||||
this.sendFailed = true;
|
||||
throw new IllegalStateException("Failed to send " + items, ex);
|
||||
}
|
||||
}
|
||||
else {
|
||||
this.earlySendAttempts.addAll(items);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Complete request processing by performing a dispatch into the servlet
|
||||
* container, where Spring MVC is invoked once more, and completes the
|
||||
@@ -302,8 +333,17 @@ public class ResponseBodyEmitter {
|
||||
*/
|
||||
interface Handler {
|
||||
|
||||
/**
|
||||
* Immediately write and flush the given data to the network.
|
||||
*/
|
||||
void send(Object data, @Nullable MediaType mediaType) throws IOException;
|
||||
|
||||
/**
|
||||
* Immediately write all data items then flush to the network.
|
||||
* @since 6.0.12
|
||||
*/
|
||||
void send(Set<DataWithMediaType> items) throws IOException;
|
||||
|
||||
void complete();
|
||||
|
||||
void completeWithError(Throwable failure);
|
||||
|
||||
@@ -20,6 +20,7 @@ import java.io.IOException;
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.Set;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
import jakarta.servlet.ServletRequest;
|
||||
@@ -202,6 +203,15 @@ public class ResponseBodyEmitterReturnValueHandler implements HandlerMethodRetur
|
||||
@Override
|
||||
public void send(Object data, @Nullable MediaType mediaType) throws IOException {
|
||||
sendInternal(data, mediaType);
|
||||
this.outputMessage.flush();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void send(Set<ResponseBodyEmitter.DataWithMediaType> items) throws IOException {
|
||||
for (ResponseBodyEmitter.DataWithMediaType item : items) {
|
||||
sendInternal(item.getData(), item.getMediaType());
|
||||
}
|
||||
this.outputMessage.flush();
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
@@ -209,7 +219,6 @@ public class ResponseBodyEmitterReturnValueHandler implements HandlerMethodRetur
|
||||
for (HttpMessageConverter<?> converter : ResponseBodyEmitterReturnValueHandler.this.sseMessageConverters) {
|
||||
if (converter.canWrite(data.getClass(), mediaType)) {
|
||||
((HttpMessageConverter<T>) converter).write(data, mediaType, this.outputMessage);
|
||||
this.outputMessage.flush();
|
||||
return;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2021 the original author or authors.
|
||||
* Copyright 2002-2023 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
@@ -123,9 +123,7 @@ public class SseEmitter extends ResponseBodyEmitter {
|
||||
public void send(SseEventBuilder builder) throws IOException {
|
||||
Set<DataWithMediaType> dataToSend = builder.build();
|
||||
synchronized (this) {
|
||||
for (DataWithMediaType entry : dataToSend) {
|
||||
super.send(entry.getData(), entry.getMediaType());
|
||||
}
|
||||
super.send(dataToSend);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user