Spec-Zone.ru › Spring Boot

Сообщения

Фреймворк Spring предоставляет обширную поддержку интеграции с системами обмена сообщениями, от упрощенного использования API JMS с помощью JmsTemplate до полной инфраструктуры для асинхронного получения сообщений. Spring AMQP предоставляет аналогичный набор функций для протокола Advanced Message Queuing Protocol. Spring Boot также предоставляет параметры автоконфигурации для RabbitTemplate и RabbitMQ. Spring WebSocket по умолчанию поддерживает обмен сообщениями STOMP, и Spring Boot поддерживает это через стартеры и небольшое количество автоконфигурации. Spring Boot также поддерживает Apache Kafka.

1. JMS

Интерфейс jakarta.jms.ConnectionFactory предоставляет стандартный метод создания jakarta.jms.Connection для взаимодействия с брокером JMS. Хотя Spring требует ConnectionFactory для работы с JMS, вам, как правило, не нужно использовать его напрямую, а вместо этого можно полагаться на более высокие уровни абстракций обмена сообщениями. (Подробности см. в соответствующем разделе документации по Spring Framework.) Spring Boot также автоматически настраивает необходимую инфраструктуру для отправки и получения сообщений.

1.1. Поддержка ActiveMQ

Когда ActiveMQ доступен в классе, Spring Boot может настроить ConnectionFactory.

Если вы используете spring-boot-starter-activemq, необходимые зависимости для подключения к экземпляру ActiveMQ предоставляются, как и инфраструктура Spring для интеграции с JMS.

Настройка ActiveMQ контролируется внешними свойствами конфигурации в spring.activemq.*. По умолчанию ActiveMQ автоматически настраивается для использования транспорта TCP, подключаясь по умолчанию к tcp://localhost:61616. Следующий пример демонстрирует, как изменить URL-адрес брокера по умолчанию:

Свойства
spring.activemq.broker-url=tcp://192.168.1.210:9876
spring.activemq.user=admin
spring.activemq.password=secret
Yaml
spring:
  activemq:
    broker-url: "tcp://192.168.1.210:9876"
    user: "admin"
    password: "secret"

По умолчанию CachingConnectionFactory оборачивает родной ConnectionFactory с разумными настройками, которые вы можете контролировать с помощью внешних свойств конфигурации в spring.jms.*:

Свойства
spring.jms.cache.session-cache-size=5
Yaml
spring:
  jms:
    cache:
      session-cache-size: 5

Если вы предпочитаете использовать родной пуллинг, вы можете сделать это, добавив зависимость от org.messaginghub:pooled-jms и настроив JmsPoolConnectionFactory соответственно, как показано в следующем примере:

Свойства
spring.activemq.pool.enabled=true
spring.activemq.pool.max-connections=50
Yaml
spring:
  activemq:
    pool:
      enabled: true
      max-connections: 50
См. ActiveMQProperties для получения дополнительных поддерживаемых параметров. Вы также можете зарегистрировать любое количество бинов, реализующих ActiveMQConnectionFactoryCustomizer для более сложных кастомизаций.

По умолчанию ActiveMQ создает пункт назначения, если он еще не существует, чтобы пункты назначения разрешались по их предоставленным именам.

1.2. Поддержка ActiveMQ Artemis

Spring Boot может автоматически настроить ConnectionFactory при обнаружении ActiveMQ Artemis в классе. Если брокер присутствует, встроенный брокер автоматически запускается и настраивается (если свойство режима не было явно задано). Поддерживаемые режимы — embedded (чтобы сделать явным, что встроенный брокер необходим и что должна произойти ошибка, если брокер не доступен в классе) и native (для подключения к брокеру с использованием протокола транспорта netty). В последнем случае Spring Boot настраивает ConnectionFactory, который подключается к брокеру, работающему на локальной машине, с настройками по умолчанию.

Если вы используете spring-boot-starter-artemis, предоставляются необходимые зависимости для подключения к существующему экземпляру ActiveMQ Artemis, а также инфраструктура Spring для интеграции с JMS. Добавление org.apache.activemq:artemis-jakarta-server в ваше приложение позволяет использовать режим встраивания.

Настройка ActiveMQ Artemis контролируется внешними свойствами конфигурации в spring.artemis.*. Например, вы можете объявить следующий раздел в application.properties:

Свойства
spring.artemis.mode=native
spring.artemis.broker-url=tcp://192.168.1.210:9876
spring.artemis.user=admin
spring.artemis.password=secret
Yaml
spring:
  artemis:
    mode: native
    broker-url: "tcp://192.168.1.210:9876"
    user: "admin"
    password: "secret"

При встраивании брокера вы можете выбрать, хотите ли вы включить сохранение и перечислить пункты назначения, которые должны быть доступны. Их можно указать как список, разделенный запятыми, чтобы создать их с параметрами по умолчанию, или вы можете определить бины типа org.apache.activemq.artemis.jms.server.config.JMSQueueConfiguration или org.apache.activemq.artemis.jms.server.config.TopicConfiguration, для расширенной настройки очередей и тем, соответственно.

По умолчанию CachingConnectionFactory оборачивает родной ConnectionFactory с разумными настройками, которые можно контролировать с помощью внешних свойств конфигурации в spring.jms.*:

