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

Spring Boot Kafka消费者连接Schema Registry消费Avro报错咨询

问题根因

你遇到的返回<字符的报错,本质是Schema Registry客户端发起请求时未携带SSL认证配置,服务端/网关返回了HTML格式的错误页面(如认证失败页、登录页),无法被JSON解析器识别,导致拉取Schema失败。
Spring Boot自动装配Kafka消费者时,默认只会将SSL配置传递给Kafka客户端,不会同步给Confluent的Schema Registry客户端,这就是手动写Properties的POC能正常运行,但自动装配的@KafkaListener消费Avro失败的核心原因。

修复步骤

1 修正application.yaml配置

先清理配置里的冗余注释错误,补充Avro反序列化必要参数:

spring:
  kafka:
    bootstrap-servers: your-kafka-server:port
    client-id: your-client-id
    consumer:
      ssl:
        key-password: your-password
        key-store-password: your-password
        key-store-location: classpath:your-keystore.jks # 放在resources下用classpath前缀,无需手动加载
        key-store-type: jks
        trust-store-password: your-password
        trust-store-location: classpath:your-truststore.jks
        trust-store-type: jks
      auto-offset-reset: earliest
      group-id: your-group-id # 统一配置group-id避免注解重复写
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: io.confluent.kafka.serializers.KafkaAvroDeserializer
    properties:
      schema.registry.url: https://your-schema-registry-url # 去掉多余的//前缀,写完整合法的URL
      sasl.mechanism: PLAIN
      security.protocol: SSL
      ssl.endpoint.identification.algorithm: "" # 自签证书需关闭域名校验时明确写空字符串
      # 补充Avro反序列化配置,使用GenericRecord时必须设为false
      specific.avro.reader: false

2 自定义Kafka消费者配置类

新增配置类,将Kafka的SSL配置同步传递给Schema Registry客户端,无需硬编码敏感信息:

import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.core.ConsumerFactory;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
import io.confluent.kafka.serializers.KafkaAvroDeserializer;
import io.confluent.kafka.serializers.KafkaAvroDeserializerConfig;

import java.util.HashMap;
import java.util.Map;

@Configuration
public class KafkaAvroConsumerConfig {

    @Value("${spring.kafka.bootstrap-servers}")
    private String bootstrapServers;

    @Value("${spring.kafka.consumer.group-id}")
    private String groupId;

    @Value("${spring.kafka.consumer.auto-offset-reset}")
    private String autoOffsetReset;

    @Value("${spring.kafka.properties.schema.registry.url}")
    private String schemaRegistryUrl;

    @Value("${spring.kafka.consumer.ssl.key-store-location}")
    private String keyStoreLocation;

    @Value("${spring.kafka.consumer.ssl.key-store-password}")
    private String keyStorePassword;

    @Value("${spring.kafka.consumer.ssl.trust-store-location}")
    private String trustStoreLocation;

    @Value("${spring.kafka.consumer.ssl.trust-store-password}")
    private String trustStorePassword;

    @Value("${spring.kafka.consumer.ssl.key-password}")
    private String keyPassword;

    @Bean
    public ConsumerFactory<String, Object> avroConsumerFactory() {
        Map<String, Object> props = new HashMap<>();
        // Kafka基础配置
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, autoOffsetReset);
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, KafkaAvroDeserializer.class);
        
        // Schema Registry配置
        props.put(KafkaAvroDeserializerConfig.SCHEMA_REGISTRY_URL_CONFIG, schemaRegistryUrl);
        props.put(KafkaAvroDeserializerConfig.SPECIFIC_AVRO_READER_CONFIG, false);
        
        // 同步SSL配置到Schema Registry客户端
        props.put("schema.registry.ssl.keystore.location", keyStoreLocation);
        props.put("schema.registry.ssl.keystore.password", keyStorePassword);
        props.put("schema.registry.ssl.truststore.location", trustStoreLocation);
        props.put("schema.registry.ssl.truststore.password", trustStorePassword);
        props.put("schema.registry.ssl.key.password", keyPassword);
        props.put("schema.registry.ssl.endpoint.identification.algorithm", "");

        // Kafka SSL配置保留
        props.put("ssl.keystore.location", keyStoreLocation);
        props.put("ssl.keystore.password", keyStorePassword);
        props.put("ssl.truststore.location", trustStoreLocation);
        props.put("ssl.truststore.password", trustStorePassword);
        props.put("ssl.key.password", keyPassword);
        props.put("security.protocol", "SSL");
        props.put("sasl.mechanism", "PLAIN");
        props.put("ssl.endpoint.identification.algorithm", "");

        return new DefaultKafkaConsumerFactory<>(props);
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, Object> avroKafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(avroConsumerFactory());
        return factory;
    }
}

3 调整@KafkaListener注解指定使用自定义工厂

@Service
public class ListenerService {
    // 明确指定使用我们自定义的avroKafkaListenerContainerFactory
    @KafkaListener(topics = "topic_name", containerFactory = "avroKafkaListenerContainerFactory")
    public void consumer(org.apache.avro.generic.GenericRecord record) {
        System.out.println(record);
    }
}
验证要点
  • 确保jks文件放在项目src/main/resources目录下,yaml里配置的classpath:前缀可以被Spring正确识别加载,无需手动处理文件路径
  • 生产环境可配合Spring Cloud Config、配置中心或环境变量注入敏感配置,避免明文泄露

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 14:54:04