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
相关产品推荐
相关产品推荐

