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

如何为kafka-clients-1.0.0的KafkaProducer编写基于Mock Broker的JUnit测试

使用Mock Kafka Broker测试KafkaProducer(kafka-clients-1.0.0)

核心思路

由于无法直接用MockProducer替代真实的KafkaProducer,我们可以通过启动嵌入式Mock Kafka Broker(基于官方kafka-test模块),让真实的KafkaProducer连接到这个Broker发送消息,再通过消费者验证消息投递结果。

依赖准备

需要引入kafka-test依赖(版本与kafka-clients保持一致),以Maven为例:

<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka_2.11</artifactId>
    <version>1.0.0</version>
    <classifier>test</classifier>
    <scope>test</scope>
</dependency>

注:artifactId中的2.11为Scala版本,需根据项目实际使用的Scala版本调整

JUnit测试用例实现

import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.apache.kafka.common.serialization.StringSerializer;
import org.junit.After;
import org.junit.Before;
import org.junit.Test;
import org.apache.kafka.streams.integration.utils.EmbeddedKafkaCluster;

import java.util.Collections;
import java.util.Properties;
import java.util.concurrent.TimeUnit;

import static org.junit.Assert.assertEquals;

public class KafkaProducerIntegrationTest {

    private EmbeddedKafkaCluster embeddedKafka;
    private KafkaProducer<String, String> producer;
    private KafkaConsumer<String, String> consumer;
    private static final String TEST_TOPIC = "test-topic";

    @Before
    public void setUp() {
        // 启动单节点Mock Kafka Broker
        embeddedKafka = new EmbeddedKafkaCluster(1);
        embeddedKafka.start();

        // 配置真实KafkaProducer,连接到Mock Broker
        Properties producerProps = new Properties();
        producerProps.put("bootstrap.servers", embeddedKafka.bootstrapServers());
        producerProps.put("key.serializer", StringSerializer.class.getName());
        producerProps.put("value.serializer", StringSerializer.class.getName());
        producerProps.put("acks", "all"); // 确保消息被Broker确认
        producer = new KafkaProducer<>(producerProps);

        // 创建测试用Topic
        embeddedKafka.createTopic(TEST_TOPIC);

        // 配置消费者,用于验证消息投递结果
        Properties consumerProps = new Properties();
        consumerProps.put("bootstrap.servers", embeddedKafka.bootstrapServers());
        consumerProps.put("group.id", "test-consumer-group");
        consumerProps.put("key.deserializer", StringDeserializer.class.getName());
        consumerProps.put("value.deserializer", StringDeserializer.class.getName());
        consumerProps.put("auto.offset.reset", "earliest");
        consumer = new KafkaConsumer<>(consumerProps);
        consumer.subscribe(Collections.singletonList(TEST_TOPIC));
    }

    @Test
    public void testSendMessageSuccess() throws InterruptedException {
        // 构造并发送消息
        String testKey = "sample-key";
        String testValue = "sample-value";
        ProducerRecord<String, String> record = new ProducerRecord<>(TEST_TOPIC, testKey, testValue);
        // 同步等待消息发送完成,避免异步导致的测试时序问题
        producer.send(record).get(5, TimeUnit.SECONDS);

        // 拉取消息并验证
        ConsumerRecords<String, String> records = consumer.poll(TimeUnit.SECONDS.toMillis(5));
        assertEquals("消息数量不匹配", 1, records.count());
        
        ConsumerRecord<String, String> receivedRecord = records.iterator().next();
        assertEquals("Key不匹配", testKey, receivedRecord.key());
        assertEquals("Value不匹配", testValue, receivedRecord.value());
        assertEquals("Topic不匹配", TEST_TOPIC, receivedRecord.topic());
    }

    @After
    public void tearDown() {
        // 清理资源,避免测试污染
        if (producer != null) producer.close();
        if (consumer != null) consumer.close();
        if (embeddedKafka != null) embeddedKafka.stop();
    }
}

关键说明

  • EmbeddedKafkaCluster是官方提供的嵌入式Broker实现,完全兼容真实Kafka的协议,能和KafkaProducer/KafkaConsumer无缝交互
  • 测试中使用真实的KafkaProducer实例,无需修改原有业务代码的Producer逻辑
  • 通过同步调用send(record).get(...)确保消息发送完成后再执行验证,避免异步操作导致的测试失败
  • 在@After方法中关闭所有资源,保证每个测试用例的独立性

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 09:10:33