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

如何在同一测试规约中创建多个配置不同的EmbeddedKafka实例

要实现两个独立、配置不同的EmbeddedKafka实例,不能直接在类上叠加多个@EmbeddedKafka注解,需要通过编程式手动创建EmbeddedKafkaBroker实例的方式实现,具体操作如下:


实现步骤

1. 移除类上的@EmbeddedKafka注解,改为手动注册两个独立的Broker Bean

在测试配置类中定义两个不同名称的EmbeddedKafkaBroker Bean,分别配置各自的主题、Broker参数:

@Configuration
public class DualKafkaTestConfig {

    // 第一个Kafka实例:绑定TEST_TOPIC_one,offsets复制因子为1
    @Bean("kafkaBrokerOne")
    public EmbeddedKafkaBroker kafkaBrokerOne() {
        return new EmbeddedKafkaBroker(1)
                .topics(Constants.TEST_TOPIC_one)
                .brokerProperty("offsets.topic.replication.factor", 1)
                .brokerProperty("transaction.state.log.replication.factor", 1)
                .brokerProperty("transaction.state.log.min.isr", 1)
                .controlledShutdown(false);
    }

    // 第二个Kafka实例:绑定TEST_TOPIC_two,offsets复制因子为2
    @Bean("kafkaBrokerTwo")
    public EmbeddedKafkaBroker kafkaBrokerTwo() {
        return new EmbeddedKafkaBroker(1)
                .topics(Constants.TEST_TOPIC_two)
                .brokerProperty("offsets.topic.replication.factor", 2)
                .brokerProperty("transaction.state.log.replication.factor", 1)
                .brokerProperty("transaction.state.log.min.isr", 1)
                .controlledShutdown(false);
    }
}

2. 为两个Broker分别配置独立的生产者、消费者组件

每个实例需要对应一套独立的ProducerFactory、KafkaTemplate、ConsumerFactory、监听容器工厂,绑定对应Broker的地址:

// ========== 第一个Kafka实例对应的组件 ==========
@Bean
public ProducerFactory<String, String> producerFactoryOne(@Qualifier("kafkaBrokerOne") EmbeddedKafkaBroker brokerOne) {
    Map<String, Object> props = new HashMap<>();
    props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, brokerOne.getBrokersAsString());
    props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
    props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
    return new DefaultKafkaProducerFactory<>(props);
}

@Bean
public KafkaTemplate<String, String> kafkaTemplateOne(@Qualifier("producerFactoryOne") ProducerFactory<String, String> factory) {
    return new KafkaTemplate<>(factory);
}

@Bean
public ConsumerFactory<String, String> consumerFactoryOne(@Qualifier("kafkaBrokerOne") EmbeddedKafkaBroker brokerOne) {
    Map<String, Object> props = new HashMap<>();
    props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, brokerOne.getBrokersAsString());
    props.put(ConsumerConfig.GROUP_ID_CONFIG, "group-one");
    props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
    props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
    return new DefaultKafkaConsumerFactory<>(props);
}

@Bean
public ConcurrentKafkaListenerContainerFactory<String, String> listenerFactoryOne(@Qualifier("consumerFactoryOne") ConsumerFactory<String, String> factory) {
    ConcurrentKafkaListenerContainerFactory<String, String> containerFactory = new ConcurrentKafkaListenerContainerFactory<>();
    containerFactory.setConsumerFactory(factory);
    return containerFactory;
}

// ========== 第二个Kafka实例对应的组件 ==========
// 按上面的结构复制一套,修改Bean名称、注入kafkaBrokerTwo即可

3. 使用时指定对应组件即可

  • 监听指定实例的主题:在@KafkaListener注解中指定对应监听工厂
@KafkaListener(topics = Constants.TEST_TOPIC_one, containerFactory = "listenerFactoryOne")
public void handleTopicOneMsg(String msg) {
    // 处理第一个实例的消息
}
  • 向指定实例发消息:注入对应KafkaTemplate调用即可
@Autowired
@Qualifier("kafkaTemplateOne")
private KafkaTemplate<String, String> kafkaTemplateOne;

注意事项

两个EmbeddedKafkaBroker默认会自动分配随机端口,不会出现端口冲突,通过getBrokersAsString()可以拿到正确的服务地址。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 03:15:03