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

