Ява 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.