Свойства
spring.jms.cache.session-cache-size=5
Yaml
spring:
  jms:
    cache:
      session-cache-size: 5

Если вы предпочитаете использовать родной пуллинг, вы можете сделать это, добавив зависимость от org.messaginghub:pooled-jms и настроив JmsPoolConnectionFactory соответственно, как показано в следующем примере:

Свойства
spring.artemis.pool.enabled=true
spring.artemis.pool.max-connections=50
Yaml
spring:
  artemis:
    pool:
      enabled: true
      max-connections: 50

См. ArtemisProperties для получения дополнительных поддерживаемых параметров.

JNDI-поиск не используется, а пункты назначения разрешаются по их именам, используя атрибут name в настройке Artemis или имена, предоставленные через конфигурацию.

1.3. Использование ConnectionFactory JNDI

Если ваше приложение выполняется в прикладном сервере, Spring Boot пытается найти JMS ConnectionFactory с помощью JNDI. По умолчанию проверяются расположения java:/JmsXA и java:/XAConnectionFactory. Вы можете использовать свойство spring.jms.jndi-name, если нужно указать альтернативное расположение, как показано в следующем примере:

Свойства
spring.jms.jndi-name=java:/MyConnectionFactory
Yaml
spring:
  jms:
    jndi-name: "java:/MyConnectionFactory"

1.4. Отправка сообщения

Автоконфигурируется JmsTemplate Spring, и вы можете напрямую автоматизировать его в свои собственные бины, как показано в следующем примере:

Java
import org.springframework.jms.core.JmsTemplate;
import org.springframework.stereotype.Component;

@Component
public class MyBean {

    private final JmsTemplate jmsTemplate;

    public MyBean(JmsTemplate jmsTemplate) {
        this.jmsTemplate = jmsTemplate;
    }

    // ...

    public void someMethod() {
        this.jmsTemplate.convertAndSend("hello");
    }

}
Kotlin
import org.springframework.jms.core.JmsTemplate
import org.springframework.stereotype.Component

@Component
class MyBean(private val jmsTemplate: JmsTemplate) {

    // ...

    fun someMethod() {
        jmsTemplate.convertAndSend("hello")
    }

}
JmsMessagingTemplate можно вводить аналогичным образом. Если определен бины DestinationResolver или MessageConverter, он автоматически связывается с автоматически настроенным JmsTemplate.

1.5. Прием сообщения

При наличии инфраструктуры JMS любой бин может быть аннотирован с помощью @JmsListener для создания точки входа слушателя. Если JmsListenerContainerFactory не определено, то оно автоматически настраивается по умолчанию. Если бины DestinationResolver, MessageConverter или jakarta.jms.ExceptionListener определены, они автоматически ассоциируются с фабрикой по умолчанию.

По умолчанию, фабрика по умолчанию транзакционная. Если вы работаете в инфраструктуре, где присутствует JtaTransactionManager, она по умолчанию связывается с контейнером слушателя. В противном случае, включается флаг sessionTransacted. В последнем случае, вы можете связать транзакцию вашего локального хранилища данных с обработкой входящего сообщения, добавив @Transactional в метод слушателя (или делегат). Это гарантирует подтверждение входящего сообщения после завершения локальной транзакции. Это также включает отправку ответных сообщений, которые были выполнены в той же сессии JMS.

Следующий компонент создает точку входа слушателя на целевом пункте someQueue:

Java
import org.springframework.jms.annotation.JmsListener;
import org.springframework.stereotype.Component;

@Component
public class MyBean {

    @JmsListener(destination = "someQueue")
    public void processMessage(String content) {
        // ...
    }

}
Kotlin
import org.springframework.jms.annotation.JmsListener
import org.springframework.stereotype.Component

@Component
class MyBean {

    @JmsListener(destination = "someQueue")
    fun processMessage(content: String?) {
        // ...
    }

}
Смотрите Javadoc @EnableJms для получения более подробной информации.

Если вам нужно создать больше экземпляров JmsListenerContainerFactory или если вы хотите переопределить значение по умолчанию, Spring Boot предоставляет DefaultJmsListenerContainerFactoryConfigurer, который вы можете использовать для инициализации DefaultJmsListenerContainerFactory с теми же настройками, что и автоматически настроенный экземпляр.

Например, следующий пример экспонирует другую фабрику, использующую определенный MessageConverter:

Java
import jakarta.jms.ConnectionFactory;

import org.springframework.boot.autoconfigure.jms.DefaultJmsListenerContainerFactoryConfigurer;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.jms.config.DefaultJmsListenerContainerFactory;

@Configuration(proxyBeanMethods = false)
public class MyJmsConfiguration {

    @Bean
    public DefaultJmsListenerContainerFactory myFactory(DefaultJmsListenerContainerFactoryConfigurer configurer) {
        DefaultJmsListenerContainerFactory factory = new DefaultJmsListenerContainerFactory();
        ConnectionFactory connectionFactory = getCustomConnectionFactory();
        configurer.configure(factory, connectionFactory);
        factory.setMessageConverter(new MyMessageConverter());
        return factory;
    }

