如何覆盖ReceiverOptions.addAssignListener()中Consumer函数的代码覆盖率?
问题:如何测试Reactive Kafka中分区分配监听器的代码逻辑
我一直尝试提升这段代码的覆盖率,但始终无法解决。该私有方法由公共方法调用,我认为由于主题订阅操作滞后,没有分区被分配,因此单元测试时addAssignListener()中的逻辑从未被触发。
目标代码
class SomeClass { public ReactiveKafkaConsumerTemplate somePublicMethod(Set<String> topics) { // ... ReceiverOptions<Object, Object> receiverOptions = somePrivateMethod(customConsumerProperties); receiverOptions = receiverOptions.subscription(topics); return new ReactiveKafkaConsumerTemplate<>(receiverOptions); } private ReceiverOptions<Object, Object> somePrivateMethod(Map<String, Object> customConsumerProperties) { boolean featureFlag = someConfig.isEnabled(); ReceiverOptions<Object, Object> basicReceiverOptions = ReceiverOptions.create(customConsumerProperties); basicReceiverOptions = basicReceiverOptions.addAssignListener(partitions -> { partitions.forEach(receiverPartition -> Mono.just(receiverPartition) .doOnNext(receiverPartition1 -> { if (featureFlag) { receiverPartition1.seekToTimestamp(someConfig.getStartTimestamp()); } }) .subscribe()); }); return basicReceiverOptions; } }
当前单元测试代码
@Test void test() { // ... ReactiveKafkaConsumerTemplate<String, String> template = someClass.somePublicMethod(Set.of("sample-topic")); assertNotNull(template); }
解决方案
方案1:反射触发监听器逻辑
通过反射获取ReceiverOptions中存储的分区分配监听器,手动传入模拟的分区集合触发逻辑:
@Test void testAssignListenerLogic() throws NoSuchFieldException, IllegalAccessException { // 初始化配置,开启功能并指定时间戳 SomeConfig someConfig = new SomeConfig(); someConfig.setEnabled(true); someConfig.setStartTimestamp(123456789L); SomeClass someClass = new SomeClass(someConfig); // 获取ReceiverOptions实例 ReactiveKafkaConsumerTemplate template = someClass.somePublicMethod(Set.of("sample-topic")); ReceiverOptions<Object, Object> receiverOptions = template.receiverOptions(); // 反射获取监听器列表(需对应ReceiverOptions的实际实现类,比如Spring Kafka的ReceiverOptionsImpl) Field listenersField = ReceiverOptionsImpl.class.getDeclaredField("assignListeners"); listenersField.setAccessible(true); List<Consumer<List<ReceiverPartition>>> assignListeners = (List<Consumer<List<ReceiverPartition>>>) listenersField.get(receiverOptions); // 模拟分区并触发监听器 ReceiverPartition mockPartition = Mockito.mock(ReceiverPartition.class); assignListeners.forEach(listener -> listener.accept(List.of(mockPartition))); // 验证seek逻辑是否执行 Mockito.verify(mockPartition).seekToTimestamp(123456789L); }
方案2:重构代码提升可测试性
把监听器内的核心逻辑抽成独立的可访问方法,直接测试该方法:
// 重构后的目标类 class SomeClass { private final SomeConfig someConfig; public SomeClass(SomeConfig someConfig) { this.someConfig = someConfig; } public ReactiveKafkaConsumerTemplate somePublicMethod(Set<String> topics) { ReceiverOptions<Object, Object> receiverOptions = somePrivateMethod(customConsumerProperties); receiverOptions = receiverOptions.subscription(topics); return new ReactiveKafkaConsumerTemplate<>(receiverOptions); } private ReceiverOptions<Object, Object> somePrivateMethod(Map<String, Object> customConsumerProperties) { ReceiverOptions<Object, Object> basicReceiverOptions = ReceiverOptions.create(customConsumerProperties); return basicReceiverOptions.addAssignListener(partitions -> partitions.forEach(this::handlePartitionAssignment) ); } // 抽离公共方法,方便测试 void handlePartitionAssignment(ReceiverPartition receiverPartition) { if (someConfig.isEnabled()) { receiverPartition.seekToTimestamp(someConfig.getStartTimestamp()); } } }
对应的测试代码:
@Test void testHandlePartitionWhenFeatureEnabled() { SomeConfig someConfig = new SomeConfig(); someConfig.setEnabled(true); someConfig.setStartTimestamp(123456789L); SomeClass someClass = new SomeClass(someConfig); ReceiverPartition mockPartition = Mockito.mock(ReceiverPartition.class); someClass.handlePartitionAssignment(mockPartition); Mockito.verify(mockPartition).seekToTimestamp(123456789L); } @Test void testHandlePartitionWhenFeatureDisabled() { SomeConfig someConfig = new SomeConfig(); someConfig.setEnabled(false); SomeClass someClass = new SomeClass(someConfig); ReceiverPartition mockPartition = Mockito.mock(ReceiverPartition.class); someClass.handlePartitionAssignment(mockPartition); Mockito.verify(mockPartition, Mockito.never()).seekToTimestamp(Mockito.anyLong()); }
方案3:使用嵌入式Kafka模拟真实场景
借助Spring Kafka的EmbeddedKafkaBroker搭建本地测试环境,触发真实的分区分配流程:
@SpringBootTest @EmbeddedKafka(topics = "sample-topic", partitions = 1) class SomeClassIntegrationTest { @Autowired private EmbeddedKafkaBroker embeddedKafkaBroker; @Test void testAssignListenerWithRealPartitionAssignment() { // 配置消费者属性指向嵌入式Kafka Map<String, Object> consumerProps = KafkaTestUtils.consumerProps("test-group", "true", embeddedKafkaBroker); SomeConfig someConfig = new SomeConfig(); someConfig.setEnabled(true); someConfig.setStartTimestamp(123456789L); SomeClass someClass = new SomeClass(someConfig); // 创建消费者模板并触发消费,触发分区分配 ReactiveKafkaConsumerTemplate template = someClass.somePublicMethod(Set.of("sample-topic")); template.receive() .take(1) .block(Duration.ofSeconds(5)); // 可通过验证offset或结合Mock配置确认seek逻辑生效 } }
内容的提问来源于stack exchange,提问作者tripletk
相关产品推荐
相关产品推荐

