Spring Boot 3中替代@EnableBinding实现Kafka消息收发方案咨询
Spring Boot 3 适配Spring Cloud Stream Kafka 收发消息方案(基于《Spring Microservices In Action (2021)》示例改造)
我正在阅读《Spring Microservices In Action (2021)》巩固微服务知识,但书中基于Spring Boot 2的Kafka消息收发示例无法在Spring Boot 3中运行,原项目对应章节代码位于https://github.com/ihuaylupo/manning-smia/tree/master/chapter10。以下是针对两个核心示例的适配改造方案:
示例1:消息发送(organization-service)
原实现问题
原代码基于Spring Cloud Stream的绑定注解模型(依赖Source接口、@Autowired注入绑定器),但Spring Boot 3对应的Spring Cloud Stream 4.x已废弃该模型,全面切换为函数式编程模型,同时移除了ZK节点配置(Kafka Broker管理不再依赖ZK)。
Spring Boot 3 适配方案
1. 配置文件调整
替换原output绑定配置,适配函数式模型,并移除ZK相关配置:
# 函数式输出绑定:函数名sendOrgChange对应绑定后缀-out-0 spring.cloud.stream.bindings.sendOrgChange-out-0.destination=orgChangeTopic spring.cloud.stream.bindings.sendOrgChange-out-0.content-type=application/json # Kafka Broker地址(与Docker Compose中的网络别名一致) spring.cloud.stream.kafka.binder.brokers=kafka # 开启自动创建主题(首次发送消息时自动生成orgChangeTopic) spring.cloud.stream.kafka.binder.auto-create-topics=true
2. 代码改造
移除Source注入逻辑,改用两种主流实现方式:
方式一:函数式Supplier(标准模型)
import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.context.annotation.Bean; import org.springframework.messaging.Message; import org.springframework.messaging.support.MessageBuilder; import org.springframework.stereotype.Component; import java.util.function.Supplier; @Component public class OrganizationEventPublisher { private static final Logger logger = LoggerFactory.getLogger(OrganizationEventPublisher.class); // 定义消息发送的Supplier函数,函数名需与配置中的sendOrgChange完全匹配 @Bean public Supplier<Message<OrganizationChangeModel>> sendOrgChange() { return () -> MessageBuilder.withPayload(new OrganizationChangeModel()).build(); } // 封装业务发送逻辑,触发Supplier发送消息 public void publishOrganizationChange(String action, String organizationId) { logger.debug("Sending Kafka message {} for Organization Id: {}", action, organizationId); OrganizationChangeModel change = new OrganizationChangeModel( OrganizationChangeModel.class.getTypeName(), action, organizationId, UserContext.getCorrelationId()); // 构造消息并发送 Message<OrganizationChangeModel> message = MessageBuilder.withPayload(change).build(); sendOrgChange().get(); } }
方式二:StreamBridge(推荐,更灵活)
无需提前定义函数,支持动态指定主题发送,适合业务场景多变的情况:
import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.cloud.stream.function.StreamBridge; import org.springframework.stereotype.Component; import org.springframework.beans.factory.annotation.Autowired; @Component public class OrganizationEventPublisher { private static final Logger logger = LoggerFactory.getLogger(OrganizationEventPublisher.class); @Autowired private StreamBridge streamBridge; public void publishOrganizationChange(String action, String organizationId) { logger.debug("Sending Kafka message {} for Organization Id: {}", action, organizationId); OrganizationChangeModel change = new OrganizationChangeModel( OrganizationChangeModel.class.getTypeName(), action, organizationId, UserContext.getCorrelationId()); // 直接指定主题发送,无需提前配置绑定 streamBridge.send("orgChangeTopic", change); } }
使用StreamBridge时,配置可简化为:
spring.cloud.stream.kafka.binder.brokers=kafka spring.cloud.stream.kafka.binder.auto-create-topics=true
示例2:消息消费(license-service)
原实现问题
原代码使用@EnableBinding(Sink.class)和@StreamListener注解,这两个API在Spring Cloud Stream 4.x中已被标记为废弃,需改用函数式消费模型。
Spring Boot 3 适配方案
1. 配置文件调整
替换原input绑定为函数式输入绑定:
# 函数式输入绑定:函数名logOrgChange对应绑定后缀-in-0 spring.cloud.stream.bindings.logOrgChange-in-0.destination=orgChangeTopic spring.cloud.stream.bindings.logOrgChange-in-0.content-type=application/json # 消费组配置,确保消息被可靠消费 spring.cloud.stream.bindings.logOrgChange-in-0.group=licensingGroup # Kafka Broker地址 spring.cloud.stream.kafka.binder.brokers=kafka
2. 代码改造
移除@EnableBinding和@StreamListener,定义消费函数:
import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; import org.springframework.cloud.client.discovery.EnableDiscoveryClient; import org.springframework.cloud.openfeign.EnableFeignClients; import org.springframework.context.annotation.Bean; import org.springframework.cloud.context.config.annotation.RefreshScope; import java.util.function.Consumer; @SpringBootApplication @RefreshScope @EnableDiscoveryClient @EnableFeignClients public class LicenseServiceApplication { private static final Logger log = LoggerFactory.getLogger(LicenseServiceApplication.class); public static void main(String[] args) { SpringApplication.run(LicenseServiceApplication.class, args); } // 定义消息消费的Consumer函数,函数名与配置中的logOrgChange完全匹配 @Bean public Consumer<OrganizationChangeModel> logOrgChange() { return orgChange -> { log.info("Received an {} event for organization id {}", orgChange.getAction(), orgChange.getOrganizationId()); }; } }
内容的提问来源于stack exchange,提问作者Yll
相关产品推荐
相关产品推荐