    private ConnectionFactory getCustomConnectionFactory() {
        return ...
    }

}
Kotlin
import jakarta.jms.ConnectionFactory
import org.springframework.boot.autoconfigure.jms.DefaultJmsListenerContainerFactoryConfigurer
import org.springframework.context.annotation.Bean
import org.springframework.context.annotation.Configuration
import org.springframework.jms.config.DefaultJmsListenerContainerFactory

@Configuration(proxyBeanMethods = false)
class MyJmsConfiguration {

    @Bean
    fun myFactory(configurer: DefaultJmsListenerContainerFactoryConfigurer): DefaultJmsListenerContainerFactory {
        val factory = DefaultJmsListenerContainerFactory()
        val connectionFactory = getCustomConnectionFactory()
        configurer.configure(factory, connectionFactory)
        factory.setMessageConverter(MyMessageConverter())
        return factory
    }

    fun getCustomConnectionFactory() : ConnectionFactory? {
        return ...
    }

}

Затем вы можете использовать фабрику в любом @JmsListener-аннотированном методе следующим образом:

Java
import org.springframework.jms.annotation.JmsListener;
import org.springframework.stereotype.Component;

@Component
public class MyBean {

    @JmsListener(destination = "someQueue", containerFactory = "myFactory")
    public void processMessage(String content) {
        // ...
    }

}
Kotlin
import org.springframework.jms.annotation.JmsListener
import org.springframework.stereotype.Component

@Component
class MyBean {

    @JmsListener(destination = "someQueue", containerFactory = "myFactory")
    fun processMessage(content: String?) {
        // ...
    }

}

2. AMQP

Протокол Advanced Message Queuing Protocol (AMQP) — это платформенно-нейтральный протокол на уровне проводов для ориентированных на сообщения систем middleware. Проект Spring AMQP применяет основные концепции Spring к разработке решений обмена сообщениями на основе AMQP. Spring Boot предлагает несколько удобств для работы с AMQP через RabbitMQ, включая spring-boot-starter-amqp «Starter».

2.1. Поддержка RabbitMQ

RabbitMQ — это лёгкий, надёжный, масштабируемый и переносимый брокер сообщений, основанный на протоколе AMQP. Spring использует RabbitMQ для взаимодействия через протокол AMQP.

Настройка RabbitMQ контролируется внешними свойствами конфигурации в spring.rabbitmq.*. Например, вы можете объявить следующий раздел в application.properties:

Свойства
spring.rabbitmq.host=localhost
spring.rabbitmq.port=5672
spring.rabbitmq.username=admin
spring.rabbitmq.password=secret
Yaml
spring:
  rabbitmq:
    host: "localhost"
    port: 5672
    username: "admin"
    password: "secret"

Альтернативно, вы можете настроить то же соединение с помощью атрибута addresses:

Свойства
spring.rabbitmq.addresses=amqp://admin:secret@localhost
Yaml
spring:
  rabbitmq:
    addresses: "amqp://admin:secret@localhost"
При указании адресов таким образом, свойства host и port игнорируются. Если адрес использует протокол amqps, поддержка SSL автоматически включается.

См. RabbitProperties для получения дополнительных вариантов конфигурации на основе свойств. Для настройки более низкоуровневых деталей экземпляра RabbitMQ ConnectionFactory, используемого Spring AMQP, определите бин ConnectionFactoryCustomizer.

Если бин ConnectionNameStrategy существует в контексте, он будет автоматически использован для именования соединений, созданных автоматически настроенным CachingConnectionFactory.

Для добавления изменений в RabbitTemplate на уровне всего приложения используйте бин RabbitTemplateCustomizer.

См. Понимание AMQP, протокола, используемого RabbitMQ для получения дополнительных сведений.

2.2. Отправка сообщения

Spring’s AmqpTemplate и AmqpAdmin автоматически настраиваются, и вы можете напрямую их внедрить в собственные бины, как показано в следующем примере:

Java
import org.springframework.amqp.core.AmqpAdmin;
import org.springframework.amqp.core.AmqpTemplate;
import org.springframework.stereotype.Component;

@Component
public class MyBean {

    private final AmqpAdmin amqpAdmin;

    private final AmqpTemplate amqpTemplate;

    public MyBean(AmqpAdmin amqpAdmin, AmqpTemplate amqpTemplate) {
        this.amqpAdmin = amqpAdmin;
        this.amqpTemplate = amqpTemplate;
    }

    // ...

    public void someMethod() {
        this.amqpAdmin.getQueueInfo("someQueue");
    }

    public void someOtherMethod() {
        this.amqpTemplate.convertAndSend("hello");
    }

}
Kotlin
import org.springframework.amqp.core.AmqpAdmin
import org.springframework.amqp.core.AmqpTemplate
import org.springframework.stereotype.Component

@Component
class MyBean(private val amqpAdmin: AmqpAdmin, private val amqpTemplate: AmqpTemplate) {

    // ...

    fun someMethod() {
        amqpAdmin.getQueueInfo("someQueue")
    }

    fun someOtherMethod() {
        amqpTemplate.convertAndSend("hello")
    }

}
RabbitMessagingTemplate может быть внедрён аналогичным образом. Если бин MessageConverter определён, он автоматически связан с автоматически настроенным AmqpTemplate.

Если необходимо, любой org.springframework.amqp.core.Queue бин, определённый как бин, автоматически используется для объявления соответствующего очереди на экземпляре RabbitMQ.

