消息

1. JMS的

jakarta.jms.ConnectionFactoryinterface 提供了一种创建jakarta.jms.Connection用于与 JMS 代理交互。 尽管 Spring 需要一个ConnectionFactory要使用 JMS,您通常不需要自己直接使用它,而是可以依赖更高级别的消息传递抽象。 (有关详细信息,请参阅 Spring Framework 参考文档的相关部分。 Spring Boot 还自动配置了发送和接收消息所需的基础设施。spring-doc.cadn.net.cn

1.1. ActiveMQ “Classic” 支持

ActiveMQ “Classic” 在 Classpath 上可用时, Spring Boot 可以配置ConnectionFactory.spring-doc.cadn.net.cn

如果您使用spring-boot-starter-activemq,提供了连接到 ActiveMQ “Classic” 实例所需的依赖项,以及与 JMS 集成的 Spring 基础架构。

ActiveMQ “Classic” 配置由spring.activemq.*. 默认情况下,ActiveMQ “Classic” 自动配置为使用 TCP 传输,默认情况下连接到tcp://localhost:61616.以下示例显示如何更改默认代理 URL:spring-doc.cadn.net.cn

性能
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-doc.cadn.net.cn

性能
spring.jms.cache.session-cache-size=5
Yaml
spring:
  jms:
    cache:
      session-cache-size: 5

如果您更愿意使用原生池,可以通过向org.messaginghub:pooled-jms并配置JmsPoolConnectionFactory因此,如以下示例所示:spring-doc.cadn.net.cn

性能
spring.activemq.pool.enabled=true
spring.activemq.pool.max-connections=50
Yaml
spring:
  activemq:
    pool:
      enabled: true
      max-connections: 50
ActiveMQProperties了解更多支持的选项。 您还可以注册任意数量的 bean 来实现ActiveMQConnectionFactoryCustomizer以获取更高级的自定义。

默认情况下,ActiveMQ “Classic” 会创建一个目标(如果尚不存在),以便根据提供的名称解析目标。spring-doc.cadn.net.cn

1.2. ActiveMQ Artemis 支持

Spring Boot 可以自动配置ConnectionFactory当它检测到 ActiveMQ Artemis 在 Classpath 上可用时。 如果存在代理,则会自动启动并配置嵌入式代理(除非已明确设置 mode 属性)。 支持的模式包括embedded(明确表示需要一个嵌入式代理,如果代理在 Classpath 上不可用,则应该发生错误)和native(要使用nettytransport 协议)。 配置后者时, Spring Boot 会配置一个ConnectionFactory连接到使用默认设置在本地计算机上运行的代理。spring-doc.cadn.net.cn

如果您使用spring-boot-starter-artemis,提供了连接到现有 ActiveMQ Artemis 实例所需的依赖项,以及与 JMS 集成的 Spring 基础架构。 添加org.apache.activemq:artemis-jakarta-server添加到您的应用程序中,允许您使用嵌入式模式。

ActiveMQ Artemis 配置由spring.artemis.*. 例如,您可以在application.properties:spring-doc.cadn.net.cn

性能
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"

在嵌入代理时,您可以选择是否要启用持久性并列出应可用的目标。 这些可以指定为逗号分隔的列表,以便使用默认选项创建它们,或者您可以定义 bean 类型的org.apache.activemq.artemis.jms.server.config.JMSQueueConfigurationorg.apache.activemq.artemis.jms.server.config.TopicConfiguration,分别用于高级队列和主题配置。spring-doc.cadn.net.cn

默认情况下,CachingConnectionFactory包装本机ConnectionFactory具有合理的设置,您可以通过spring.jms.*:spring-doc.cadn.net.cn

性能
spring.jms.cache.session-cache-size=5
Yaml
spring:
  jms:
    cache:
      session-cache-size: 5

如果您更愿意使用本机池,可以通过添加对org.messaginghub:pooled-jms并配置JmsPoolConnectionFactory因此,如以下示例所示:spring-doc.cadn.net.cn

性能
spring.artemis.pool.enabled=true
spring.artemis.pool.max-connections=50
Yaml
spring:
  artemis:
    pool:
      enabled: true
      max-connections: 50

ArtemisProperties以获取更多支持的选项。spring-doc.cadn.net.cn

不涉及 JNDI 查找,目标是根据其名称解析的,使用name属性或通过配置提供的名称。spring-doc.cadn.net.cn

1.3. 使用 JNDI ConnectionFactory

如果您在应用程序服务器中运行应用程序,则 Spring Boot 会尝试查找 JMSConnectionFactory通过使用 JNDI。 默认情况下,java:/JmsXAjava:/XAConnectionFactorylocation 进行检查。 您可以使用spring.jms.jndi-nameproperty (如果需要指定备用位置),如以下示例所示:spring-doc.cadn.net.cn

性能
spring.jms.jndi-name=java:/MyConnectionFactory
Yaml
spring:
  jms:
    jndi-name: "java:/MyConnectionFactory"

1.4. 发送消息

Spring的JmsTemplate是自动配置的,你可以将其直接自动连接到你自己的 bean 中,如以下示例所示:spring-doc.cadn.net.cn

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可以以类似的方式注射。 如果DestinationResolverMessageConverterbean 时,它会自动关联到自动配置的JmsTemplate.

1.5. 接收消息

当 JMS 基础架构存在时,任何 bean 都可以使用@JmsListener以创建侦听器终端节点。 如果没有JmsListenerContainerFactory,则会自动配置默认 1 个。 如果DestinationResolver一个MessageConverterjakarta.jms.ExceptionListenerbean,它们会自动与默认工厂相关联。spring-doc.cadn.net.cn

默认情况下,默认工厂是事务性的。 如果您在JtaTransactionManager存在,则默认情况下它与侦听器容器相关联。 如果没有,则sessionTransacted标志。 在后一种情况下,您可以通过添加@Transactional在你的侦听器方法(或其委托)上。 这可确保在本地事务完成后确认传入消息。 这还包括发送已在同一 JMS 会话上执行的响应消息。spring-doc.cadn.net.cn

以下组件在someQueue目的地:spring-doc.cadn.net.cn

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的设置与自动配置的设置相同。spring-doc.cadn.net.cn

例如,以下示例公开了另一个工厂,该工厂使用特定的MessageConverter:spring-doc.cadn.net.cn

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-annotated 方法,如下所示:spring-doc.cadn.net.cn

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

高级消息队列协议 (AMQP) 是一种平台中立的有线级协议,适用于面向消息的中间件。 Spring AMQP 项目将核心 Spring 概念应用于基于 AMQP 的消息传递解决方案的开发。 Spring Boot 为通过 RabbitMQ 使用 AMQP 提供了多种便利,包括spring-boot-starter-amqp“Starters”。spring-doc.cadn.net.cn

2.1. RabbitMQ 支持

RabbitMQ 是基于 AMQP 协议的轻量级、可靠、可扩展且可移植的消息代理。 Spring 使用 RabbitMQ 通过 AMQP 协议进行通信。spring-doc.cadn.net.cn

RabbitMQ 配置由spring.rabbitmq.*. 例如,您可以在application.properties:spring-doc.cadn.net.cn

性能
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-doc.cadn.net.cn

性能
spring.rabbitmq.addresses=amqp://admin:secret@localhost
Yaml
spring:
  rabbitmq:
    addresses: "amqp://admin:secret@localhost"
以这种方式指定地址时,hostportproperties 将被忽略。 如果地址使用amqps协议,则会自动启用 SSL 支持。

RabbitProperties,了解更多受支持的基于属性的配置选项。 配置 RabbitMQ 的较低级别详细信息ConnectionFactory,定义一个ConnectionFactoryCustomizer豆。spring-doc.cadn.net.cn

如果ConnectionNameStrategybean 存在于上下文中,它将自动用于命名由自动配置的CachingConnectionFactory.spring-doc.cadn.net.cn

要对RabbitTemplate,请使用RabbitTemplateCustomizer豆。spring-doc.cadn.net.cn

有关更多详细信息,请参阅了解 RabbitMQ 使用的协议 AMQP

2.2. 发送消息

Spring的AmqpTemplateAmqpAdmin是自动配置的,你可以将它们直接自动连接到你自己的 bean 中,如以下示例所示:spring-doc.cadn.net.cn

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可以以类似的方式注射。 如果MessageConverterbean 时,它会自动关联到自动配置的AmqpTemplate.

如有必要,任何org.springframework.amqp.core.Queue定义为 bean 时,它会自动用于在 RabbitMQ 实例上声明相应的队列。spring-doc.cadn.net.cn

要重试作,您可以在AmqpTemplate(例如,如果代理连接丢失):spring-doc.cadn.net.cn

性能
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豆。spring-doc.cadn.net.cn

如果您需要创建更多RabbitTemplate实例,或者如果您想覆盖默认值,Spring Boot 会提供RabbitTemplateConfigurerBean,可用于初始化RabbitTemplate使用与 auto-configuration 使用的工厂相同的设置。spring-doc.cadn.net.cn

2.3. 向流发送消息

要向特定流发送消息,请指定流的名称,如以下示例所示:spring-doc.cadn.net.cn

性能
spring.rabbitmq.stream.name=my-stream
Yaml
spring:
  rabbitmq:
    stream:
      name: "my-stream"

如果MessageConverter,StreamMessageConverterProducerCustomizerbean 时,它会自动关联到自动配置的RabbitStreamTemplate.spring-doc.cadn.net.cn

如果您需要创建更多RabbitStreamTemplate实例,或者如果您想覆盖默认值,Spring Boot 会提供RabbitStreamTemplateConfigurerBean,可用于初始化RabbitStreamTemplate使用与 auto-configuration 使用的工厂相同的设置。spring-doc.cadn.net.cn

2.4. 接收消息

当 Rabbit 基础结构存在时,任何 bean 都可以使用@RabbitListener以创建侦听器终端节点。 如果没有RabbitListenerContainerFactory,则默认的SimpleRabbitListenerContainerFactory会自动配置,您可以使用spring.rabbitmq.listener.type财产。 如果MessageConverterMessageRecovererbean,它会自动与默认工厂相关联。spring-doc.cadn.net.cn

以下示例组件在someQueue队列:spring-doc.cadn.net.cn

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?) {
        // ...
    }

}
的 Javadoc@EnableRabbit了解更多详情。

