如何从字符串构造Kafka Consumer Record编写Java Kafka消费者JUnit测试用例
编写Kafka消费者processConsumerRecord方法的JUnit测试
咱们一步步来搞定这个测试,先解决最基础的——如何从字符串构造ConsumerRecord和ConsumerRecords,再写具体的测试用例。
1. 从字符串构造ConsumerRecord(含GenericRecord)
因为你的方法处理的是ConsumerRecords<String, GenericRecord>,而GenericRecord是Avro的通用记录,所以得先构造符合你业务Schema的GenericRecord,再包装成ConsumerRecord。
步骤1:构造GenericRecord
假设你的业务用的Avro Schema是这样的(举个实际例子):
{ "type": "record", "name": "UserEvent", "fields": [ {"name": "userId", "type": "string"}, {"name": "eventType", "type": "string"} ] }
你可以用Avro的API把字符串数据转换成GenericRecord:
import org.apache.avro.Schema; import org.apache.avro.generic.GenericData; import org.apache.avro.generic.GenericRecord; // 加载Schema(也可以从文件读取,这里直接用字符串构造) Schema schema = new Schema.Parser().parse("{\n" + " \"type\": \"record\",\n" + " \"name\": \"UserEvent\",\n" + " \"fields\": [\n" + " {\"name\": \"userId\", \"type\": \"string\"},\n" + " {\"name\": \"eventType\", \"type\": \"string\"}\n" + " ]\n" + "}"); // 用字符串数据填充GenericRecord GenericRecord userEvent = new GenericData.Record(schema); userEvent.put("userId", "user_123"); userEvent.put("eventType", "login");
步骤2:构造ConsumerRecord
有了GenericRecord,就可以创建ConsumerRecord了,指定topic、分区、偏移量等必要参数:
import org.apache.kafka.clients.consumer.ConsumerRecord; ConsumerRecord<String, GenericRecord> consumerRecord = new ConsumerRecord<>( "test-topic", // Kafka Topic名称 0, // 分区号 1L, // 偏移量 "key_123", // 消息Key(字符串类型) userEvent // 消息Value(GenericRecord类型) );
步骤3:包装成ConsumerRecords
因为你的方法接收的是ConsumerRecords,需要把单个或多个ConsumerRecord包装进去:
import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.common.TopicPartition; import java.util.List; import java.util.Map; // 创建TopicPartition对象 TopicPartition topicPartition = new TopicPartition("test-topic", 0); // 把ConsumerRecord放到List里 List<ConsumerRecord<String, GenericRecord>> recordList = List.of(consumerRecord); // 构建ConsumerRecords ConsumerRecords<String, GenericRecord> consumerRecords = new ConsumerRecords<>(Map.of(topicPartition, recordList));
2. 编写processConsumerRecord方法的JUnit测试用例
咱们用JUnit 5 + Mockito来写测试,这样不需要依赖真实的Kafka集群,只需要模拟必要的组件。
步骤1:添加必要依赖(如果还没加)
确保你的pom.xml(Maven)里有这些测试依赖:
<dependencies> <!-- JUnit 5 --> <dependency> <groupId>org.junit.jupiter</groupId> <artifactId>junit-jupiter-api</artifactId> <version>5.9.2</version> <scope>test</scope> </dependency> <dependency> <groupId>org.junit.jupiter</groupId> <artifactId>junit-jupiter-engine</artifactId> <version>5.9.2</version> <scope>test</scope> </dependency> <!-- Mockito --> <dependency> <groupId>org.mockito</groupId> <artifactId>mockito-core</artifactId> <version>4.11.0</version> <scope>test</scope> </dependency> <dependency> <groupId>org.mockito</groupId> <artifactId>mockito-junit-jupiter</artifactId> <version>4.11.0</version> <scope>test</scope> </dependency> <!-- Avro --> <dependency> <groupId>org.apache.avro</groupId> <artifactId>avro</artifactId> <version>1.11.1</version> <scope>test</scope> </dependency> <!-- Kafka Clients --> <dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> <version>3.4.0</version> <scope>test</scope> </dependency> </dependencies>
步骤2:编写测试类
假设你的消费者类叫KafkaEventConsumer,里面包含processConsumerRecord方法。测试类示例如下:
import org.apache.avro.Schema; import org.apache.avro.generic.GenericData; import org.apache.avro.generic.GenericRecord; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.common.TopicPartition; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; import org.mockito.Mock; import org.mockito.junit.jupiter.MockitoExtension; import java.util.List; import java.util.Map; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.times; @ExtendWith(MockitoExtension.class) class KafkaEventConsumerTest { // 模拟Kafka Consumer对象,不需要真实实例 @Mock private Consumer<String, GenericRecord> mockConsumer; // 实例化要测试的消费者类 private final KafkaEventConsumer consumerUnderTest = new KafkaEventConsumer(); @Test void processConsumerRecord_WhenEventProcessedAndOffsetCommitted_ShouldCommitOffset() { // 1. 准备测试数据:构造ConsumerRecords ConsumerRecords<String, GenericRecord> records = buildTestConsumerRecords("user_123", "login"); // 2. 调用要测试的方法:事件已处理、偏移量要提交、错误计数为0 consumerUnderTest.processConsumerRecord(records, true, true, 0, 0, mockConsumer); // 3. 验证逻辑:是否调用了Consumer的提交方法(根据你的业务逻辑调整,比如commitAsync) verify(mockConsumer, times(1)).commitSync(); } @Test void processConsumerRecord_WhenEventNotProcessed_ShouldNotCommitOffset() { // 1. 准备测试数据 ConsumerRecords<String, GenericRecord> records = buildTestConsumerRecords("user_456", "logout"); // 2. 调用方法:事件未处理,不提交偏移量 consumerUnderTest.processConsumerRecord(records, false, false, 0, 0, mockConsumer); // 3. 验证:没有调用提交方法 verify(mockConsumer, times(0)).commitSync(); } // 抽取复用的方法,用来构造测试用的ConsumerRecords private ConsumerRecords<String, GenericRecord> buildTestConsumerRecords(String userId, String eventType) { Schema schema = new Schema.Parser().parse("{\n" + " \"type\": \"record\",\n" + " \"name\": \"UserEvent\",\n" + " \"fields\": [\n" + " {\"name\": \"userId\", \"type\": \"string\"},\n" + " {\"name\": \"eventType\", \"type\": \"string\"}\n" + " ]\n" + "}"); GenericRecord userEvent = new GenericData.Record(schema); userEvent.put("userId", userId); userEvent.put("eventType", eventType); ConsumerRecord<String, GenericRecord> record = new ConsumerRecord<>("test-topic", 0, 1L, "key_" + userId, userEvent); TopicPartition tp = new TopicPartition("test-topic", 0); return new ConsumerRecords<>(Map.of(tp, List.of(record))); } }
3. 测试更多业务场景
根据你的processConsumerRecord方法的逻辑,还可以覆盖这些场景:
- 当
totalErrorCountFromSinkService或totalErrorCount达到阈值时的处理逻辑 - 空
ConsumerRecords(无消息)的情况 - 多个
ConsumerRecord批量处理的场景 - 消息处理抛出异常时的容错逻辑
小提示
- 如果你的Schema比较复杂,可以把Schema文件放到
src/test/resources目录下,用new Schema.Parser().parse(new File("src/test/resources/user-event.avsc"))加载,更简洁。 - 用Mockito的
ArgumentCaptor可以捕获方法中传递的参数,比如验证提交的偏移量是否正确。 - 可以把构造测试数据的逻辑进一步封装成工具类,减少测试代码的重复。
内容的提问来源于stack exchange,提问作者Srinivas Bandaru
相关产品推荐
相关产品推荐

