如何在同一测试规约中创建多个配置不同的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
相关产品推荐
相关产品推荐

