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

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.

  • @Mock creates 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 real MyService instance from the context.
  • @MockBean, on the other hand, replaces the existing MyService bean 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 @KafkaListener so Spring can register it.
  • The auto-offset-reset=earliest property ensures the consumer picks up the test message we send (since it's a new consumer group).
  • If you're using a CountDownLatch instead of Thread.sleep, add a latch to your listener class and count down in onMessage, then await it in the test for cleaner synchronization.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 07:23:53