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

如何编写考虑事件头的Kafka反序列化单元测试?覆盖包迁移报错场景

Kafka事件跨包迁移后的反序列化单元测试方案

核心测试目标

聚焦验证消费者对带有新包路径__TypeId__头的事件的反序列化逻辑,覆盖生产中出现的两类异常场景,同时验证修复后的兼容性。


1. 测试前置准备

  • 依赖:Spring Kafka Test、JUnit 5、Immutables 编译插件
  • 定义测试用的Immutables事件类(模拟迁移前后的包结构):
    // 旧包路径
    package com.our.company.package;
    import org.immutables.value.Value;
    @Value.Immutable
    public interface MyEventClass {
        String getEventId();
        String getContent();
    }
    
    // 新包路径
    package com.our.company.otherPackage;
    import org.immutables.value.Value;
    @Value.Immutable
    public interface MyEventClass {
        String getEventId();
        String getContent();
    }
    

2. 单元测试场景(直接测试反序列化器)

场景1:未配置新包到可信列表,触发可信包异常

验证当__TypeId__头为新包路径,但消费者仅信任旧包时,是否抛出预期的非法参数异常:

import org.junit.jupiter.api.Test;
import org.springframework.kafka.support.serializer.JsonDeserializer;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.common.header.internals.RecordHeader;
import org.apache.kafka.common.header.internals.RecordHeaders;
import java.nio.charset.StandardCharsets;
import java.util.Map;
import static org.junit.jupiter.api.Assertions.assertThrows;

class KafkaEventDeserializerTest {
    private static final String OLD_TYPE_ID = "com.our.company.package.ImmutableMyEventClass";
    private static final String NEW_TYPE_ID = "com.our.company.otherPackage.ImmutableMyEventClass";
    private static final String EVENT_JSON = "{\"eventId\":\"test-123\",\"content\":\"test-content\"}";

    @Test
    void shouldThrowTrustedPackageExceptionForNewTypeId() {
        // 初始化反序列化器,仅信任旧包
        JsonDeserializer<com.our.company.package.MyEventClass> deserializer = new JsonDeserializer<>();
        deserializer.configure(
            Map.of(
                JsonDeserializer.TRUSTED_PACKAGES, "com.our.company.package",
                JsonDeserializer.USE_TYPE_INFO_HEADERS, "true"
            ),
            false
        );

        // 构造带新包TypeId头的消费记录
        ConsumerRecord<String, String> record = new ConsumerRecord<>(
            "test-topic", 0, 0L, "key", EVENT_JSON,
            new RecordHeaders().add(new RecordHeader(JsonDeserializer.TYPE_ID_HEADER, NEW_TYPE_ID.getBytes(StandardCharsets.UTF_8))),
            0
        );

        // 验证异常抛出
        assertThrows(IllegalArgumentException.class, () -> deserializer.deserialize("test-topic", record));
    }
}

场景2:引用旧包类但TypeId为新包,触发类找不到异常

验证当消费者代码仍引用旧包类,但收到新包TypeId头的事件时,是否抛出类找不到异常:

@Test
void shouldThrowClassNotFoundExceptionForMismatchedTypeId() {
    // 初始化反序列化器,信任所有包(排除可信包干扰)
    JsonDeserializer<com.our.company.package.MyEventClass> deserializer = new JsonDeserializer<>();
    deserializer.configure(
        Map.of(
            JsonDeserializer.TRUSTED_PACKAGES, "*",
            JsonDeserializer.USE_TYPE_INFO_HEADERS, "true"
        ),
        false
    );

    // 构造带新包TypeId头的消费记录
    ConsumerRecord<String, String> record = new ConsumerRecord<>(
        "test-topic", 0, 0L, "key", EVENT_JSON,
        new RecordHeaders().add(new RecordHeader(JsonDeserializer.TYPE_ID_HEADER, NEW_TYPE_ID.getBytes(StandardCharsets.UTF_8))),
        0
    );

    // 验证异常抛出
    assertThrows(ClassNotFoundException.class, () -> deserializer.deserialize("test-topic", record));
}

场景3:配置正确后正常反序列化

子场景3.1:消费者改用新包类+配置新包可信

验证消费者切换到新包类,并将新包加入可信列表后,能正常处理新TypeId头的事件:

