如何测试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
相关产品推荐
相关产品推荐