Чтобы повторить операции, можно включить повторы на AmqpTemplate (например, в случае потери соединения с брокером):

Свойства
spring.rabbitmq.template.retry.enabled=true
spring.rabbitmq.template.retry.initial-interval=2s
Yaml
spring:
  rabbitmq:
    template:
      retry:
        enabled: true
        initial-interval: "2s"

Повторные попытки отключены по умолчанию. Вы также можете настроить RetryTemplate программно, объявив бин RabbitRetryTemplateCustomizer.

Если вам нужно создать больше RabbitTemplate экземпляров или если вы хотите переопределить значения по умолчанию, Spring Boot предоставляет бин RabbitTemplateConfigurer, который вы можете использовать для инициализации RabbitTemplate с теми же настройками, что и фабрики, используемые автонастройкой.

2.3. Отправка сообщения в поток

Чтобы отправить сообщение в определённый поток, укажите имя потока, как показано в следующем примере:

Свойства
spring.rabbitmq.stream.name=my-stream
Yaml
spring:
  rabbitmq:
    stream:
      name: "my-stream"

Если определён бин MessageConverter, StreamMessageConverter или ProducerCustomizer, он автоматически связан с автоматически настроенным RabbitStreamTemplate.

Если вам нужно создать больше RabbitStreamTemplate экземпляров или если вы хотите переопределить значения по умолчанию, Spring Boot предоставляет бин RabbitStreamTemplateConfigurer, который вы можете использовать для инициализации RabbitStreamTemplate с теми же настройками, что и фабрики, используемые автонастройкой.

2.4. Получение сообщения

При наличии инфраструктуры Rabbit, любой бин может быть аннотирован @RabbitListener, чтобы создать точку входа для слушателя. Если RabbitListenerContainerFactory не определён, по умолчанию настраивается SimpleRabbitListenerContainerFactory и вы можете переключиться на прямой контейнер с помощью свойства spring.rabbitmq.listener.type . Если определён бин MessageConverter или MessageRecoverer, он автоматически ассоциируется с фабрикой по умолчанию.

Следующий компонент создаёт точку входа для слушателя в очереди someQueue:

Java
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;

@Component
public class MyBean {

    @RabbitListener(queues = "someQueue")
    public void processMessage(String content) {
        // ...
    }

}
Kotlin
import org.springframework.amqp.rabbit.annotation.RabbitListener
import org.springframework.stereotype.Component

@Component
class MyBean {

    @RabbitListener(queues = ["someQueue"])
    fun processMessage(content: String?) {
        // ...
    }

}
См. документацию @EnableRabbit для получения дополнительных сведений.

Если вам нужно создать больше RabbitListenerContainerFactory экземпляров или если вы хотите переопределить значения по умолчанию, Spring Boot предоставляет бин SimpleRabbitListenerContainerFactoryConfigurer и DirectRabbitListenerContainerFactoryConfigurer, которые вы можете использовать для инициализации SimpleRabbitListenerContainerFactory и DirectRabbitListenerContainerFactory с теми же настройками, что и фабрики, используемые автонастройкой.

Тип контейнера, который вы выбрали, не имеет значения. Эти два бинa доступны через автонастройку.

Например, следующий класс конфигурации предоставляет другую фабрику, использующую определённый MessageConverter:

Java
import org.springframework.amqp.rabbit.config.SimpleRabbitListenerContainerFactory;
import org.springframework.amqp.rabbit.connection.ConnectionFactory;
import org.springframework.boot.autoconfigure.amqp.SimpleRabbitListenerContainerFactoryConfigurer;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

@Configuration(proxyBeanMethods = false)
public class MyRabbitConfiguration {

    @Bean
    public SimpleRabbitListenerContainerFactory myFactory(SimpleRabbitListenerContainerFactoryConfigurer configurer) {
        SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory();
        ConnectionFactory connectionFactory = getCustomConnectionFactory();
        configurer.configure(factory, connectionFactory);
        factory.setMessageConverter(new MyMessageConverter());
        return factory;
    }

    private ConnectionFactory getCustomConnectionFactory() {
        return ...
    }

}
Kotlin
import org.springframework.amqp.rabbit.config.SimpleRabbitListenerContainerFactory
import org.springframework.amqp.rabbit.connection.ConnectionFactory
import org.springframework.boot.autoconfigure.amqp.SimpleRabbitListenerContainerFactoryConfigurer
import org.springframework.context.annotation.Bean
import org.springframework.context.annotation.Configuration

@Configuration(proxyBeanMethods = false)
class MyRabbitConfiguration {

    @Bean
    fun myFactory(configurer: SimpleRabbitListenerContainerFactoryConfigurer): SimpleRabbitListenerContainerFactory {
        val factory = SimpleRabbitListenerContainerFactory()
        val connectionFactory = getCustomConnectionFactory()
        configurer.configure(factory, connectionFactory)
        factory.setMessageConverter(MyMessageConverter())
        return factory
    }

    fun getCustomConnectionFactory() : ConnectionFactory? {
        return ...
    }

}

Затем вы можете использовать фабрику в любом методе, аннотированном @RabbitListener:

Java
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;

@Component
public class MyBean {

    @RabbitListener(queues = "someQueue", containerFactory = "myFactory")
    public void processMessage(String content) {
        // ...
    }

}
Kotlin
import org.springframework.amqp.rabbit.annotation.RabbitListener
import org.springframework.stereotype.Component

