【问题标题】:EnableBinding, Output, Input deprecated Since version of 3.1 of Spring Cloud Stream从 Spring Cloud Stream 3.1 版本开始不推荐使用 EnableBinding、Output、Input
【发布时间】:2021-05-04 18:21:28
【问题描述】:

从 3.1 版开始,用于处理队列的主要 API 已被弃用。 在课堂评论中它说:

已弃用 从 3.1 开始支持函数式编程模型

我在网上搜索了很多解决方案,但没有找到关于我应该如何迁移的可靠的 E2E 解释。

寻找以下例子:

  1. 从队列中读取
  2. 写入队列

如果有几种方法可以做到这一点(正如我在网络上看到的那样),我会很高兴得到解释以及每个选项的典型用例。

【问题讨论】:

    标签: spring spring-cloud-stream


    【解决方案1】:
    1. 我假设您已经熟悉主要概念,并将专注于迁移。
    2. 我在演示代码中使用了 kotlin,以减少冗长

    首先,一些可能有帮助的参考资料:

    • 这是最初的相关文档:link
    • 这是对新功能格式中命名方案的解释:link
    • 这是一个更详细的解释,带有一些更高级的场景:link

    TL;DR

    spring 现在不再使用基于注解的配置,而是使用检测到的 Consumer/Function/Supplier bean 为您定义流。

    输入/消费者

    而在您的代码看起来像这样之前:

    interface BindableGradesChannel {
        @Input
        fun gradesChannel(): SubscribableChannel
    
        companion object {
            const val INPUT = "gradesChannel"
        }
    }
    

    用法类似:

    @Service
    @EnableBinding(BindableGradesChannel::class)
    class GradesListener {
        private val log = LoggerFactory.getLogger(GradesListener::class.java)
        
        @StreamListener(BindableScoresChannel.INPUT)
        fun listen(grade: Grade) {
            log.info("Received $grade")
            // do something
        }
    }
    

    现在整个定义都无关紧要了,可以这样做:

    @Service
    class GradesListener {
        private val log = LoggerFactory.getLogger(GradesListener::class.java)
    
        @Bean
        fun gradesChannel(): Consumer<Grade> {
            return Consumer { listen(grade = it) }
        }
        
        fun listen(grade: Grade) {
            log.info("Received $grade")
            // do something
        }
    }
    

    注意Consumer bean 如何替换@StreamListener@Input

    关于配置,如果之前为了配置你有一个 application.yml 看起来像这样:

    spring:
      cloud:
        stream:
          bindings:
            gradesChannel:
              destination: GradesExchange
              group: grades-updates
              consumer:
                concurrency: 10
                max-attempts: 3
    

    现在应该是这样的:

    spring:
      cloud:
        stream:
          bindings:
            gradesChannel-in-0:
              destination: GradesExchange
              group: grades-updates
              consumer:
                concurrency: 10
                max-attempts: 3
    

    注意gradesChannel 是如何被gradesChannel-in-0 替换的 - 要了解完整的命名约定,请参阅顶部的命名约定链接。

    一些细节:

    1. 如果您的应用程序中有多个此类 bean,则需要定义 spring.cloud.function.definition 属性。
    2. 您可以选择自定义频道名称,因此如果您想继续使用 gradesChannel,可以设置 spring.cloud.stream.function.bindings.gradesChannel-in-0=gradesChannel 并在配置中的任何位置使用 gradesChannel

    输出/供应商

    这里的概念类似,你替换配置和代码看起来像这样:

    interface BindableStudentsChannel {
        @Output
        fun studentsChannel(): MessageChannel
    }
    

    @Service
    @EnableBinding(BindableStudentsChannel::class)
    class StudentsQueueWriter(private val studentsChannel: BindableStudentsChannel) {
        fun publish(message: Message<Student>) {
            studentsChannel.studentsChannel().send(message)
        }
    }
    

    现在可以替换为:

    @Service
    class StudentsQueueWriter {
        @Bean
        fun studentsChannel(): Supplier<Student> {
            return Supplier { Student("Adam") }
        }
    }
    

    如您所见,我们有一个主要区别 - 何时调用以及由谁调用?

    以前我们可以手动触发它,但现在它是由 spring 触发的,每秒触发一次(默认情况下)。这对于用例来说很好,例如当您需要每秒发布传感器数据时,但是当您想要在事件上发送消息时这并不好。除了出于任何原因使用 Function 之外,spring 还提供了 2 种选择:

    StreamBridge - link

    使用StreamBridge 即可。像这样明确定义目标:

    @Service
    class StudentsQueueWriter(private val streamBridge: StreamBridge) {
        fun publish(message: Message<Student>) {
            streamBridge.send("studentsChannel-out-0", message)
        }
    }
    

    这样您就不会将目标通道定义为 bean,但您仍然可以发送消息。缺点是你的类中有一些明确的配置。

    反应堆 API - link

    另一种方法是使用某种反应机制,例如EmitterProcessor,并返回它。使用它,您的代码将类似于:

    @Service
    class StudentsQueueWriter {
        val students: EmitterProcessor<Student> = EmitterProcessor.create()
        @Bean
        fun studentsChannel(): Supplier<Flux<Student>> {
            return Supplier { students }
        }
    }
    

    用法可能类似于:

    class MyClass(val studentsQueueWriter: StudentsQueueWriter) {
        fun newStudent() {
            studentsQueueWriter.students.onNext(Student("Adam"))
        }
    }
    

    【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2019-06-28
    • 1970-01-01
    • 2021-12-18
    • 1970-01-01
    • 2021-11-15
    • 1970-01-01
    • 1970-01-01
    • 2010-09-17
    相关资源
    最近更新 更多