@Test
void shouldDeserializeSuccessfullyWithNewClass() {
    // 初始化反序列化器,使用新包类并信任新包
    JsonDeserializer<com.our.company.otherPackage.MyEventClass> deserializer = new JsonDeserializer<>(
        com.our.company.otherPackage.MyEventClass.class
    );
    deserializer.configure(
        Map.of(
            JsonDeserializer.TRUSTED_PACKAGES, "com.our.company.otherPackage",
            JsonDeserializer.USE_TYPE_INFO_HEADERS, "true"
        ),
        false
    );

    // 构造带新包TypeId头的消费记录
    ConsumerRecord<String, String> record = new ConsumerRecord<>(
        "test-topic", 0, 0L, "key", EVENT_JSON,
        new RecordHeaders().add(new RecordHeader(JsonDeserializer.TYPE_ID_HEADER, NEW_TYPE_ID.getBytes(StandardCharsets.UTF_8))),
        0
    );

    // 验证反序列化成功
    com.our.company.otherPackage.MyEventClass event = deserializer.deserialize("test-topic", record);
    assert event.getEventId().equals("test-123");
    assert event.getContent().equals("test-content");
}

子场景3.2:保留旧包类+配置类型映射兼容

验证通过TYPE_MAPPINGS配置,让旧包类兼容新TypeId头的事件:

@Test
void shouldDeserializeSuccessfullyWithTypeMapping() {
    // 初始化反序列化器,使用旧包类,配置新TypeId到旧类的映射
    JsonDeserializer<com.our.company.package.MyEventClass> deserializer = new JsonDeserializer<>(
        com.our.company.package.MyEventClass.class
    );
    deserializer.configure(
        Map.of(
            JsonDeserializer.TRUSTED_PACKAGES, "*",
            JsonDeserializer.USE_TYPE_INFO_HEADERS, "true",
            JsonDeserializer.TYPE_MAPPINGS, 
            "com.our.company.otherPackage.ImmutableMyEventClass:com.our.company.package.ImmutableMyEventClass"
        ),
        false
    );

    // 构造带新包TypeId头的消费记录
    ConsumerRecord<String, String> record = new ConsumerRecord<>(
        "test-topic", 0, 0L, "key", EVENT_JSON,
        new RecordHeaders().add(new RecordHeader(JsonDeserializer.TYPE_ID_HEADER, NEW_TYPE_ID.getBytes(StandardCharsets.UTF_8))),
        0
    );

    // 验证反序列化成功
    com.our.company.package.MyEventClass event = deserializer.deserialize("test-topic", record);
    assert event.getEventId().equals("test-123");
}

3. 集成测试补充(验证完整KafkaListener逻辑)

用@EmbeddedKafka启动嵌入式Kafka,验证完整的消费者监听逻辑:

import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.kafka.test.context.EmbeddedKafka;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.LinkedBlockingQueue;
import static org.awaitility.Awaitility.await;

@SpringBootTest
@EmbeddedKafka(partitions = 1, topics = "${com.our.company.topic}")
class KafkaListenerIntegrationTest {

    @Autowired
    private KafkaTemplate<String, Object> kafkaTemplate;

    private final BlockingQueue<com.our.company.otherPackage.MyEventClass> eventQueue = new LinkedBlockingQueue<>();

    // 模拟生产环境的KafkaListener
    @KafkaListener(topics = "${com.our.company.topic}", containerFactory = "containerFactory")
    void onEvent(ConsumerRecord<String, com.our.company.otherPackage.MyEventClass> record) {
        eventQueue.add(record.value());
    }

    @Test
    void shouldConsumeEventWithNewTypeIdHeader() {
        // 构造带新TypeId头的消息
        com.our.company.otherPackage.MyEventClass event = com.our.company.otherPackage.ImmutableMyEventClass.builder()
            .eventId("test-456")
            .content("integration-test-content")
            .build();

        Message<com.our.company.otherPackage.MyEventClass> message = MessageBuilder.withPayload(event)
            .setHeader(KafkaHeaders.TYPE_ID, "com.our.company.otherPackage.ImmutableMyEventClass")
            .build();

        // 发送消息
        kafkaTemplate.send(message);

        // 验证消费者成功接收并反序列化
        await().until(() -> !eventQueue.isEmpty());
        com.our.company.otherPackage.MyEventClass receivedEvent = eventQueue.poll();
        assert receivedEvent.getEventId().equals("test-456");
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 13:37:57