如果您需要创建更多RabbitListenerContainerFactory实例,或者如果您想覆盖默认值,Spring Boot 会提供SimpleRabbitListenerContainerFactoryConfigurer以及DirectRabbitListenerContainerFactoryConfigurer,可用于初始化SimpleRabbitListenerContainerFactory以及DirectRabbitListenerContainerFactory使用与 auto-configuration 使用的工厂相同的设置。spring-doc.cadn.net.cn

选择哪种容器类型并不重要。 这两个 bean 由自动配置公开。

例如,下面的配置类公开了另一个使用特定MessageConverter:spring-doc.cadn.net.cn

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-annotated 方法,如下所示:spring-doc.cadn.net.cn

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豆。spring-doc.cadn.net.cn

默认情况下,如果禁用重试并且侦听器引发异常,则会无限期地重试投放。 您可以通过两种方式修改此行为:将defaultRequeueRejectedproperty 设置为false,以便尝试零重新投递,或者抛出AmqpRejectAndDontRequeueException来表示消息应该被拒绝。 后者是启用重试并达到最大投放尝试次数时使用的机制。

3. Apache Kafka 支持

Apache Kafka 通过提供spring-kafka项目。spring-doc.cadn.net.cn

Kafka 配置由spring.kafka.*. 例如,您可以在application.properties:spring-doc.cadn.net.cn

