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

如何覆盖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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 10:24:51