@Component
class MyBean {

    @RabbitListener(queues = ["someQueue"], containerFactory = "myFactory")
    fun processMessage(content: String?) {
        // ...
    }

}

Вы можете включить повторы для обработки ситуаций, когда ваш слушатель выбрасывает исключение. По умолчанию используется RejectAndDontRequeueRecoverer , но вы можете определить свой собственный MessageRecoverer. Когда попытки повтора исчерпаны, сообщение отклоняется и либо удаляется, либо направляется в обменник для сообщений, если брокером настроено такое поведение. По умолчанию повторы отключены. Вы также можете настроить RetryTemplate программно, объявив бин RabbitRetryTemplateCustomizer.

По умолчанию, если повторы отключены и слушатель выбрасывает исключение, доставка повторяется неограниченное количество раз. Вы можете изменить это поведение двумя способами: установите свойство defaultRequeueRejected в false, чтобы не предпринималось никаких повторных попыток доставки, или выбросьте AmqpRejectAndDontRequeueException, чтобы указать, что сообщение должно быть отклонено. Последний вариант используется, когда повторы включены, и максимальное количество попыток доставки достигнуто.

3. Поддержка Apache Kafka

Apache Kafka поддерживается за счёт автоматической конфигурации проекта spring-kafka.

Конфигурация Kafka управляется внешними свойствами конфигурации в spring.kafka.*. Например, вы можете объявить следующий раздел в application.properties:

Свойства
spring.kafka.bootstrap-servers=localhost:9092
spring.kafka.consumer.group-id=myGroup
Yaml
spring:
  kafka:
    bootstrap-servers: "localhost:9092"
    consumer:
      group-id: "myGroup"
Для создания темы при запуске добавьте bean типа NewTopic. Если тема уже существует, bean игнорируется.

См. KafkaProperties для получения более подробной информации о поддерживаемых вариантах.

3.1. Отправка сообщения

Spring’s KafkaTemplate автоматически настраивается, и вы можете напрямую внедрить его в свои собственные компоненты, как показано в следующем примере:

Java
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.stereotype.Component;

@Component
public class MyBean {

    private final KafkaTemplate<String, String> kafkaTemplate;

    public MyBean(KafkaTemplate<String, String> kafkaTemplate) {
        this.kafkaTemplate = kafkaTemplate;
    }

    // ...

    public void someMethod() {
        this.kafkaTemplate.send("someTopic", "Hello");
    }

}
Kotlin
import org.springframework.kafka.core.KafkaTemplate
import org.springframework.stereotype.Component

@Component
class MyBean(private val kafkaTemplate: KafkaTemplate<String, String>) {

    // ...

    fun someMethod() {
        kafkaTemplate.send("someTopic", "Hello")
    }

}
Если свойство spring.kafka.producer.transaction-id-prefix определено, автоматически настраивается KafkaTransactionManager. Кроме того, если bean RecordMessageConverter определён, он автоматически связывается с автоматически настроенным KafkaTemplate .

3.2. Приём сообщения

При наличии инфраструктуры Apache Kafka любой компонент может быть аннотирован с помощью @KafkaListener для создания точки входа для прослушивания. Если KafkaListenerContainerFactory не определено, автоматически настраивается значение по умолчанию с ключами, определёнными в spring.kafka.listener.*.

Следующий компонент создаёт точку входа для прослушивания на теме someTopic:

Java
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.stereotype.Component;

@Component
public class MyBean {

    @KafkaListener(topics = "someTopic")
    public void processMessage(String content) {
        // ...
    }

}
Kotlin
import org.springframework.kafka.annotation.KafkaListener
import org.springframework.stereotype.Component

@Component
class MyBean {

    @KafkaListener(topics = ["someTopic"])
    fun processMessage(content: String?) {
        // ...
    }

}

Если bean KafkaTransactionManager определён, он автоматически связывается с фабрикой контейнера. Аналогично, если определён bean RecordFilterStrategy, CommonErrorHandler, AfterRollbackProcessor или ConsumerAwareRebalanceListener, он автоматически связывается с фабрикой по умолчанию.

В зависимости от типа прослушивателя, bean RecordMessageConverter или BatchMessageConverter связывается с фабрикой по умолчанию. Если для прослушивателя пакетной обработки присутствует только bean RecordMessageConverter, он оборачивается в BatchMessageConverter.

Настраиваемый ChainedKafkaTransactionManager должен быть помечен как @Primary, так как он обычно ссылается на автоматически настроенный bean KafkaTransactionManager .

3.3. Kafka Streams

Spring для Apache Kafka предоставляет фабричный bean для создания объекта StreamsBuilder и управления жизненным циклом его потоков. Spring Boot автоматически настраивает необходимый bean KafkaStreamsConfiguration, если kafka-streams находится в классе, и Kafka Streams включено аннотацией @EnableKafkaStreams.

Включение Kafka Streams означает, что необходимо установить идентификатор приложения и серверы Bootstrap. Первый можно настроить с помощью spring.kafka.streams.application-id, по умолчанию spring.application.name, если не указано другое. Последний можно установить глобально или переопределить только для потоков.

