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

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注册和读取的是同一个内存存储。

代码示例

  1. 定义全局Mock客户端:
private static final MockSchemaRegistryClient MOCK_REGISTRY = new MockSchemaRegistryClient();
  1. 修改消费者创建逻辑,指定使用全局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);
}
  1. 生产者端同步使用该全局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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 17:16:28