性能
spring.kafka.bootstrap-servers=localhost:9092
spring.kafka.consumer.group-id=myGroup
Yaml
spring:
  kafka:
    bootstrap-servers: "localhost:9092"
    consumer:
      group-id: "myGroup"
要在启动时创建主题,请添加NewTopic. 如果主题已存在,则忽略该 Bean。

KafkaProperties以获取更多支持的选项。spring-doc.cadn.net.cn

3.1. 发送消息

Spring的KafkaTemplate是自动配置的,你可以直接在自己的 bean 中自动装配它,如以下示例所示:spring-doc.cadn.net.cn

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将自动配置。 此外,如果RecordMessageConverterbean 时,它会自动关联到自动配置的KafkaTemplate.

3.2. 接收消息

当存在 Apache Kafka 基础架构时,任何 bean 都可以使用@KafkaListener以创建侦听器终端节点。 如果没有KafkaListenerContainerFactory已定义,则会自动使用spring.kafka.listener.*.spring-doc.cadn.net.cn

以下组件在someTopic主题:spring-doc.cadn.net.cn

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?) {
        // ...
    }

}

如果KafkaTransactionManagerbean,它会自动关联到容器工厂。 同样,如果RecordFilterStrategy,CommonErrorHandler,AfterRollbackProcessorConsumerAwareRebalanceListenerbean,它会自动关联到默认工厂。spring-doc.cadn.net.cn

