Сообщения
Фреймворк 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 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 spring:
jms:
cache:
session-cache-size: 5 Если вы предпочитаете использовать родной пуллинг, вы можете сделать это, добавив зависимость от org.messaginghub:pooled-jms и настроив JmsPoolConnectionFactory соответственно, как показано в следующем примере:
spring.activemq.pool.enabled=true
spring.activemq.pool.max-connections=50 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 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 spring:
jms:
cache:
session-cache-size: 5 Если вы предпочитаете использовать родной пуллинг, вы можете сделать это, добавив зависимость от org.messaginghub:pooled-jms и настроив JmsPoolConnectionFactory соответственно, как показано в следующем примере:
spring.artemis.pool.enabled=true
spring.artemis.pool.max-connections=50 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 spring:
jms:
jndi-name: "java:/MyConnectionFactory" 1.4. Отправка сообщения
Автоконфигурируется JmsTemplate Spring, и вы можете напрямую автоматизировать его в свои собственные бины, как показано в следующем примере:
@Component
public class MyBean {
private final JmsTemplate jmsTemplate;
public MyBean(JmsTemplate jmsTemplate) {
this.jmsTemplate = jmsTemplate;
}
}
@Component
class MyBean(private val jmsTemplate: JmsTemplate) {
}
JmsMessagingTemplate можно вводить аналогичным образом. Если определен бины DestinationResolver или MessageConverter, он автоматически связывается с автоматически настроенным JmsTemplate. |
1.5. Прием сообщения
При наличии инфраструктуры JMS любой бин может быть аннотирован с помощью @JmsListener для создания точки входа слушателя. Если JmsListenerContainerFactory не определено, то оно автоматически настраивается по умолчанию. Если бины DestinationResolver, MessageConverter или jakarta.jms.ExceptionListener определены, они автоматически ассоциируются с фабрикой по умолчанию.
По умолчанию, фабрика по умолчанию транзакционная. Если вы работаете в инфраструктуре, где присутствует JtaTransactionManager, она по умолчанию связывается с контейнером слушателя. В противном случае, включается флаг sessionTransacted. В последнем случае, вы можете связать транзакцию вашего локального хранилища данных с обработкой входящего сообщения, добавив @Transactional в метод слушателя (или делегат). Это гарантирует подтверждение входящего сообщения после завершения локальной транзакции. Это также включает отправку ответных сообщений, которые были выполнены в той же сессии JMS.
Следующий компонент создает точку входа слушателя на целевом пункте someQueue:
@Component
public class MyBean {
@JmsListener(destination = "someQueue")
public void processMessage(String content) {
// ...
}
}
@Component
class MyBean {
@JmsListener(destination = "someQueue")
fun processMessage(content: String?) {
// ...
}
}
Смотрите Javadoc @EnableJms для получения более подробной информации. |
Если вам нужно создать больше экземпляров JmsListenerContainerFactory или если вы хотите переопределить значение по умолчанию, Spring Boot предоставляет DefaultJmsListenerContainerFactoryConfigurer, который вы можете использовать для инициализации DefaultJmsListenerContainerFactory с теми же настройками, что и автоматически настроенный экземпляр.
Например, следующий пример экспонирует другую фабрику, использующую определенный MessageConverter:
@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 ...
}
}
@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-аннотированном методе следующим образом:
@Component
public class MyBean {
@JmsListener(destination = "someQueue", containerFactory = "myFactory")
public void processMessage(String content) {
// ...
}
}
@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 spring:
rabbitmq:
host: "localhost"
port: 5672
username: "admin"
password: "secret" Альтернативно, вы можете настроить то же соединение с помощью атрибута addresses:
spring.rabbitmq.addresses=amqp://admin:secret@localhost 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 автоматически настраиваются, и вы можете напрямую их внедрить в собственные бины, как показано в следующем примере:
@Component
public class MyBean {
private final AmqpAdmin amqpAdmin;
private final AmqpTemplate amqpTemplate;
public MyBean(AmqpAdmin amqpAdmin, AmqpTemplate amqpTemplate) {
this.amqpAdmin = amqpAdmin;
this.amqpTemplate = amqpTemplate;
}
}
@Component
class MyBean(private val amqpAdmin: AmqpAdmin, private val amqpTemplate: AmqpTemplate) {
}
RabbitMessagingTemplate может быть внедрён аналогичным образом. Если бин MessageConverter определён, он автоматически связан с автоматически настроенным AmqpTemplate. |
Если необходимо, любой org.springframework.amqp.core.Queue бин, определённый как бин, автоматически используется для объявления соответствующего очереди на экземпляре RabbitMQ.
Чтобы повторить операции, можно включить повторы на AmqpTemplate (например, в случае потери соединения с брокером):
spring.rabbitmq.template.retry.enabled=true
spring.rabbitmq.template.retry.initial-interval=2s spring:
rabbitmq:
template:
retry:
enabled: true
initial-interval: "2s" Повторные попытки отключены по умолчанию. Вы также можете настроить RetryTemplate программно, объявив бин RabbitRetryTemplateCustomizer.
Если вам нужно создать больше RabbitTemplate экземпляров или если вы хотите переопределить значения по умолчанию, Spring Boot предоставляет бин RabbitTemplateConfigurer, который вы можете использовать для инициализации RabbitTemplate с теми же настройками, что и фабрики, используемые автонастройкой.
2.3. Отправка сообщения в поток
Чтобы отправить сообщение в определённый поток, укажите имя потока, как показано в следующем примере:
spring.rabbitmq.stream.name=my-stream 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:
@Component
public class MyBean {
@RabbitListener(queues = "someQueue")
public void processMessage(String content) {
// ...
}
}
@Component
class MyBean {
@RabbitListener(queues = ["someQueue"])
fun processMessage(content: String?) {
// ...
}
}
См. документацию @EnableRabbit для получения дополнительных сведений. |
Если вам нужно создать больше RabbitListenerContainerFactory экземпляров или если вы хотите переопределить значения по умолчанию, Spring Boot предоставляет бин SimpleRabbitListenerContainerFactoryConfigurer и DirectRabbitListenerContainerFactoryConfigurer, которые вы можете использовать для инициализации SimpleRabbitListenerContainerFactory и DirectRabbitListenerContainerFactory с теми же настройками, что и фабрики, используемые автонастройкой.
| Тип контейнера, который вы выбрали, не имеет значения. Эти два бинa доступны через автонастройку. |
Например, следующий класс конфигурации предоставляет другую фабрику, использующую определённый MessageConverter:
@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 ...
}
}
@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:
@Component
public class MyBean {
@RabbitListener(queues = "someQueue", containerFactory = "myFactory")
public void processMessage(String content) {
// ...
}
}
@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 spring:
kafka:
bootstrap-servers: "localhost:9092"
consumer:
group-id: "myGroup" Для создания темы при запуске добавьте bean типа NewTopic. Если тема уже существует, bean игнорируется. |
См. KafkaProperties для получения более подробной информации о поддерживаемых вариантах.
3.1. Отправка сообщения
Spring’s KafkaTemplate автоматически настраивается, и вы можете напрямую внедрить его в свои собственные компоненты, как показано в следующем примере:
@Component
public class MyBean {
private final KafkaTemplate<String, String> kafkaTemplate;
public MyBean(KafkaTemplate<String, String> kafkaTemplate) {
this.kafkaTemplate = kafkaTemplate;
}
}
@Component
class MyBean(private val kafkaTemplate: KafkaTemplate<String, String>) {
}
Если свойство spring.kafka.producer.transaction-id-prefix определено, автоматически настраивается KafkaTransactionManager. Кроме того, если bean RecordMessageConverter определён, он автоматически связывается с автоматически настроенным KafkaTemplate . |
3.2. Приём сообщения
При наличии инфраструктуры Apache Kafka любой компонент может быть аннотирован с помощью @KafkaListener для создания точки входа для прослушивания. Если KafkaListenerContainerFactory не определено, автоматически настраивается значение по умолчанию с ключами, определёнными в spring.kafka.listener.*.
Следующий компонент создаёт точку входа для прослушивания на теме someTopic:
@Component
public class MyBean {
@KafkaListener(topics = "someTopic")
public void processMessage(String content) {
// ...
}
}
@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, как показано в следующем примере:
@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());
}
}
@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 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 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 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в классе теста:
static {
System.setProperty(EmbeddedKafkaBroker.BROKER_LIST_PROPERTY, "spring.kafka.bootstrap-servers");
}
init {
System.setProperty(EmbeddedKafkaBroker.BROKER_LIST_PROPERTY, "spring.kafka.bootstrap-servers")
}
-
Настройте имя свойства в аннотации
@EmbeddedKafka:
@SpringBootTest
@EmbeddedKafka(topics = "someTopic", bootstrapServersProperty = "spring.kafka.bootstrap-servers")
class MyTest {
// ...
}
@SpringBootTest
@EmbeddedKafka(topics = ["someTopic"], bootstrapServersProperty = "spring.kafka.bootstrap-servers")
class MyTest {
// ...
}
-
Используйте плейсхолдер в свойствах конфигурации:
spring.kafka.bootstrap-servers=${spring.embedded.kafka.brokers} 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. По умолчанию автоконфигурация попытается настроить следующее (в порядке приоритета):
-
CBOR кодеки с Jackson
-
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 spring:
rsocket:
server:
mapping-path: "/rsocket"
transport: "websocket" | Подключение RSocket к веб-серверу поддерживается только с Reactor Netty, так как сам RSocket построен с использованием этой библиотеки. |
В качестве альтернативы, RSocket TCP- или WebSocket-сервер запускается как независимый встроенный сервер. Помимо требований к зависимостям, единственная необходимая конфигурация — определение порта для этого сервера:
spring.rsocket.server.port=9898 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'ы с различными настройками, используя один и тот же экземпляр.
Следующий код показывает типичный пример:
@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);
}
}
@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 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 # 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 # 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