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

Spring应用中从application.yml提取replyTopic配置用于@KafkaListener

提取Spring配置中的replyTopic用于@KafkaListener

原始配置(YAML)

integration
mapping:
    - producer:
        name: somename
        topic: sometopic
        replyTopic: sometopic
        replyPartition: 1
    consumer:
        name: somename
        field_1: ...
        field_2: ...
    - producer:
        name: somename
        topic: sometopic
        replyTopic: sometopic
        replyPartition: 1
    consumer:
        name: somename
        field_1: ...
        field_2: ...

需求

无需绑定整个配置结构,仅提取所有replyTopic字段值,作为@KafkaListener注解的监听主题,要求启动时即可获取这些主题名称。

实现方案

1. 绑定配置属性类

创建仅包含所需字段的配置类,通过@ConfigurationProperties绑定integration.mapping下的producer.replyTopic:

@ConfigurationProperties(prefix = "integration")
public class KafkaReplyTopicProperties {
    private List<Mapping> mapping;

    public static class Mapping {
        private Producer producer;

        public static class Producer {
            private String replyTopic;

            public String getReplyTopic() {
                return replyTopic;
            }
        }

        public Producer getProducer() {
            return producer;
        }
    }

    public List<Mapping> getMapping() {
        return mapping;
    }
}

2. 启用配置属性并编写KafkaListener

在配置类中启用配置属性,通过SpEL表达式动态提取所有replyTopic值:

@Configuration
@EnableConfigurationProperties(KafkaReplyTopicProperties.class)
public class KafkaListenerConfig {

    @Bean
    public ConsumerFactory<String, Object> consumerFactory() {
        // 配置消费者工厂,示例省略具体参数
        Map<String, Object> props = new HashMap<>();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        // 其他必要配置...
        return new DefaultKafkaConsumerFactory<>(props);
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, Object> kafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        return factory;
    }

    @KafkaListener(topics = "#{@kafkaReplyTopicProperties.mapping.stream()
            .filter(m -> m.getProducer() != null)
            .map(m -> m.getProducer().getReplyTopic())
            .toArray(String[]::new)}")
    public void handleReplyMessages(String message) {
        // 处理监听消息的业务逻辑
        System.out.println("Received message: " + message);
    }
}

关键说明

  • SpEL表达式#{@kafkaReplyTopicProperties.mapping.stream().map(...).toArray(...)}会在启动时遍历配置中的所有mapping项,提取有效replyTopic值并转为字符串数组,直接传入@KafkaListener的topics参数。
  • 添加filter(m -> m.getProducer() != null)避免因配置中存在无producer的项导致空指针异常。
  • 确保KafkaReplyTopicProperties被Spring容器管理,默认Bean名为类名首字母小写(kafkaReplyTopicProperties),与SpEL中的引用一致。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 04:01:19