Несколько дополнительных свойств доступны с помощью специальных свойств; другие произвольные свойства Kafka могут быть установлены с помощью пространства имён spring.kafka.streams.properties. См. также Дополнительные свойства Kafka для получения дополнительной информации.

Чтобы использовать фабричный bean, подключите StreamsBuilder к вашему @Bean, как показано в следующем примере:

Java
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.KeyValue;
import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.kstream.KStream;
import org.apache.kafka.streams.kstream.Produced;

import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.annotation.EnableKafkaStreams;
import org.springframework.kafka.support.serializer.JsonSerde;

@Configuration(proxyBeanMethods = false)
@EnableKafkaStreams
public class MyKafkaStreamsConfiguration {

    @Bean
    public KStream<Integer, String> kStream(StreamsBuilder streamsBuilder) {
        KStream<Integer, String> stream = streamsBuilder.stream("ks1In");
        stream.map(this::uppercaseValue).to("ks1Out", Produced.with(Serdes.Integer(), new JsonSerde<>()));
        return stream;
    }

    private KeyValue<Integer, String> uppercaseValue(Integer key, String value) {
        return new KeyValue<>(key, value.toUpperCase());
    }

}
Kotlin
import org.apache.kafka.common.serialization.Serdes
import org.apache.kafka.streams.KeyValue
import org.apache.kafka.streams.StreamsBuilder
import org.apache.kafka.streams.kstream.KStream
import org.apache.kafka.streams.kstream.Produced
import org.springframework.context.annotation.Bean
import org.springframework.context.annotation.Configuration
import org.springframework.kafka.annotation.EnableKafkaStreams
import org.springframework.kafka.support.serializer.JsonSerde

@Configuration(proxyBeanMethods = false)
@EnableKafkaStreams
class MyKafkaStreamsConfiguration {

    @Bean
    fun kStream(streamsBuilder: StreamsBuilder): KStream<Int, String> {
        val stream = streamsBuilder.stream<Int, String>("ks1In")
        stream.map(this::uppercaseValue).to("ks1Out", Produced.with(Serdes.Integer(), JsonSerde()))
        return stream
    }

    private fun uppercaseValue(key: Int, value: String): KeyValue<Int?, String?> {
        return KeyValue(key, value.uppercase())
    }

}

По умолчанию потоки, управляемые объектом StreamBuilder, запускаются автоматически. Вы можете настроить это поведение с помощью свойства spring.kafka.streams.auto-startup.

3.4. Дополнительные свойства Kafka

Поддерживаемые свойства автоматической настройки показаны в разделе “Свойства интеграции” приложения. Обратите внимание, что в большинстве случаев эти свойства (с дефисом или в camelCase) напрямую соответствуют свойствам Apache Kafka с точками. См. документацию Apache Kafka для получения подробностей.

Свойства, не содержащие тип клиента (producer, consumer, admin, или streams) в своём имени, считаются общими и применяются ко всем клиентам. Большинство этих общих свойств могут быть переопределены для одного или нескольких типов клиентов, если это необходимо.

Apache Kafka назначает свойства с важностью HIGH, MEDIUM или LOW. Spring Boot автоконфигурация поддерживает все свойства с высокой важностью, некоторые выбранные свойства со средней и низкой важностью, а также любые свойства, у которых нет значения по умолчанию.

Только подмножество свойств, поддерживаемых Kafka, доступно непосредственно через класс KafkaProperties. Если вы хотите настроить отдельные типы клиентов с дополнительными свойствами, которые не поддерживаются напрямую, используйте следующие свойства:

Свойства
spring.kafka.properties[prop.one]=first
spring.kafka.admin.properties[prop.two]=second
spring.kafka.consumer.properties[prop.three]=third
spring.kafka.producer.properties[prop.four]=fourth
spring.kafka.streams.properties[prop.five]=fifth
Yaml
spring:
  kafka:
    properties:
      "[prop.one]": "first"
    admin:
      properties:
        "[prop.two]": "second"
    consumer:
      properties:
        "[prop.three]": "third"
    producer:
      properties:
        "[prop.four]": "fourth"
    streams:
      properties:
        "[prop.five]": "fifth"

Это устанавливает общее свойство prop.one Kafka со значением first (применяется к производителям, потребителям, администраторам и потокам), свойство администратора prop.two со значением second, свойство потребителя prop.three со значением third, свойство производителя prop.four со значением fourth и свойство потоков prop.five со значением fifth.

Также можно настроить Spring Kafka JsonDeserializer следующим образом:

Свойства
spring.kafka.consumer.value-deserializer=org.springframework.kafka.support.serializer.JsonDeserializer
spring.kafka.consumer.properties[spring.json.value.default.type]=com.example.Invoice
spring.kafka.consumer.properties[spring.json.trusted.packages]=com.example.main,com.example.another
Yaml
spring:
  kafka:
    consumer:
      value-deserializer: "org.springframework.kafka.support.serializer.JsonDeserializer"
      properties:
        "[spring.json.value.default.type]": "com.example.Invoice"
        "[spring.json.trusted.packages]": "com.example.main,com.example.another"

Аналогично, можно отключить поведение по умолчанию JsonSerializer отправки информации о типе в заголовках:

