Ява 8+
Весенний ботинок 2
Весенняя сеть MVC
Обзор
Всем привет. Сегодня я хочу рассказать вам, как реализовать SSE (Server Sent Events) в Java.
Что такое события, отправленные сервером?
SSE — это технология, позволяющая передавать данные от сервера к клиенту в рамках одного HTTP-соединения в одном направлении. Давайте создадим приложение Spring Boot, используя инициализатор Spring.
Перейти на https://start.spring.io/
Выберите параметры проекта:
Тип проекта — Maven
Язык — Java
Версия Spring Boot — последняя стабильная версия (на момент написания Spring Boot 3.0.1)
Выберите зависимость Spring Web.

Щелкните Создать. Сохраните созданный шаблон проекта и откройте его в своей среде IDE. В моем случае это IntelliJ IDEA.
Представим, что у нас есть система, отображающая котировки акций. Соответственно, при изменении котировок мы должны отправлять обновленную информацию клиентам. Создадим простой класс, отражающий информацию об акции. Этот класс будет содержать только одно поле — цену акции.
import java.math.BigDecimal;
public class Stock {
private final BigDecimal price;
Stock(BigDecimal price) {
this.price = price;
}
public BigDecimal getPrice() {
return price;
}
}
Теперь давайте создадим класс, который будет имитировать изменение цены акции.
import org.springframework.context.ApplicationEventPublisher;
import org.springframework.stereotype.Component;
import java.math.BigDecimal;
import java.util.Random;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import static java.util.concurrent.TimeUnit.MILLISECONDS;
import static java.util.concurrent.TimeUnit.SECONDS;
@Component
public class StockPriceChanger {
private final ApplicationEventPublisher publisher; //1.1
private final Random random = new Random(); //1.2
private final ScheduledExecutorService executor = Executors.newSingleThreadScheduledExecutor(); //1.3
public StockPriceChanger(ApplicationEventPublisher publisher) {
this.publisher = publisher;
this.executor.schedule(this::changePrice, 1, SECONDS); //1.4
}
private void changePrice() {
BigDecimal price = BigDecimal.valueOf(random.nextGaussian()); //1.5
publisher.publishEvent(new Stock(price)); //1.6
executor.schedule(this::changePrice, random.nextInt(10000), MILLISECONDS); //1.7
}
}
У нас есть зависимость от ApplicationEventPublisher, внедренная через конструктор (1.1). Создайте экземпляр класса Random для генерации случайной цены акций (1.2).
Мы будем генерировать изменение цены в отдельном потоке через случайные промежутки времени. Для этого создайте однопоточный исполнитель, который может планировать выполнение команд после заданной задержки или периодически (1.3).
В конструкторе нашего класса запустим поток генерации случайной цены акции (1.4). В методе changePrice мы получаем случайное значение цены (1,5), публикуем событие Stock для всех подписчиков (1,6) и планируем создание следующего значения со случайной задержкой (1,7).
Как вы заметили, у нас есть зависимость класса от ApplicationEventPublisher. Spring «из коробки» предоставляет простой механизм обработки событий, уменьшающий связность компонентов системы. Событие, происходящее в одной точке приложения, может быть перехвачено и обработано в любой другой части приложения благодаря таким сущностям, как издатель и прослушиватель событий.
Затем создайте контроллер для обработки запросов от клиентов.
import org.springframework.context.event.EventListener;
import org.springframework.http.MediaType;
import org.springframework.scheduling.annotation.Async;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RestController;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
import java.util.ArrayList;
import java.util.List;
import java.util.Set;
import java.util.concurrent.CopyOnWriteArraySet;
@RestController
public class StockController {
private final Set<SseEmitter> clients = new CopyOnWriteArraySet<>(); //2.1
@GetMapping("/stocks-stream") //2.2
public SseEmitter stocksStream() {
SseEmitter sseEmitter = new SseEmitter();
clients.add(sseEmitter);
sseEmitter.onTimeout(() -> clients.remove(sseEmitter));
sseEmitter.onError(throwable -> clients.remove(sseEmitter));
return sseEmitter;
}
@Async
@EventListener
public void stockMessageHandler(Stock stock) { //2.3
List<SseEmitter> errorEmitters = new ArrayList<>();
clients.forEach(emitter -> {
try {
emitter.send(stock, MediaType.APPLICATION_JSON); //2.4
} catch (Exception e) {
errorEmitters.add(emitter);
}
});
errorEmitters.forEach(clients::remove); //2.5
}
}
Spring 4.2 представил новый класс ResponseBodyEmitter, который используется в качестве возвращаемого типа в контроллерах Spring Web MVC для асинхронных запросов. ResponseBodyEmitter можно использовать для отправки нескольких объектов, где каждый объект записывается с помощью совместимого HttpMessageConverter.
SseEmitter расширяет ResponseBodyEmitter и позволяет отправлять множество сообщений в ответ на один запрос, как того требует протокол SSE.
Рассмотрим этот класс поближе:
Создадим набор, в котором будем хранить подключенных клиентов (1). Конечная точка, которая регистрирует клиента для получения уведомлений об изменении цен (2). В этом методе мы создаем новый SseEmitter, сохраняем его, чтобы мы могли отправлять ему уведомления в будущем. И регистрируем два обработчика, onTimeout и onError. При возникновении ошибки или тайм-аута мы удалим клиента из списка зарегистрированных клиентов.
Обработчик сообщений (2.3). Этот метод имеет две аннотации @Async и @EventListener. Аннотация @EventListener используется для указания того, что метод является обработчиком событий. Тип событий, которые он обрабатывает, определяется типом аргумента, в данном случае Stock. По умолчанию прослушиватель вызывается синхронно. Однако мы можем легко сделать его асинхронным, добавив аннотацию @Async. Аннотирование метода компонента с помощью @Async заставит его выполняться в отдельном потоке. Другими словами, вызывающая сторона не будет ждать завершения вызываемого метода. Наш метод принимает новое событие со значением цены акции и асинхронно отправляет его всем клиентам в формате JSON (2.4). Если при отправке возникает ошибка, мы сохраняем этот неудачный эмиттер, а затем удаляем его из списка активных клиентов (2.5).
Чтобы Spring распознавал методы, помеченные аннотацией @Async, и запускал эти методы в пуле фоновых потоков, аннотацию @EnableAsync необходимо добавить над классом конфигурации (3.1)
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.scheduling.annotation.EnableAsync;
@EnableAsync //3.1
@SpringBootApplication
public class SsewebmvcApplication {
public static void main(String[] args) {
SpringApplication.run(SsewebmvcApplication.class, args);
}
}
Теперь давайте создадим простой пользовательский интерфейс, чтобы увидеть, что у нас есть.
Давайте добавим файл index.html в наш проект по пути src/main/resources/static.

