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

SpringBoot Kafka生产者单元测试编写及异常用例报错排查

问题根因分析

你写的测试用例无法运行主要有4个核心错误:

  • Topic名称不匹配:生产代码中实际使用的消息topic是硬编码的output-flow,但测试代码打桩时用的topic是test,导致mock规则完全不匹配,调用kafkaProducer.send时不会触发你预设的返回/抛异常逻辑。
  • 错误mock被测类:MessageProducer是你要测试的业务类,必须真实实例化,你对mockProducer.sendMessage做打桩会直接跳过真实业务逻辑,测试完全无效。
  • KafkaTemplate返回值不符合要求:KafkaTemplate.send()方法返回的是Future类型对象,生产代码后续调用了.get()方法获取发送结果,你打桩时没有返回符合类型的Future对象,会直接报类型错误或者逻辑不生效。
  • 异常场景断言逻辑错误:生产代码的sendMessage方法内部已经catch了所有Exception,异常发生时只会返回false,不会向外抛出异常,所以用assertThrows判断抛异常必然失败;同时异常用例里调用方法传了null,和你打桩时的参数message不匹配,也会导致stub不生效。

注:你贴出的生产代码中private String topicName = "output-flow";定义在sendMessage方法内部会直接编译报错,属于代码粘贴时的格式问题,实际使用时请将该字段定义为类成员变量。

正确单元测试实现

以下是基于Mockito + JUnit5的可运行测试代码,覆盖发送成功、send接口直接抛异常、future.get()抛异常三个核心场景:

import org.junit.jupiter.api.extension.ExtendWith;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.MethodSource;
import org.mockito.InjectMocks;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.support.SendResult;
import org.springframework.util.concurrent.ListenableFuture;
import java.util.stream.Stream;
import org.junit.jupiter.params.provider.Arguments;

import static org.junit.jupiter.api.Assertions.*;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.*;

@ExtendWith(MockitoExtension.class)
public class MessageProducerTest {

    // 和生产代码硬编码的topic保持一致
    private static final String EXPECTED_TOPIC = "output-flow";

    @Mock
    private KafkaTemplate<String, ReceivedMessage> kafkaProducer;

    // 自动注入mock的依赖,不用手动new被测类
    @InjectMocks
    private MessageProducer messageProducer;

    @ParameterizedTest
    @MethodSource("getTransactionProvider")
    public void sendMessage_Success_ReturnsTrue(ReceivedMessage message) throws Exception {
        // 构造mock的Future对象,模拟future.get()正常返回
        ListenableFuture<SendResult<String, ReceivedMessage>> mockFuture = mock(ListenableFuture.class);
        when(mockFuture.get()).thenReturn(mock(SendResult.class));
        
        // 给kafkaTemplate.send打桩,参数严格匹配topic和消息
        when(kafkaProducer.send(eq(EXPECTED_TOPIC), eq(message))).thenReturn(mockFuture);

        // 调用被测方法
        boolean result = messageProducer.sendMessage(message);

        // 断言返回结果,验证send方法被正确调用
        assertTrue(result);
        verify(kafkaProducer, times(1)).send(EXPECTED_TOPIC, message);
    }

    @ParameterizedTest
    @MethodSource("getTransactionProvider")
    public void sendMessage_SendThrowsException_ReturnsFalse(ReceivedMessage message) {
        // 模拟send方法直接抛出异常
        when(kafkaProducer.send(eq(EXPECTED_TOPIC), eq(message))).thenThrow(new RuntimeException("kafka send error"));

        boolean result = messageProducer.sendMessage(message);

        // 业务代码捕获异常后返回false,不会向外抛异常
        assertFalse(result);
        verify(kafkaProducer, times(1)).send(EXPECTED_TOPIC, message);
    }

    @ParameterizedTest
    @MethodSource("getTransactionProvider")
    public void sendMessage_FutureGetThrowsException_ReturnsFalse(ReceivedMessage message) throws Exception {
        // 模拟send正常返回,但调用future.get()时抛异常(比如Broker返回错误、线程中断等场景)
        ListenableFuture<SendResult<String, ReceivedMessage>> mockFuture = mock(ListenableFuture.class);
        when(mockFuture.get()).thenThrow(new RuntimeException("future get error"));
        when(kafkaProducer.send(eq(EXPECTED_TOPIC), eq(message))).thenReturn(mockFuture);

        boolean result = messageProducer.sendMessage(message);

        assertFalse(result);
        verify(kafkaProducer, times(1)).send(EXPECTED_TOPIC, message);
    }

    // 参数构造方法,根据你的ReceivedMessage结构调整即可
    static Stream<Arguments> getTransactionProvider() {
        return Stream.of(
                Arguments.of(new ReceivedMessage(/* 构造第一组测试消息 */)),
                Arguments.of(new ReceivedMessage(/* 构造第二组测试消息 */))
        );
    }
}
适配说明
  • 如果你使用的Spring Kafka版本在2.9及以上,KafkaTemplate.send()返回值已经改为CompletableFuture<SendResult>,打桩时将Future类型替换为CompletableFuture即可:成功场景用CompletableFuture.completedFuture(mock(SendResult.class))构造返回值,异常场景用CompletableFuture.failedFuture(new RuntimeException("error"))构造返回值。
  • 所有打桩操作只针对外部依赖(也就是这里的KafkaTemplate),绝对不要对被测类的方法做mock,否则会跳过真实业务逻辑,测试没有任何意义。
  • 打桩时的参数必须和生产代码实际传入的参数严格匹配,比如topic名称、消息对象,如果不需要严格匹配可以用any()、anyString()这类宽松匹配器,但核心场景建议用精确匹配保证测试有效性。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 23:27:17