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

如何测试Spring-Kafka消息发送失败时的onFailure回调逻辑?

解决方案

1. 修正MockBean类型不匹配问题

你的MessageSending类中注入的是Map<String, KafkaTemplate<String, String>>,但测试类里的@MockBean声明为Map<Country, KafkaTemplate<String, String>>,类型不匹配会导致mock无法正确注入,首先修正这个类型:

@MockBean
Map<String, KafkaTemplate<String, String>> producerByCountry;

2. 正确Mock KafkaTemplate与ListenableFuture

不需要依赖真实的EmbeddedKafka模拟失败场景,直接通过Mock控制KafkaTemplate的行为,手动触发onFailure回调:

@SpringBootTest
public class MessageSendingTest {

    @MockBean
    Map<String, KafkaTemplate<String, String>> producerByCountry;

    @Mock
    KafkaTemplate<String, String> mockKafkaTemplate;

    @Mock
    ListenableFuture<SendResult<String, String>> mockListenableFuture;

    @Autowired
    MessageSending messageSending;

    @Test
    void failTest(CapturedOutput capturedOutput) {
        // 指定map返回mock的KafkaTemplate
        given(producerByCountry.get("countryName")).willReturn(mockKafkaTemplate);
        // 指定send方法返回mock的ListenableFuture
        given(mockKafkaTemplate.send(anyString(), anyString())).willReturn(mockListenableFuture);

        // 拦截addCallback调用,手动触发onFailure
        doAnswer(invocation -> {
            ListenableFutureCallback<SendResult<String, String>> callback = invocation.getArgument(0);
            KafkaProducerException exception = new KafkaProducerException(
                    new ProducerRecord<>("countryTopic", "data"),
                    "模拟发送失败",
                    new RuntimeException("测试异常")
            );
            callback.onFailure(exception);
            return null;
        }).when(mockListenableFuture).addCallback(any(ListenableFutureCallback.class));

        // 执行发送逻辑
        messageSending.sendMessage("data");

        // 验证日志输出
        assertThat(capturedOutput).contains("failed");
    }
}

3. 解释之前的Mockito异常原因

你之前直接mock(ListenableFuture.class)并stub它的addCallback方法,但这个mock实例从未被KafkaTemplate的send方法返回,属于未被使用的无效Stub,因此Mockito抛出UnnecessaryStubbingException。正确的做法是让KafkaTemplate返回你预先mock好的ListenableFuture实例,再针对这个实例做回调拦截。

4. 额外注意事项

  • 确保MessageSending类中的日志对象正确注入(比如使用@Slf4j注解),否则CapturedOutput无法捕获日志输出。
  • 若使用JUnit 5,需确保CapturedOutput通过@ExtendWith(OutputCaptureExtension.class)或@SpringBootTest自动启用了输出捕获功能。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 13:00:50