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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 14:50:12