You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spring Cloud Stream v2至v4:EnableBinding与StreamListener迁移咨询

迁移Spring Cloud Stream v2到v4:替代EnableBinding和StreamListener

Spring Cloud Stream从v3版本开始转向函数式编程模型,v4彻底移除了EnableBinding、StreamListener等旧注解。以下是针对你的代码的具体迁移方案:

1. 核心修改步骤

(1)移除旧注解与调整类结构

删除@EnableBinding(Sink.class)注解,将原来的StreamListener方法重构为函数式Bean,并改用构造注入替代字段注入(符合Spring最佳实践)。

(2)重构消息消费逻辑

将原来的handleEnergyProductionMessage方法转换为Consumer<EnergyProductionMessage>类型的Bean,Spring Cloud Stream会自动将其绑定到输入通道。

(3)调整错误消息处理

原来的errorChannel监听逻辑,改用@EventListener监听ErrorMessage事件来实现。

2. 迁移后的完整代码

import java.time.Clock;
import java.time.Duration;
import java.time.Instant;

import org.springframework.context.event.EventListener;
import org.springframework.messaging.support.ErrorMessage;
import org.springframework.stereotype.Service;

import com.enterprise.production.model.exceptions.ProductionStoreException;
import com.enterprise.production.model.message.EnergyProductionMessage;
import com.enterprise.production.stream.service.ProductionStoreService;

import lombok.extern.slf4j.Slf4j;

@Slf4j
@Service
public class ProductionMessageConsumer {

    private final ProductionStoreService productionService;
    private final Clock clock;

    // 构造注入替代@Autowired字段注入
    public ProductionMessageConsumer(ProductionStoreService productionService, Clock clock) {
        this.productionService = productionService;
        this.clock = clock;
    }

    // 定义消费消息的函数式Bean,对应原Sink.INPUT的处理逻辑
    @Bean
    public Consumer<EnergyProductionMessage> handleEnergyProductionMessage() {
        return energyProductionMessage -> {
            try {
                Instant start = clock.instant();

                log.debug("Processing energy productions message with original interval: {}|{}|{}", 
                    energyProductionMessage.getTenantId(), 
                    energyProductionMessage.getUsername(),
                    energyProductionMessage.getDeviceId());
                log.info("Processing {} energy productions ", energyProductionMessage.getSolarEnergies().size());

                productionService.saveProductions(energyProductionMessage);
                log.debug("Ending energy productions message with original interval: {}|{}|{}: ended in {}ms", 
                    energyProductionMessage.getTenantId(),
                    energyProductionMessage.getUsername(), 
                    energyProductionMessage.getDeviceId(), 
                    Duration.between(start, clock.instant()).toMillis());
                
                Instant startNormalization = clock.instant();
                log.debug("Processing energy productions message with normalization 30m: {}|{}|{}", 
                    energyProductionMessage.getTenantId(), 
                    energyProductionMessage.getUsername(),
                    energyProductionMessage.getDeviceId());
                productionService.saveProductions30m(energyProductionMessage);
                log.debug("Ending energy productions message with normalization 30m: {}|{}|{}: ended in {}ms", 
                    energyProductionMessage.getTenantId(),
                    energyProductionMessage.getUsername(), 
                    energyProductionMessage.getDeviceId(), 
                    Duration.between(startNormalization, clock.instant()).toMillis());
            } catch (ProductionStoreException e) {
                // 将业务异常转为RuntimeException,让框架处理重试/死信(需配合配置)
                throw new RuntimeException("Failed to process production message", e);
            }
        };
    }

    // 替代原errorChannel的监听逻辑
    @EventListener
    public void handleError(ErrorMessage errorMessage) {
        log.error("Fail to read message with error '{}'", errorMessage.getPayload());
        // 如需获取原始失败消息,可调用errorMessage.getOriginalMessage()
    }
}

3. 必要的配置调整

在application.yml(或application.properties)中配置函数定义与输入通道绑定,替换原来的Sink相关配置:

spring:
  cloud:
    stream:
      # 绑定函数的输入通道,命名规则:函数名-in-0
      bindings:
        handleEnergyProductionMessage-in-0:
          destination: your-input-topic-or-queue  # 替换为你原来Sink.INPUT对应的目标名称
          group: your-consumer-group              # 设置消费者组,避免重复消费
      # 指定要激活的函数Bean名称
      function:
        definition: handleEnergyProductionMessage

4. 关键说明

  • 函数式模型核心:每个消费逻辑对应一个Consumer类型的Bean,生产者对应Supplier,双向处理对应Function,框架自动完成通道绑定。
  • 通道命名规则:函数的输入通道命名为{函数名}-in-0,输出通道为{函数名}-out-0,多输入/输出可按序号递增。
  • 异常处理:业务异常需转为RuntimeException抛出,配合Spring Cloud Stream的重试、死信队列配置,实现消息容错。
  • Sink接口废弃:v4不再需要自定义或使用内置的Sink/Source接口,完全通过函数Bean和配置来声明通道。

内容的提问来源于stack exchange,提问作者ThrowsError

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.04 16:55:55