Spring Boot集成测试:Kafka Test Container与Mock Schema Registry消费异常
解决Mock Schema Registry集成测试中消费者找不到Schema的问题
问题核心原因
你遇到的Subject Not Found错误,本质是Mock Schema Registry是进程内独立实例:当你用mock://URL配置时,生产者和消费者会各自创建一个独立的MockSchemaRegistryClient对象。生产者发送消息时注册的Schema只存在于它自己的Mock实例中,消费者的Mock实例根本没有这个Schema的记录,自然会返回404。
解决方案1:共享全局MockSchemaRegistryClient实例
放弃使用mock://URL,手动创建一个全局的MockSchemaRegistryClient实例,让生产者和消费者共用它,确保Schema注册和读取的是同一个内存存储。
代码示例
- 定义全局Mock客户端:
private static final MockSchemaRegistryClient MOCK_REGISTRY = new MockSchemaRegistryClient();
- 修改消费者创建逻辑,指定使用全局Mock实例:
public static KafkaConsumer<String, CarDTO> createEventConsumer() { Properties props = new Properties(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaContainer.getBootstrapServers()); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); props.put(KafkaAvroDeserializerConfig.SPECIFIC_AVRO_READER_CONFIG, "true"); props.put(ConsumerConfig.GROUP_ID_CONFIG, "kafkatest"); // 手动注入全局Mock实例到反序列化器 KafkaAvroDeserializer valueDeserializer = new KafkaAvroDeserializer(MOCK_REGISTRY); valueDeserializer.configure(props, false); return new KafkaConsumer<>(props, new StringDeserializer(), valueDeserializer); }
- 生产者端同步使用该全局Mock实例:
// 生产者配置时同样注入全局Mock客户端 KafkaAvroSerializer valueSerializer = new KafkaAvroSerializer(MOCK_REGISTRY); valueSerializer.configure(props, false); KafkaProducer<String, CarDTO> producer = new KafkaProducer<>(props, new StringSerializer(), valueSerializer);
解决方案2:改用Testcontainer的Schema Registry容器
如果Mock Registry的局限性太多,不如直接使用Confluent提供的Schema Registry Testcontainer,和你的Kafka容器搭配使用,完全模拟真实环境,避免Mock带来的不一致问题。
代码示例
// 启动Schema Registry容器并关联Kafka容器 SchemaRegistryContainer schemaRegistryContainer = new SchemaRegistryContainer( DockerImageName.parse("confluentinc/cp-schema-registry").withTag("5.4.3") ).withKafka(kafkaContainer); schemaRegistryContainer.start(); // 生产者和消费者统一配置真实的Schema Registry地址 String schemaRegistryUrl = schemaRegistryContainer.getSchemaRegistryUrl(); props.put(AbstractKafkaAvroSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG, schemaRegistryUrl);
额外注意点
- 你提供的轮询代码存在泛型不一致问题:
ConsumerRecords<String, CarDTO>和循环里的ConsumerRecord<String, StoreReplenOrderDTO>会导致类型转换异常,需修正为一致的DTO类型。 - 使用
ErrorHandlingDeserializer能拿到key是因为它捕获了反序列化错误,但并未解决根本问题,只是跳过了value的反序列化失败逻辑。
内容的提问来源于stack exchange,提问作者user1555190
相关产品推荐
相关产品推荐

