如何编写考虑事件头的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
相关产品推荐
相关产品推荐

