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

如何从字符串构造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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 10:11:59