Свойства
spring.kafka.producer.value-serializer=org.springframework.kafka.support.serializer.JsonSerializer
spring.kafka.producer.properties[spring.json.add.type.headers]=false
Yaml
spring:
  kafka:
    producer:
      value-serializer: "org.springframework.kafka.support.serializer.JsonSerializer"
      properties:
        "[spring.json.add.type.headers]": false
Свойства, заданные таким образом, переопределяют любые конфигурационные элементы, которые Spring Boot явно поддерживает.

3.5. Тестирование с встроенным Kafka

Spring для Apache Kafka предоставляет удобный способ тестирования проектов с встроенным брокером Apache Kafka. Чтобы использовать эту функцию, аннотируйте класс теста с помощью @EmbeddedKafka из модуля spring-kafka-test. Более подробная информация представлена в справочном руководстве Spring для Apache Kafka здесь.

Чтобы заставить Spring Boot автоконфигурацию работать с указанным выше встроенным брокером Apache Kafka, нужно перемапить системную переменную для адресов встроенного брокера (заполненную EmbeddedKafkaBroker) в свойство конфигурации Spring Boot для Apache Kafka. Существует несколько способов сделать это:

  • Укажите системную переменную для сопоставления адресов встроенного брокера со свойством spring.kafka.bootstrap-servers в классе теста:

Java
static {
    System.setProperty(EmbeddedKafkaBroker.BROKER_LIST_PROPERTY, "spring.kafka.bootstrap-servers");
}
Kotlin
init {
    System.setProperty(EmbeddedKafkaBroker.BROKER_LIST_PROPERTY, "spring.kafka.bootstrap-servers")
}
  • Настройте имя свойства в аннотации @EmbeddedKafka:

Java
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.kafka.test.context.EmbeddedKafka;

@SpringBootTest
@EmbeddedKafka(topics = "someTopic", bootstrapServersProperty = "spring.kafka.bootstrap-servers")
class MyTest {

    // ...

}
Kotlin
import org.springframework.boot.test.context.SpringBootTest
import org.springframework.kafka.test.context.EmbeddedKafka

@SpringBootTest
@EmbeddedKafka(topics = ["someTopic"], bootstrapServersProperty = "spring.kafka.bootstrap-servers")
class MyTest {

    // ...

}
  • Используйте плейсхолдер в свойствах конфигурации:

Свойства
spring.kafka.bootstrap-servers=${spring.embedded.kafka.brokers}
Yaml
spring:
  kafka:
    bootstrap-servers: "${spring.embedded.kafka.brokers}"

4. RSocket

RSocket — это двоичный протокол для использования на транспортных каналах потоков байтов. Он позволяет использовать симметричные модели взаимодействия посредством асинхронного обмена сообщениями по одному соединению.

Модуль spring-messaging Spring Framework предоставляет поддержку RSocket-запрашивателей и -отвечателей, как на стороне клиента, так и на стороне сервера. Дополнительные сведения, включая обзор протокола RSocket, см. в разделе RSocket справочника Spring Framework.

4.1. Автоконфигурация стратегий RSocket

Spring Boot автоматически настраивает bean RSocketStrategies, который предоставляет всю необходимую инфраструктуру для кодирования и декодирования полезной нагрузки RSocket. По умолчанию автоконфигурация попытается настроить следующее (в порядке приоритета):

  1. CBOR кодеки с Jackson

  2. JSON кодеки с Jackson

Начальный модуль spring-boot-starter-rsocket предоставляет обе зависимости. Дополнительную информацию о возможностях настройки см. в разделе поддержки Jackson.

Разработчики могут настроить компонент RSocketStrategies, создавая beans, которые реализуют интерфейс RSocketStrategiesCustomizer. Обратите внимание, что их @Order важен, так как он определяет порядок кодеков.

4.2. Автоконфигурация RSocket-сервера

Spring Boot предоставляет автоконфигурацию RSocket-сервера. Необходимые зависимости предоставляются модулем spring-boot-starter-rsocket.

Spring Boot позволяет экспонировать RSocket через WebSocket из сервера WebFlux или запустить независимый RSocket-сервер. Это зависит от типа приложения и его конфигурации.

Для приложения WebFlux (которое имеет тип WebApplicationType.REACTIVE), RSocket-сервер будет подключен к веб-серверу только в том случае, если следующие свойства совпадают:

Свойства
spring.rsocket.server.mapping-path=/rsocket
spring.rsocket.server.transport=websocket
Yaml
spring:
  rsocket:
    server:
      mapping-path: "/rsocket"
      transport: "websocket"
Подключение RSocket к веб-серверу поддерживается только с Reactor Netty, так как сам RSocket построен с использованием этой библиотеки.

В качестве альтернативы, RSocket TCP- или WebSocket-сервер запускается как независимый встроенный сервер. Помимо требований к зависимостям, единственная необходимая конфигурация — определение порта для этого сервера:

Свойства
spring.rsocket.server.port=9898
Yaml
spring:
  rsocket:
    server:
      port: 9898

4.3. Поддержка RSocket в Spring Messaging

Spring Boot автоматически настраивает инфраструктуру Spring Messaging для RSocket.

Это означает, что Spring Boot создаст bean RSocketMessageHandler, который будет обрабатывать запросы RSocket к вашему приложению.

4.4. Вызов RSocket-сервисов с помощью RSocketRequester