根据侦听器类型,RecordMessageConverterBatchMessageConverterbean 与默认工厂相关联。 如果只有一个RecordMessageConverterbean 的 bean 被打包在BatchMessageConverter.spring-doc.cadn.net.cn

自定义ChainedKafkaTransactionManager必须标记@Primary因为它通常引用 auto-configuredKafkaTransactionManager豆。

3.3. Kafka 流

Spring for Apache Kafka 提供了一个工厂 Bean 来创建StreamsBuilder对象并管理其流的生命周期。 Spring Boot 会自动配置所需的KafkaStreamsConfigurationbean 只要kafka-streams位于 Classpath 上,并且 Kafka Streams 由@EnableKafkaStreams注解。spring-doc.cadn.net.cn

启用 Kafka Streams 意味着必须设置应用程序 ID 和引导服务器。 前者可以使用spring.kafka.streams.application-id,默认为spring.application.name如果未设置。 后者可以全局设置,也可以仅针对流专门覆盖。spring-doc.cadn.net.cn

使用专用属性可以使用多个其他属性;其他任意 Kafka 属性可以使用spring.kafka.streams.propertiesNamespace。 有关更多信息,另请参阅其他 Kafka 属性spring-doc.cadn.net.cn

要使用工厂 Bean,请将 wireStreamsBuilder放入您的@Bean如以下示例所示:spring-doc.cadn.net.cn

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财产。spring-doc.cadn.net.cn

3.4. 其他 Kafka 属性

自动配置支持的属性显示在附录的 “集成属性” 部分中。 请注意,在大多数情况下,这些属性(带连字符或 camelCase)直接映射到 Apache Kafka 点分隔属性。 有关详细信息,请参阅 Apache Kafka 文档。spring-doc.cadn.net.cn

不包含客户端类型 (producer,consumer,adminstreams) 被视为通用的,适用于所有客户端。 如果需要,可以覆盖一个或多个客户端类型的大多数公共属性。spring-doc.cadn.net.cn

Apache Kafka 指定重要性为 HIGH、MEDIUM 或 LOW 的属性。 Spring Boot 自动配置支持所有 HIGH 重要性属性,一些选定的 MEDIUM 和 LOW 属性,以及任何没有默认值的属性。spring-doc.cadn.net.cn

只有 Kafka 支持的属性子集可以直接通过KafkaProperties类。 如果要使用不直接支持的其他属性配置各个客户端类型,请使用以下属性:spring-doc.cadn.net.cn