И давайте напишем простой код для отображения изменений котировок акций:
<!DOCTYPE html>
<html lang="en">
<head>
<meta charset="UTF-8">
<title>SSE example</title>
</head>
<body>
<ul id="stockChanges"></ul>
<script type="application/javascript">
function add(message) {
const li = document.createElement("li");
li.innerHTML = message;
document.getElementById("stockChanges").appendChild(li);
}
const eventSource = new EventSource("/stocks-stream"); //4.1
eventSource.onmessage = e => {
const response = JSON.parse(e.data);
add('Stock price was changed, new price: ' + response.price + ' $'); //4.2
}
eventSource.onopen = e => add('Connection opened');
eventSource.onerror = e => add('Connection closed');
</script>
</body>
</html>
Чтобы начать получать данные, мы создаем новый EventSource(“/stocks-stream”) (4.1). Браузер подключится к /stocks-stream и будет держать его открытым в ожидании события.
По умолчанию объект EventSource генерирует 3 события:
- message — сообщение получено, доступно как e.data.
- open — соединение открыто.
- error — соединение не удалось, т.е. сервер вернул статус 500.
В обработчике события message читаем полученную котировку акций и создаем новый элемент списка (4.2)
Наконец, давайте запустим наше приложение и посмотрим, что произойдет. После запуска приложения нам нужно перейти по адресу http://localhost:8080 и мы увидим примерно следующую картину.

Спасибо за внимание.
Полный код примера вы можете получить на GitHub.