Spring Boot Kafka消费者使用Glue Schema Registry反序列化AVRO GenericRecord异常
问题:Spring Boot消费Glue Schema Registry的Avro GenericRecord Topic失败
我有一个通过Kafka Connect写入、采用Glue Schema Registry的AVRO GENERIC_RECORD格式的Topic,普通Java程序能正常消费,但在Spring Boot应用中消费时遇到困难。
配置类代码
@EnableKafka @Configuration public class KafkaAvroConsumerConfig { @Value("${spring.kafka.bootstrap-servers}") private String brokers; @Value("${spring.kafka.consumer.group-id}") private String groupId; // 创建监听器 @Bean public ConcurrentKafkaListenerContainerFactory<GenericRecord, GenericRecord> concurrentKafkaListenerContainerFactory() { ConcurrentKafkaListenerContainerFactory<GenericRecord, GenericRecord> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); return factory; } @Bean public ConsumerFactory<GenericRecord, GenericRecord> consumerFactory() { return new DefaultKafkaConsumerFactory<>(consumerConfigs()); } @Bean public Map<String, Object> consumerConfigs() { // 创建字符串-对象映射 Map<String, Object> config = new HashMap<>(); // 添加配置项 config.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, brokers); config.put(ConsumerConfig.GROUP_ID_CONFIG, groupId); config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, GlueSchemaRegistryKafkaDeserializer.class.getName()); config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, GlueSchemaRegistryKafkaDeserializer.class.getName()); config.put(AWSSchemaRegistryConstants.AWS_REGION, region); config.put(AWSSchemaRegistryConstants.REGISTRY_NAME, registryName); config.put(AWSSchemaRegistryConstants.AVRO_RECORD_TYPE, AvroRecordType.GENERIC_RECORD.getName()); config.put(AWSSchemaRegistryConstants.SCHEMA_NAMING_GENERATION_CLASS, MySchemaNamingStrategy.class.getName()); return config; } }
监听器类代码
@Component public class KafkaAvroConsumer { @Autowired KafkaTemplate<GenericRecord, GenericRecord> kafkaTemplate; @KafkaListener(topics = "gsr1.HR.DEPARTMENTS") public void listenDepartment(ConsumerRecord<GenericRecord, GenericRecord> record) { //System.out.println("DEPARTMENTS key schema = " + record.key().getSchema().toString()); GenericRecord key = record.key(); GenericRecord value = record.value(); System.out.println(" record.key() = " + key); System.out.println(" record.value() = " + value); System.out.println(" Key DEPARTMENT_ID = " + key.get("DEPARTMENT_ID")); System.out.println(" DEPARTMENT_NAME = " + (String) value.get("DEPARTMENT_NAME")); } }
报错信息
执行GenericRecord key = record.key();时抛出类型转换异常,数据未被反序列化为GenericRecord,而是原始字节:
Caused by: java.lang.ClassCastException: class java.lang.String cannot be cast to class org.apache.avro.generic.GenericRecord (java.lang.String is in module java.base of loader 'bootstrap'; org.apache.avro.generic.GenericRecord is in unnamed module of loader 'app')
尝试的解决方案(未成功)
参考Spring文档尝试另一种写法,但无法编译,因为GlueSchemaRegistryKafkaDeserializer不支持类型参数:
public ConsumerFactory<GenericRecord, GenericRecord> consumerFactory() { Deserializer<GenericRecord> avroDeser = new GlueSchemaRegistryKafkaDeserializer(); avroDeser.configure(consumerConfigs(), false); return new DefaultKafkaConsumerFactory<>(consumerConfigs(), avroDeser, avroDeser); }
POM文件
<?xml version="1.0" encoding="UTF-8"?> <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion> <parent> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-parent</artifactId> <version>3.0.1</version> <relativePath/> </parent> <groupId>com.test</groupId> <artifactId>SpringBootKafkaAvro</artifactId> <version>0.0.1-SNAPSHOT</version> <name>SpringBootKafkaAvro</name> <description>Spring boot Kafka Avro using Glue Schema registry</description> <properties> <java.version>17</java.version> </properties> <dependencies> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter</artifactId> </dependency> <dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka</artifactId> </dependency> <dependency> <groupId>software.amazon.glue</groupId> <artifactId>schema-registry-serde</artifactId> <version>1.1.14</version> </dependency> <dependency> <groupId>com.fasterxml.jackson.core</groupId> <artifactId>jackson-databind</artifactId> </dependency> <dependency> <groupId>com.fasterxml.jackson.datatype</groupId> <artifactId>jackson-datatype-jsr310</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-test</artifactId> <scope>test</scope> </dependency> <dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka-test</artifactId> <scope>test</scope> </dependency> </dependencies> <build> <plugins> <plugin> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-maven-plugin</artifactId> </plugin> </plugins> </build> </project>
内容的提问来源于stack exchange,提问作者Anand K
相关产品推荐
相关产品推荐