性能
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.oneKafka 属性设置为first(适用于创建者、使用者、管理员和流),则prop.twoadmin 属性设置为secondprop.threeconsumer 属性设置为thirdprop.fourproducer 属性设置为fourthprop.fivestreams 属性设置为fifth.spring-doc.cadn.net.cn

您还可以配置 Spring KafkaJsonDeserializer如下:spring-doc.cadn.net.cn

性能
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在 Headers 中发送类型信息的默认行为:spring-doc.cadn.net.cn

性能
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 for Apache Kafka 提供了一种使用嵌入式 Apache Kafka 代理测试项目的便捷方法。 要使用此功能,请使用@EmbeddedKafkaspring-kafka-test模块。 有关更多信息,请参阅 Spring for Apache Kafka 参考手册spring-doc.cadn.net.cn

要使 Spring Boot 自动配置与上述嵌入式 Apache Kafka 代理一起工作,您需要重新映射嵌入式代理地址的系统属性(由EmbeddedKafkaBroker) 添加到 Apache Kafka 的 Spring Boot 配置属性中。 有几种方法可以做到这一点:spring-doc.cadn.net.cn

  • 提供一个系统属性,用于将嵌入式代理地址映射到spring.kafka.bootstrap-servers在 test 类中:spring-doc.cadn.net.cn

Java
static {
    System.setProperty(EmbeddedKafkaBroker.BROKER_LIST_PROPERTY, "spring.kafka.bootstrap-servers");
}
Kotlin
init {
    System.setProperty(EmbeddedKafkaBroker.BROKER_LIST_PROPERTY, "spring.kafka.bootstrap-servers")
}
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. RS锁

RSocket 是一种用于字节流传输的二进制协议。 它通过单个连接上的异步消息传递来实现对称交互模型。spring-doc.cadn.net.cn

spring-messagingModule 在客户端和服务器端都为 RSocket 请求者和响应者提供支持。 有关更多详细信息,包括 RSocket 协议的概述,请参见 Spring Framework 参考的 RSocket 部分spring-doc.cadn.net.cn

4.1. RSocket 策略自动配置

Spring Boot 会自动配置RSocketStrategiesbean 提供编码和解码 RSocket 有效负载所需的所有基础结构。 默认情况下,自动配置将尝试配置以下内容(按顺序):spring-doc.cadn.net.cn

  1. 使用 Jackson 的 CBOR 编解码器spring-doc.cadn.net.cn

  2. 使用 Jackson 的 JSON 编解码器spring-doc.cadn.net.cn

spring-boot-starter-rsocketstarter 提供了这两个依赖项。 请参阅 Jackson 支持部分以了解有关自定义可能性的更多信息。spring-doc.cadn.net.cn

开发人员可以自定义RSocketStrategies组件,方法是创建实现RSocketStrategiesCustomizer接口。 请注意,他们的@Order很重要,因为它决定了编解码器的顺序。spring-doc.cadn.net.cn

4.2. RSocket 服务器自动配置

Spring Boot 提供 RSocket 服务器自动配置。 所需的依赖项由spring-boot-starter-rsocket.spring-doc.cadn.net.cn

Spring Boot 允许从 WebFlux 服务器通过 WebSocket 公开 RSocket,或者建立独立的 RSocket 服务器。 这取决于应用程序的类型及其配置。spring-doc.cadn.net.cn

对于 WebFlux 应用程序(类型为WebApplicationType.REACTIVE),则仅当以下属性匹配时,RSocket 服务器才会插入 Web 服务器:spring-doc.cadn.net.cn

性能
spring.rsocket.server.mapping-path=/rsocket
spring.rsocket.server.transport=websocket
Yaml
spring:
  rsocket:
    server:
      mapping-path: "/rsocket"
      transport: "websocket"
只有 Reactor Netty 才支持将 RSocket 插入 Web 服务器,因为 RSocket 本身就是使用该库构建的。

或者,RSocket TCP 或 websocket 服务器作为独立的嵌入式服务器启动。 除了依赖项要求之外,唯一需要的配置是为该服务器定义一个端口:spring-doc.cadn.net.cn

