Spring Kafka消费者单元测试:onMessage未触发且Mock服务失效
Alright, let's tackle your Kafka consumer unit testing issues head-on. You're facing two key problems: your mocked MyService isn't being used (Spring keeps spinning up the real instance) and the onMessage method of your MessageListener isn't triggering. Here's how to fix both:
1. Why Your Mocked Service Isn't Being Injected
The most likely culprit is that you're using plain Mockito's @Mock instead of Spring Boot's @MockBean.
@Mockcreates a mock object but doesn't register it with the Spring application context. So when your Kafka consumer bean is created, it still gets the realMyServiceinstance from the context.@MockBean, on the other hand, replaces the existingMyServicebean in the Spring context with your mock. This ensures the consumer injects the mock instead of the real thing.
2. Getting onMessage to Trigger
For your listener's onMessage method to run, your consumer needs to actually receive a Kafka message. In unit tests, you don't want to rely on a real Kafka cluster—use an embedded Kafka server instead.
Fix 1: Integration Test with Embedded Kafka (Spring Context)
This approach loads a lightweight Spring context and spins up an embedded Kafka server to test the full consumer flow:
@SpringBootTest @EmbeddedKafka(topics = "your-topic-name", partitions = 1) @TestPropertySource(properties = { "spring.kafka.consumer.bootstrap-servers=${spring.embedded.kafka.brokers}", "spring.kafka.consumer.group-id=test-consumer-group", "spring.kafka.consumer.auto-offset-reset=earliest" }) public class KafkaConsumerIntegrationTest { @Autowired private KafkaTemplate<String, String> kafkaTemplate; @MockBean // This replaces the real MyService in the Spring context private MyService myService; @Autowired private KafkaListenerEndpointRegistry endpointRegistry; @BeforeEach void setUp() { // Ensure all Kafka listener containers are started endpointRegistry.getListenerContainers().forEach(container -> { if (!container.isRunning()) { container.start(); } }); } @Test void testConsumerProcessesMessage() throws InterruptedException { // 1. Send a test message to the embedded Kafka topic String testPayload = "Hello, Kafka!"; kafkaTemplate.send("your-topic-name", testPayload); // 2. Wait for the consumer to process the message (adjust sleep time if needed) Thread.sleep(1000); // For production tests, use a CountDownLatch for better reliability // 3. Verify your mock service was called with the correct payload verify(myService, Mockito.times(1)).processMessage(testPayload); } }
Fix 2: Pure Unit Test (No Spring Context)
If you want to test just the MessageListener logic without loading the Spring context, manually inject the mocked service and call onMessage directly:
public class KafkaListenerUnitTest { @Mock private MyService myService; private YourMessageListener consumer; // Your MessageListener implementation @BeforeEach void setUp() { // Initialize Mockito mocks MockitoAnnotations.openMocks(this); // Manually inject the mocked service into your listener consumer = new YourMessageListener(myService); } @Test void testOnMessageCallsService() { // Create a mock ConsumerRecord ConsumerRecord<String, String> mockRecord = new ConsumerRecord<>( "your-topic-name", 0, 1L, "test-key", "test-payload" ); // Call onMessage directly consumer.onMessage(mockRecord); // Verify the service method was invoked verify(myService, Mockito.times(1)).processMessage("test-payload"); } }
Key Notes
- For the integration test, make sure your consumer class is annotated with
@KafkaListenerso Spring can register it. - The
auto-offset-reset=earliestproperty ensures the consumer picks up the test message we send (since it's a new consumer group). - If you're using a
CountDownLatchinstead ofThread.sleep, add a latch to your listener class and count down inonMessage, then await it in the test for cleaner synchronization.
内容的提问来源于stack exchange,提问作者Sahil Chhabra