После того, как канал RSocket установлен между сервером и клиентом, любая сторона может отправлять или получать запросы друг другу.

В качестве сервера вы можете получить инъекцию экземпляра RSocketRequester в любой обработчик метода RSocket @Controller. В качестве клиента вам сначала нужно настроить и установить RSocket-соединение. Spring Boot автоматически настраивает RSocketRequester.Builder для таких случаев с ожидаемыми кодеками и применяет любой bean RSocketConnectorConfigurer.

Экземпляр RSocketRequester.Builder — это bean-прототип, что означает, что каждый пункт инъекции предоставит вам новый экземпляр. Это сделано намеренно, так как этот билдер имеет состояние, и вы не должны создавать requester'ы с различными настройками, используя один и тот же экземпляр.

Следующий код показывает типичный пример:

Java
import reactor.core.publisher.Mono;

import org.springframework.messaging.rsocket.RSocketRequester;
import org.springframework.stereotype.Service;

@Service
public class MyService {

    private final RSocketRequester rsocketRequester;

    public MyService(RSocketRequester.Builder rsocketRequesterBuilder) {
        this.rsocketRequester = rsocketRequesterBuilder.tcp("example.org", 9898);
    }

    public Mono<User> someRSocketCall(String name) {
        return this.rsocketRequester.route("user").data(name).retrieveMono(User.class);
    }

}
Kotlin
import org.springframework.messaging.rsocket.RSocketRequester
import org.springframework.stereotype.Service
import reactor.core.publisher.Mono

@Service
class MyService(rsocketRequesterBuilder: RSocketRequester.Builder) {

    private val rsocketRequester: RSocketRequester

    init {
        rsocketRequester = rsocketRequesterBuilder.tcp("example.org", 9898)
    }

    fun someRSocketCall(name: String): Mono<User> {
        return rsocketRequester.route("user").data(name).retrieveMono(
            User::class.java
        )
    }

}

5. Spring Integration

Spring Boot предлагает несколько удобств для работы с Spring Integration, включая начальный модуль spring-boot-starter-integration. Spring Integration предоставляет абстракции для обмена сообщениями, а также для других транспортных средств, таких как HTTP, TCP и других. Если Spring Integration доступен в вашем классе, он инициализируется через аннотацию @EnableIntegration.

Логика опроса Spring Integration полагается на автоматически настроенный TaskScheduler. По умолчанию PollerMetadata (опрашивает неограниченное количество сообщений каждую секунду) можно настроить с помощью свойств конфигурации spring.integration.poller.*.

Spring Boot также настраивает некоторые функции, которые запускаются при наличии дополнительных модулей Spring Integration. Если spring-integration-jmx также присутствует в классе, статистика обработки сообщений публикуется через JMX. Если spring-integration-jdbc доступен, схема базы данных по умолчанию может быть создана при запуске, как показано в следующей строке:

Свойства
spring.integration.jdbc.initialize-schema=always
Yaml
spring:
  integration:
    jdbc:
      initialize-schema: "always"

Если spring-integration-rsocket доступен, разработчики могут настроить RSocket-сервер, используя свойства "spring.rsocket.server.*", и позволить ему использовать компоненты IntegrationRSocketEndpoint или RSocketOutboundGateway для обработки входящих сообщений RSocket. Эта инфраструктура может обрабатывать адаптеры каналов Spring Integration RSocket и обработчики @MessageMapping (если "spring.integration.rsocket.server.message-mapping-enabled" настроено).

Spring Boot также может автоматически настроить ClientRSocketConnector с помощью свойств конфигурации:

Свойства
# Connecting to a RSocket server over TCP
spring.integration.rsocket.client.host=example.org
spring.integration.rsocket.client.port=9898
Yaml
# Connecting to a RSocket server over TCP
spring:
  integration:
    rsocket:
      client:
        host: "example.org"
        port: 9898
Свойства
# Connecting to a RSocket Server over WebSocket
spring.integration.rsocket.client.uri=ws://example.org
Yaml
# Connecting to a RSocket Server over WebSocket
spring:
  integration:
    rsocket:
      client:
        uri: "ws://example.org"

См. классы IntegrationAutoConfiguration и IntegrationProperties для получения более подробной информации.

6. WebSockets

Spring Boot предоставляет автоконфигурацию WebSockets для встроенных Tomcat, Jetty и Undertow. Если вы разворачиваете файл war в автономном контейнере, Spring Boot предполагает, что контейнер отвечает за конфигурацию поддержки WebSocket.

Spring Framework предоставляет полную поддержку WebSocket для MVC-веб-приложений, к которой можно легко получить доступ через модуль spring-boot-starter-websocket.

Поддержка WebSocket также доступна для реактивных веб-приложений и требует включения API WebSocket вместе с spring-boot-starter-webflux.

<dependency>
    <groupId>jakarta.websocket</groupId>
    <artifactId>jakarta.websocket-api</artifactId>
</dependency>

7. Что читать дальше

Следующий раздел описывает, как включить возможности ввода-вывода в вашем приложении. Вы можете узнать о кешировании, почте, валидации, клиентах REST и многом другом в этом разделе.

Copyright © 2012-2023 VMware, Inc.
Licensed under the Apache License, Version 2.0.
https://docs.spring.io/spring-boot/docs/3.1.3/reference/html/messaging.html

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API