性能
spring.rsocket.server.port=9898
Yaml
spring:
  rsocket:
    server:
      port: 9898

4.3. Spring Messaging RSocket 支持

Spring Boot 将为 RSocket 自动配置 Spring Messaging 基础设施。spring-doc.cadn.net.cn

这意味着 Spring Boot 将创建一个RSocketMessageHandlerbean 的 bean 来处理对应用程序的 RSocket 请求。spring-doc.cadn.net.cn

4.4. 使用 RSocketRequester 调用 RSocket 服务

一旦RSocketchannel 建立在 Server 和 Client 端之间,任何一方都可以向另一方发送或接收请求。spring-doc.cadn.net.cn

作为服务器,您可以注入RSocketRequester实例在 RSocket 的任何处理程序方法上@Controller. 作为客户端,您需要先配置并建立 RSocket 连接。 Spring Boot 会自动配置RSocketRequester.Builder对于此类情况,使用预期的编解码器并应用任何RSocketConnectorConfigurer豆。spring-doc.cadn.net.cn

RSocketRequester.Builderinstance 是一个原型 bean,这意味着每个注入点都会为您提供一个新的 instance 。 这是有意为之的,因为此构建器是有状态的,您不应使用同一实例创建具有不同设置的请求者。spring-doc.cadn.net.cn

以下代码显示了一个典型的示例:spring-doc.cadn.net.cn

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 集成

Spring Boot 为使用 Spring 集成提供了多种便利,包括spring-boot-starter-integration“Starters”。 Spring 集成通过消息传递以及其他传输(如 HTTP、TCP 等)提供抽象。 如果 Spring 集成在您的类路径上可用,则通过@EnableIntegration注解。spring-doc.cadn.net.cn

Spring 集成轮询逻辑依赖于在自动配置的TaskScheduler. 默认的PollerMetadata(轮询每秒无限数量的消息)可以使用spring.integration.poller.*configuration 属性。spring-doc.cadn.net.cn

Spring Boot 还配置了一些由其他 Spring 集成模块的存在触发的功能。 如果spring-integration-jmx也在 Classpath 上,则消息处理统计信息将通过 JMX 发布。 如果spring-integration-jdbc可用,则可以在启动时创建默认数据库架构,如以下行所示:spring-doc.cadn.net.cn

性能
spring.integration.jdbc.initialize-schema=always
Yaml
spring:
  integration:
    jdbc:
      initialize-schema: "always"

如果spring-integration-rsocket可用,开发人员可以使用"spring.rsocket.server.*"属性,并让它使用IntegrationRSocketEndpointRSocketOutboundGateway组件来处理传入的 RSocket 消息。 此基础结构可以处理 Spring 集成 RSocket 通道适配器和@MessageMappinghandlers(给定"spring.integration.rsocket.server.message-mapping-enabled"已配置)。spring-doc.cadn.net.cn

Spring Boot 还可以自动配置ClientRSocketConnector使用配置属性:spring-doc.cadn.net.cn

性能
# 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"

6. 网络套接字

Spring Boot 为嵌入式 Tomcat、Jetty 和 Undertow 提供 WebSockets 自动配置。 如果将 war 文件部署到独立容器,则 Spring Boot 假定该容器负责其 WebSocket 支持的配置。spring-doc.cadn.net.cn

Spring Framework 为 MVC Web 应用程序提供了丰富的 WebSocket 支持,可以通过spring-boot-starter-websocket模块。spring-doc.cadn.net.cn

WebSocket 支持也可用于反应式 Web 应用程序,并且需要包含 WebSocket API 以及spring-boot-starter-webflux:spring-doc.cadn.net.cn

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

7. 下一步要读什么

下一节介绍如何在应用程序中启用 IO 功能。 您可以在本节中阅读有关缓存邮件验证REST 客户端等的信息。spring-doc.cadn.net.cn